factorisation de code
This commit is contained in:
@@ -11,6 +11,13 @@ use pmoaudio::AudioError;
|
||||
use pmoflac::FlacEncodedStream;
|
||||
use tokio::io::AsyncReadExt;
|
||||
|
||||
/// State for FLAC stream subscription.
|
||||
pub(crate) enum FlacStreamState {
|
||||
SendingHeader,
|
||||
Streaming,
|
||||
}
|
||||
|
||||
|
||||
/// Validate and parse FLAC block size from frame header
|
||||
///
|
||||
/// Returns the number of samples in the frame if the header is valid, or None if:
|
||||
|
||||
@@ -6,6 +6,7 @@
|
||||
|
||||
pub mod byte_stream_reader;
|
||||
pub mod chunk_to_pcm;
|
||||
pub mod streaming_icyflac_sink;
|
||||
|
||||
#[cfg(feature = "cache-sink")]
|
||||
mod flac_cache_sink;
|
||||
@@ -26,12 +27,19 @@ mod timed_broadcast;
|
||||
mod streaming_flac_sink;
|
||||
|
||||
#[cfg(feature = "http-stream")]
|
||||
pub use streaming_flac_sink::{
|
||||
FlacClientStream, IcyClientStream, MetadataSnapshot, StreamHandle, StreamingFlacSink,
|
||||
};
|
||||
mod streaming_sink_common;
|
||||
|
||||
#[cfg(feature = "http-stream")]
|
||||
pub use streaming_flac_sink::{FlacClientStream, StreamHandle, StreamingFlacSink};
|
||||
|
||||
#[cfg(feature = "http-stream")]
|
||||
pub use streaming_icyflac_sink::IcyClientStream;
|
||||
|
||||
#[cfg(feature = "http-stream")]
|
||||
mod streaming_ogg_flac_sink;
|
||||
|
||||
#[cfg(feature = "http-stream")]
|
||||
pub use streaming_ogg_flac_sink::{OggFlacClientStream, OggFlacStreamHandle, StreamingOggFlacSink};
|
||||
|
||||
#[cfg(feature = "http-stream")]
|
||||
pub use streaming_sink_common::MetadataSnapshot;
|
||||
|
||||
@@ -57,7 +57,7 @@
|
||||
use std::collections::VecDeque;
|
||||
use std::io;
|
||||
use std::pin::Pin;
|
||||
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
|
||||
use std::sync::atomic::{AtomicBool, Ordering};
|
||||
use std::sync::Arc;
|
||||
use std::task::{Context, Poll};
|
||||
use std::time::Duration;
|
||||
@@ -71,186 +71,92 @@ use async_trait::async_trait;
|
||||
use bytes::Bytes;
|
||||
use pmoaudio::{
|
||||
pipeline::{AudioPipelineNode, Node, NodeLogic, PipelineHandle, StopReason},
|
||||
AudioError, AudioSegment, SyncMarker, TypeRequirement, TypedAudioNode,
|
||||
_AudioSegment,
|
||||
AudioError, AudioSegment, SyncMarker, TypeRequirement, TypedAudioNode, _AudioSegment,
|
||||
};
|
||||
use pmoflac::{encode_flac_stream, EncoderOptions, FlacEncodedStream, PcmFormat};
|
||||
use pmometadata::TrackMetadata;
|
||||
use pmoflac::{EncoderOptions, FlacEncodedStream};
|
||||
use tokio::io::{AsyncRead, AsyncReadExt, ReadBuf};
|
||||
use tokio::sync::{mpsc, RwLock};
|
||||
use tokio_util::sync::CancellationToken;
|
||||
use tracing::{debug, error, info, trace, warn};
|
||||
|
||||
use crate::byte_stream_reader::{PcmChunk,ByteStreamReader};
|
||||
use crate::byte_stream_reader::{PcmChunk};
|
||||
use crate::chunk_to_pcm::chunk_to_pcm_bytes;
|
||||
use crate::sinks::timed_broadcast::{DEFAULT_BROADCAST_MAX_LEAD_TIME, calculate_broadcast_capacity};
|
||||
use crate::sinks::streaming_sink_common::{
|
||||
MetadataSnapshot, SharedClientStream, SharedSinkContext, SharedStreamHandleInner,
|
||||
};
|
||||
use crate::sinks::timed_broadcast::{
|
||||
calculate_broadcast_capacity, DEFAULT_BROADCAST_MAX_LEAD_TIME,
|
||||
};
|
||||
use crate::streaming_icyflac_sink::IcyClientStream;
|
||||
|
||||
/// Default ICY metadata interval (bytes of audio between metadata blocks).
|
||||
/// Standard value used by most streaming servers.
|
||||
const DEFAULT_ICY_METAINT: usize = 16000;
|
||||
|
||||
|
||||
|
||||
|
||||
/// Snapshot of track metadata at a point in time.
|
||||
///
|
||||
/// This structure is shared between the sink and clients to provide
|
||||
/// real-time metadata updates as tracks change in a continuous stream.
|
||||
#[derive(Debug, Clone, Default, serde::Serialize, serde::Deserialize)]
|
||||
pub struct MetadataSnapshot {
|
||||
/// Track title
|
||||
pub title: Option<String>,
|
||||
/// Artist name
|
||||
pub artist: Option<String>,
|
||||
/// Album name
|
||||
pub album: Option<String>,
|
||||
/// Track duration
|
||||
#[serde(skip_serializing_if = "Option::is_none")]
|
||||
pub duration: Option<Duration>,
|
||||
/// Cover image URL (external/original)
|
||||
#[serde(skip_serializing_if = "Option::is_none")]
|
||||
pub cover_url: Option<String>,
|
||||
/// Cover primary key in local cache (for constructing server URL)
|
||||
#[serde(skip_serializing_if = "Option::is_none")]
|
||||
pub cover_pk: Option<String>,
|
||||
/// Track number
|
||||
#[serde(skip_serializing_if = "Option::is_none")]
|
||||
pub track_number: Option<u32>,
|
||||
/// Album artist
|
||||
#[serde(skip_serializing_if = "Option::is_none")]
|
||||
pub album_artist: Option<String>,
|
||||
/// Genre
|
||||
#[serde(skip_serializing_if = "Option::is_none")]
|
||||
pub genre: Option<String>,
|
||||
/// Year
|
||||
#[serde(skip_serializing_if = "Option::is_none")]
|
||||
pub year: Option<u32>,
|
||||
/// Audio timestamp where this metadata became active (seconds)
|
||||
pub audio_timestamp_sec: f64,
|
||||
/// Version counter incremented on each update (for client-side change detection)
|
||||
pub version: u64,
|
||||
}
|
||||
|
||||
/// Handle for accessing the FLAC stream and metadata from HTTP handlers.
|
||||
///
|
||||
/// This handle is designed to be cloned and used by multiple HTTP clients
|
||||
/// simultaneously. Each client gets its own independent stream by subscribing.
|
||||
#[derive(Clone)]
|
||||
pub struct StreamHandle {
|
||||
/// Broadcast sender for FLAC bytes (pure mode)
|
||||
broadcast: timed_broadcast::Sender<Bytes>,
|
||||
|
||||
/// Current track metadata (read-only for consumers)
|
||||
metadata: Arc<RwLock<MetadataSnapshot>>,
|
||||
|
||||
/// Active client counter
|
||||
active_clients: Arc<AtomicUsize>,
|
||||
|
||||
/// Stop token to signal pipeline shutdown
|
||||
stop_token: CancellationToken,
|
||||
|
||||
/// Cached FLAC header (sent to new subscribers first)
|
||||
header: Arc<RwLock<Option<Bytes>>>,
|
||||
|
||||
auto_stop: Arc<AtomicBool>,
|
||||
inner: Arc<SharedStreamHandleInner>,
|
||||
}
|
||||
|
||||
impl StreamHandle {
|
||||
/// Subscribe to the FLAC stream in pure mode (no ICY metadata).
|
||||
///
|
||||
/// Returns an `AsyncRead` stream suitable for use with `tokio_util::io::ReaderStream`.
|
||||
pub fn subscribe_flac(&self) -> FlacClientStream {
|
||||
let count = self.active_clients.fetch_add(1, Ordering::SeqCst);
|
||||
debug!("New FLAC client subscribed (total: {})", count + 1);
|
||||
|
||||
FlacClientStream {
|
||||
rx: self.broadcast.subscribe(),
|
||||
buffer: VecDeque::new(),
|
||||
finished: false,
|
||||
handle: self.clone(),
|
||||
state: FlacStreamState::SendingHeader,
|
||||
current_epoch: 0,
|
||||
}
|
||||
pub fn new(inner: Arc<SharedStreamHandleInner>) -> Self {
|
||||
Self { inner }
|
||||
}
|
||||
|
||||
pub fn subscribe_flac(&self) -> FlacClientStream {
|
||||
let total = self.inner.client_connected();
|
||||
let rx = self.inner.register_client();
|
||||
debug!("New FLAC client subscribed (total: {})", total);
|
||||
FlacClientStream::new(rx, self.inner.clone())
|
||||
}
|
||||
|
||||
/// Subscribe to the FLAC stream with ICY metadata injection.
|
||||
///
|
||||
/// Returns an `AsyncRead` stream that injects ICY metadata blocks
|
||||
/// at regular intervals (default: every 16000 bytes).
|
||||
pub fn subscribe_icy(&self) -> IcyClientStream {
|
||||
self.subscribe_icy_with_interval(DEFAULT_ICY_METAINT)
|
||||
}
|
||||
|
||||
/// Subscribe to the FLAC stream with custom ICY metadata interval.
|
||||
pub fn subscribe_icy_with_interval(&self, metaint: usize) -> IcyClientStream {
|
||||
let count = self.active_clients.fetch_add(1, Ordering::SeqCst);
|
||||
let total = self.inner.client_connected();
|
||||
let rx = self.inner.register_client();
|
||||
debug!(
|
||||
"New ICY client subscribed (total: {}, metaint: {})",
|
||||
count + 1,
|
||||
metaint
|
||||
total, metaint
|
||||
);
|
||||
|
||||
IcyClientStream {
|
||||
rx: self.broadcast.subscribe(),
|
||||
metadata: self.metadata.clone(),
|
||||
metaint,
|
||||
byte_count: 0,
|
||||
buffer: VecDeque::new(),
|
||||
current_metadata_version: 0,
|
||||
cached_icy_metadata: Bytes::new(),
|
||||
finished: false,
|
||||
handle: self.clone(),
|
||||
state: FlacStreamState::SendingHeader,
|
||||
current_epoch: 0,
|
||||
}
|
||||
IcyClientStream::new(rx, self.inner.clone(), metaint)
|
||||
}
|
||||
|
||||
/// Get the current metadata snapshot.
|
||||
pub async fn get_metadata(&self) -> MetadataSnapshot {
|
||||
self.metadata.read().await.clone()
|
||||
self.inner.metadata.read().await.clone()
|
||||
}
|
||||
|
||||
/// Get the number of active clients.
|
||||
pub fn active_client_count(&self) -> usize {
|
||||
self.active_clients.load(Ordering::SeqCst)
|
||||
self.inner.active_clients.load(Ordering::SeqCst)
|
||||
}
|
||||
|
||||
/// Check if the stream should be stopped (no more clients).
|
||||
pub fn should_stop(&self) -> bool {
|
||||
self.active_clients.load(Ordering::SeqCst) == 0
|
||||
self.inner.active_clients.load(Ordering::SeqCst) == 0
|
||||
}
|
||||
|
||||
/// Enable or disable automatic pipeline shutdown when the last client disconnects.
|
||||
pub fn set_auto_stop(&self, enabled: bool) {
|
||||
self.auto_stop.store(enabled, Ordering::SeqCst);
|
||||
self.inner.auto_stop.store(enabled, Ordering::SeqCst);
|
||||
}
|
||||
}
|
||||
|
||||
/// State for FLAC stream subscription.
|
||||
enum FlacStreamState {
|
||||
SendingHeader,
|
||||
Streaming,
|
||||
}
|
||||
|
||||
/// Pure FLAC client stream (implements AsyncRead).
|
||||
///
|
||||
/// Each read pulls bytes out of a [`timed_broadcast`] receiver.
|
||||
/// If the receiver reports [`TryRecvError::Lagged`] it means the underlying
|
||||
/// queue expired packets before the client consumed them; we log the skip
|
||||
/// and immediately keep draining so that a late client can resynchronise
|
||||
/// with the latest epoch instead of stalling forever.
|
||||
pub struct FlacClientStream {
|
||||
rx: timed_broadcast::Receiver<Bytes>,
|
||||
buffer: VecDeque<u8>,
|
||||
finished: bool,
|
||||
handle: StreamHandle,
|
||||
state: FlacStreamState,
|
||||
current_epoch: u64,
|
||||
inner: SharedClientStream,
|
||||
}
|
||||
|
||||
impl FlacClientStream {
|
||||
fn new(rx: timed_broadcast::Receiver<Bytes>, handle: Arc<SharedStreamHandleInner>) -> Self {
|
||||
Self {
|
||||
inner: SharedClientStream::new(rx, handle),
|
||||
}
|
||||
}
|
||||
|
||||
pub fn current_epoch(&self) -> u64 {
|
||||
self.current_epoch
|
||||
self.inner.current_epoch()
|
||||
}
|
||||
}
|
||||
|
||||
@@ -260,496 +166,19 @@ impl AsyncRead for FlacClientStream {
|
||||
cx: &mut Context<'_>,
|
||||
buf: &mut ReadBuf<'_>,
|
||||
) -> Poll<io::Result<()>> {
|
||||
loop {
|
||||
// If in header state, send the header first
|
||||
if matches!(self.state, FlacStreamState::SendingHeader) {
|
||||
let header_opt = if let Ok(guard) = self.handle.header.try_read() {
|
||||
guard.clone()
|
||||
} else {
|
||||
None
|
||||
};
|
||||
|
||||
if let Some(header) = header_opt {
|
||||
self.buffer.extend(header.iter());
|
||||
debug!(
|
||||
"Sending cached FLAC header to new client ({} bytes)",
|
||||
header.len()
|
||||
);
|
||||
self.state = FlacStreamState::Streaming;
|
||||
continue; // Now copy header to output buffer
|
||||
} else {
|
||||
// Header not yet captured - client will receive it via broadcast
|
||||
// Skip directly to streaming to avoid blocking
|
||||
debug!("FLAC header not yet available, client will receive it via broadcast");
|
||||
self.state = FlacStreamState::Streaming;
|
||||
}
|
||||
}
|
||||
|
||||
// If we have buffered data, copy it
|
||||
if !self.buffer.is_empty() {
|
||||
let to_copy = self.buffer.len().min(buf.remaining());
|
||||
if to_copy == 0 {
|
||||
return Poll::Ready(Ok(()));
|
||||
}
|
||||
|
||||
let slice = self.buffer.make_contiguous();
|
||||
buf.put_slice(&slice[..to_copy]);
|
||||
self.buffer.drain(..to_copy);
|
||||
return Poll::Ready(Ok(()));
|
||||
}
|
||||
|
||||
if self.finished {
|
||||
return Poll::Ready(Ok(()));
|
||||
}
|
||||
|
||||
// Try to receive more data
|
||||
match self.rx.try_recv() {
|
||||
Ok(packet) => {
|
||||
self.current_epoch = packet.epoch;
|
||||
self.buffer.extend(packet.payload.iter());
|
||||
}
|
||||
Err(TryRecvError::Empty) => {
|
||||
// No data available right now.
|
||||
// Schedule a wakeup after a small delay to avoid busy-loop polling.
|
||||
let waker = cx.waker().clone();
|
||||
tokio::spawn(async move {
|
||||
tokio::time::sleep(tokio::time::Duration::from_millis(10)).await;
|
||||
waker.wake();
|
||||
});
|
||||
return Poll::Pending;
|
||||
}
|
||||
Err(TryRecvError::Lagged(skipped)) => {
|
||||
warn!("FLAC client lagged, skipped {} messages", skipped);
|
||||
// Continue to try receiving again
|
||||
}
|
||||
Err(TryRecvError::Closed) => {
|
||||
self.finished = true;
|
||||
return Poll::Ready(Ok(()));
|
||||
}
|
||||
}
|
||||
}
|
||||
Pin::new(&mut self.inner).poll_read(cx, buf)
|
||||
}
|
||||
}
|
||||
|
||||
impl Drop for FlacClientStream {
|
||||
fn drop(&mut self) {
|
||||
let count = self.handle.active_clients.fetch_sub(1, Ordering::SeqCst);
|
||||
debug!("FLAC client disconnected (remaining: {})", count - 1);
|
||||
|
||||
if count == 1 {
|
||||
if self.handle.auto_stop.load(Ordering::SeqCst) {
|
||||
debug!("Last client disconnected, signaling pipeline stop");
|
||||
self.handle.stop_token.cancel();
|
||||
} else {
|
||||
debug!("Last client disconnected, keeping pipeline alive");
|
||||
}
|
||||
}
|
||||
let remaining = self.inner.handle().client_disconnected();
|
||||
debug!("FLAC client disconnected (remaining: {})", remaining);
|
||||
}
|
||||
}
|
||||
|
||||
/// ICY-wrapped FLAC client stream (implements AsyncRead).
|
||||
///
|
||||
/// This stream injects ICY metadata blocks at regular intervals,
|
||||
/// allowing clients to display "Now Playing" information.
|
||||
/// As with [`FlacClientStream`], hitting [`TryRecvError::Lagged`]
|
||||
/// simply indicates that the timed broadcast discarded a stale chunk;
|
||||
/// the client resumes with fresh data to avoid wedging the HTTP response.
|
||||
pub struct IcyClientStream {
|
||||
rx: timed_broadcast::Receiver<Bytes>,
|
||||
metadata: Arc<RwLock<MetadataSnapshot>>,
|
||||
metaint: usize,
|
||||
byte_count: usize,
|
||||
buffer: VecDeque<u8>,
|
||||
current_metadata_version: u64,
|
||||
cached_icy_metadata: Bytes,
|
||||
finished: bool,
|
||||
handle: StreamHandle,
|
||||
state: FlacStreamState,
|
||||
current_epoch: u64,
|
||||
}
|
||||
|
||||
impl IcyClientStream {
|
||||
pub fn current_epoch(&self) -> u64 {
|
||||
self.current_epoch
|
||||
}
|
||||
}
|
||||
|
||||
impl IcyClientStream {
|
||||
/// Format metadata as ICY metadata block.
|
||||
///
|
||||
/// ICY format: StreamTitle='Artist - Title';StreamUrl='url';
|
||||
/// Padded to multiple of 16 bytes, prefixed with length byte.
|
||||
///
|
||||
/// If cover_pk is available, constructs a URL for the cover image:
|
||||
/// - If pmoserver is initialized: http://server/covers/image/{pk}/256
|
||||
/// - Otherwise: relative URL /covers/image/{pk}/256
|
||||
fn format_icy_metadata(meta: &MetadataSnapshot) -> Bytes {
|
||||
let title = meta.title.as_deref().unwrap_or("Unknown");
|
||||
let artist = meta.artist.as_deref().unwrap_or("Unknown Artist");
|
||||
|
||||
// Build ICY metadata string with cover URL if available
|
||||
let mut metadata_str = format!("StreamTitle='{} - {}';", artist, title);
|
||||
|
||||
// Add cover URL if we have a cover_pk
|
||||
if let Some(pk) = &meta.cover_pk {
|
||||
// Use relative URL /covers/image/{pk}/256
|
||||
// This works when streaming from the same server that serves covers
|
||||
// VLC and other players will resolve relative URLs correctly
|
||||
metadata_str.push_str(&format!("StreamUrl='/covers/image/{}/256';", pk));
|
||||
} else if let Some(url) = &meta.cover_url {
|
||||
// Fallback to external cover URL if no local pk
|
||||
metadata_str.push_str(&format!("StreamUrl='{}';", url));
|
||||
}
|
||||
|
||||
// ICY metadata is padded to multiple of 16 bytes
|
||||
let metadata_bytes = metadata_str.as_bytes();
|
||||
let length = metadata_bytes.len();
|
||||
let padded_length = ((length + 15) / 16) * 16;
|
||||
let length_byte = (padded_length / 16) as u8;
|
||||
|
||||
let mut result = Vec::with_capacity(1 + padded_length);
|
||||
result.push(length_byte);
|
||||
result.extend_from_slice(metadata_bytes);
|
||||
result.resize(1 + padded_length, 0); // Pad with zeros
|
||||
|
||||
Bytes::from(result)
|
||||
}
|
||||
|
||||
/// Get metadata block if it needs to be inserted.
|
||||
async fn get_metadata_if_changed(&mut self) -> Option<Bytes> {
|
||||
let meta = self.metadata.read().await;
|
||||
if meta.version > self.current_metadata_version {
|
||||
self.current_metadata_version = meta.version;
|
||||
let icy_meta = Self::format_icy_metadata(&meta);
|
||||
self.cached_icy_metadata = icy_meta.clone();
|
||||
Some(icy_meta)
|
||||
} else if self.byte_count == 0 {
|
||||
// Always send metadata at the start
|
||||
Some(self.cached_icy_metadata.clone())
|
||||
} else {
|
||||
// No change, send empty metadata block
|
||||
Some(Bytes::from(vec![0u8]))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl AsyncRead for IcyClientStream {
|
||||
fn poll_read(
|
||||
mut self: Pin<&mut Self>,
|
||||
cx: &mut Context<'_>,
|
||||
buf: &mut ReadBuf<'_>,
|
||||
) -> Poll<io::Result<()>> {
|
||||
loop {
|
||||
// If in header state, send the header first
|
||||
if matches!(self.state, FlacStreamState::SendingHeader) {
|
||||
let header_opt = if let Ok(guard) = self.handle.header.try_read() {
|
||||
guard.clone()
|
||||
} else {
|
||||
None
|
||||
};
|
||||
|
||||
if let Some(header) = header_opt {
|
||||
self.buffer.extend(header.iter());
|
||||
debug!(
|
||||
"Sending cached FLAC header to new ICY client ({} bytes)",
|
||||
header.len()
|
||||
);
|
||||
self.state = FlacStreamState::Streaming;
|
||||
continue; // Now copy header to output buffer
|
||||
} else {
|
||||
// Header not yet captured - client will receive it via broadcast
|
||||
// Skip directly to streaming to avoid blocking
|
||||
debug!(
|
||||
"FLAC header not yet available, ICY client will receive it via broadcast"
|
||||
);
|
||||
self.state = FlacStreamState::Streaming;
|
||||
}
|
||||
}
|
||||
|
||||
// If we have buffered data, copy it
|
||||
if !self.buffer.is_empty() {
|
||||
let to_copy = self.buffer.len().min(buf.remaining());
|
||||
if to_copy == 0 {
|
||||
return Poll::Ready(Ok(()));
|
||||
}
|
||||
|
||||
let slice = self.buffer.make_contiguous();
|
||||
buf.put_slice(&slice[..to_copy]);
|
||||
self.buffer.drain(..to_copy);
|
||||
return Poll::Ready(Ok(()));
|
||||
}
|
||||
|
||||
if self.finished {
|
||||
return Poll::Ready(Ok(()));
|
||||
}
|
||||
|
||||
// Check if we need to insert metadata
|
||||
if self.byte_count % self.metaint == 0 && self.byte_count > 0 {
|
||||
// Time to insert ICY metadata
|
||||
// Use try_read to avoid blocking in poll context
|
||||
let update = {
|
||||
if let Ok(meta) = self.metadata.try_read() {
|
||||
if meta.version > self.current_metadata_version {
|
||||
Some((meta.version, Self::format_icy_metadata(&meta)))
|
||||
} else {
|
||||
None
|
||||
}
|
||||
} else {
|
||||
None
|
||||
}
|
||||
};
|
||||
|
||||
if let Some((new_version, new_metadata)) = update {
|
||||
self.current_metadata_version = new_version;
|
||||
self.cached_icy_metadata = new_metadata;
|
||||
}
|
||||
|
||||
let icy_data = self.cached_icy_metadata.clone();
|
||||
self.buffer.extend(icy_data.iter());
|
||||
self.byte_count = 0; // Reset counter after metadata
|
||||
continue;
|
||||
}
|
||||
|
||||
// Try to receive audio data
|
||||
match self.rx.try_recv() {
|
||||
Ok(packet) => {
|
||||
self.current_epoch = packet.epoch;
|
||||
// Calculate how many bytes until next metadata block
|
||||
let until_metadata = self.metaint - (self.byte_count % self.metaint);
|
||||
let to_buffer = packet.payload.len().min(until_metadata);
|
||||
|
||||
self.buffer.extend(packet.payload[..to_buffer].iter());
|
||||
self.byte_count += to_buffer;
|
||||
|
||||
// If we have more data, we'll process it in the next iteration
|
||||
if to_buffer < packet.payload.len() {
|
||||
// Save remaining for next iteration
|
||||
// For now, we'll just drop it and get it again
|
||||
// TODO: Improve this
|
||||
}
|
||||
}
|
||||
Err(TryRecvError::Empty) => {
|
||||
// No data available right now.
|
||||
// Schedule a wakeup after a small delay to avoid busy-loop polling.
|
||||
let waker = cx.waker().clone();
|
||||
tokio::spawn(async move {
|
||||
tokio::time::sleep(tokio::time::Duration::from_millis(10)).await;
|
||||
waker.wake();
|
||||
});
|
||||
return Poll::Pending;
|
||||
}
|
||||
Err(TryRecvError::Lagged(skipped)) => {
|
||||
warn!("ICY client lagged, skipped {} messages", skipped);
|
||||
}
|
||||
Err(TryRecvError::Closed) => {
|
||||
self.finished = true;
|
||||
return Poll::Ready(Ok(()));
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl Drop for IcyClientStream {
|
||||
fn drop(&mut self) {
|
||||
let count = self.handle.active_clients.fetch_sub(1, Ordering::SeqCst);
|
||||
debug!("ICY client disconnected (remaining: {})", count - 1);
|
||||
|
||||
if count == 1 {
|
||||
if self.handle.auto_stop.load(Ordering::SeqCst) {
|
||||
debug!("Last client disconnected, signaling pipeline stop");
|
||||
self.handle.stop_token.cancel();
|
||||
} else {
|
||||
debug!("Last client disconnected, keeping pipeline alive");
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Internal state for encoder initialization.
|
||||
struct EncoderState {
|
||||
broadcaster_task: tokio::task::JoinHandle<()>,
|
||||
}
|
||||
|
||||
/// Logic for the streaming FLAC sink.
|
||||
struct StreamingFlacSinkLogic {
|
||||
encoder_options: EncoderOptions,
|
||||
bits_per_sample: u8,
|
||||
pcm_tx: Option<mpsc::Sender<PcmChunk>>,
|
||||
pcm_rx: Option<mpsc::Receiver<PcmChunk>>,
|
||||
metadata: Arc<RwLock<MetadataSnapshot>>,
|
||||
broadcast: timed_broadcast::Sender<Bytes>,
|
||||
header: Arc<RwLock<Option<Bytes>>>,
|
||||
encoder_state: Option<EncoderState>,
|
||||
sample_rate: Option<u32>,
|
||||
broadcast_max_lead_time: f64,
|
||||
first_chunk_timestamp_checked: bool,
|
||||
/// Accumulated timestamp offset for maintaining continuity across encoder restarts
|
||||
timestamp_offset_sec: f64,
|
||||
/// Shared timestamp reference to read the last broadcast timestamp
|
||||
current_timestamp: Arc<RwLock<f64>>,
|
||||
}
|
||||
|
||||
impl StreamingFlacSinkLogic {
|
||||
/// Initialize the FLAC encoder once we know the sample rate.
|
||||
async fn initialize_encoder(&mut self, sample_rate: u32, timestamp_offset_sec: f64) -> Result<(), AudioError> {
|
||||
if self.encoder_state.is_some() {
|
||||
return Ok(()); // Already initialized
|
||||
}
|
||||
|
||||
debug!(
|
||||
"Initializing FLAC encoder with sample rate: {} Hz",
|
||||
sample_rate
|
||||
);
|
||||
|
||||
// Take the PCM receiver (we only initialize once)
|
||||
let pcm_rx = self
|
||||
.pcm_rx
|
||||
.take()
|
||||
.ok_or_else(|| AudioError::ProcessingError("PCM receiver already consumed".into()))?;
|
||||
|
||||
// Create shared timestamp and duration for pacing
|
||||
// Reuse existing current_timestamp Arc to maintain reference for reading later
|
||||
let current_timestamp = self.current_timestamp.clone();
|
||||
let current_duration = Arc::new(RwLock::new(0.0f64));
|
||||
|
||||
// Create ByteStreamReader for the encoder
|
||||
let pcm_reader =
|
||||
ByteStreamReader::new(pcm_rx, current_timestamp.clone(), current_duration.clone());
|
||||
|
||||
// Create PCM format
|
||||
let pcm_format = PcmFormat {
|
||||
sample_rate,
|
||||
channels: 2,
|
||||
bits_per_sample: self.bits_per_sample,
|
||||
};
|
||||
|
||||
// Start the FLAC encoder
|
||||
let flac_stream = encode_flac_stream(pcm_reader, pcm_format, self.encoder_options.clone())
|
||||
.await
|
||||
.map_err(|e| {
|
||||
AudioError::ProcessingError(format!("Failed to start FLAC encoder: {}", e))
|
||||
})?;
|
||||
|
||||
debug!("FLAC encoder initialized successfully");
|
||||
|
||||
// Spawn broadcaster task with timestamp and duration for pacing
|
||||
let broadcast = self.broadcast.clone();
|
||||
let header = self.header.clone();
|
||||
let max_lead = self.broadcast_max_lead_time;
|
||||
let broadcaster_task = tokio::spawn(async move {
|
||||
if let Err(e) = broadcast_flac_stream(
|
||||
flac_stream,
|
||||
broadcast,
|
||||
header,
|
||||
current_timestamp,
|
||||
current_duration,
|
||||
max_lead,
|
||||
sample_rate,
|
||||
timestamp_offset_sec,
|
||||
)
|
||||
.await
|
||||
{
|
||||
error!("Broadcaster task error: {}", e);
|
||||
}
|
||||
});
|
||||
|
||||
self.encoder_state = Some(EncoderState { broadcaster_task });
|
||||
|
||||
debug!("Broadcaster task spawned");
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Update metadata from a TrackBoundary marker.
|
||||
async fn update_metadata(
|
||||
&mut self,
|
||||
metadata_lock: &Arc<RwLock<dyn TrackMetadata>>,
|
||||
timestamp_sec: f64,
|
||||
) -> Result<(), AudioError> {
|
||||
let metadata = metadata_lock.read().await;
|
||||
|
||||
let mut snapshot = self.metadata.write().await;
|
||||
|
||||
// Extract all metadata fields
|
||||
snapshot.title = metadata.get_title().await.ok().flatten();
|
||||
snapshot.artist = metadata.get_artist().await.ok().flatten();
|
||||
snapshot.album = metadata.get_album().await.ok().flatten();
|
||||
snapshot.duration = metadata.get_duration().await.ok().flatten();
|
||||
snapshot.cover_url = metadata.get_cover_url().await.ok().flatten();
|
||||
snapshot.cover_pk = metadata.get_cover_pk().await.ok().flatten();
|
||||
snapshot.year = metadata.get_year().await.ok().flatten();
|
||||
|
||||
// Extract extra fields
|
||||
if let Ok(Some(extra)) = metadata.get_extra().await {
|
||||
snapshot.genre = extra.get("genre").cloned();
|
||||
snapshot.track_number = extra
|
||||
.get("track_number")
|
||||
.and_then(|s| s.parse::<u32>().ok());
|
||||
}
|
||||
|
||||
snapshot.audio_timestamp_sec = timestamp_sec;
|
||||
snapshot.version += 1;
|
||||
|
||||
debug!(
|
||||
"Metadata updated: v{} @ {:.2}s - {} - {} (cover_pk: {:?})",
|
||||
snapshot.version,
|
||||
timestamp_sec,
|
||||
snapshot.artist.as_deref().unwrap_or("?"),
|
||||
snapshot.title.as_deref().unwrap_or("?"),
|
||||
snapshot.cover_pk
|
||||
);
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Restart the FLAC encoder for a new track.
|
||||
/// This causes a new "fLaC" header to be emitted and timestamps to reset to 0.
|
||||
async fn restart_encoder_for_new_track(&mut self) -> Result<(), AudioError> {
|
||||
let sample_rate = self
|
||||
.sample_rate
|
||||
.ok_or_else(|| AudioError::ProcessingError("Sample rate not initialized".into()))?;
|
||||
|
||||
debug!("Restarting FLAC encoder for new track");
|
||||
|
||||
// 1. Read current timestamp to maintain continuity
|
||||
let last_timestamp = *self.current_timestamp.read().await;
|
||||
debug!("Last timestamp before restart: {:.3}s", last_timestamp);
|
||||
|
||||
// 2. Close current PCM sender to signal encoder to finish
|
||||
if let Some(tx) = self.pcm_tx.take() {
|
||||
drop(tx);
|
||||
trace!("Dropped PCM sender to signal encoder finish");
|
||||
}
|
||||
|
||||
// 3. Wait for current broadcaster task to complete
|
||||
if let Some(state) = self.encoder_state.take() {
|
||||
trace!("Waiting for broadcaster task to finish...");
|
||||
match state.broadcaster_task.await {
|
||||
Ok(_) => {
|
||||
trace!("Broadcaster task finished successfully");
|
||||
}
|
||||
Err(e) => {
|
||||
warn!("Broadcaster task error during restart: {:?}", e);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// 4. Update timestamp offset to maintain continuity across encoder restart
|
||||
self.timestamp_offset_sec += last_timestamp;
|
||||
debug!("New timestamp offset: {:.3}s", self.timestamp_offset_sec);
|
||||
|
||||
// 5. Create new PCM channel
|
||||
let (pcm_tx, pcm_rx) = mpsc::channel::<PcmChunk>(16);
|
||||
self.pcm_tx = Some(pcm_tx);
|
||||
self.pcm_rx = Some(pcm_rx);
|
||||
|
||||
// 6. Initialize new encoder with timestamp offset (will naturally emit new "fLaC" header)
|
||||
self.initialize_encoder(sample_rate, self.timestamp_offset_sec).await?;
|
||||
|
||||
debug!("FLAC encoder restarted successfully for new track");
|
||||
Ok(())
|
||||
}
|
||||
ctx: SharedSinkContext,
|
||||
}
|
||||
|
||||
#[async_trait]
|
||||
@@ -781,8 +210,8 @@ impl NodeLogic for StreamingFlacSinkLogic {
|
||||
Some(seg) => {
|
||||
match &seg.segment {
|
||||
_AudioSegment::Chunk(chunk) => {
|
||||
if !self.first_chunk_timestamp_checked {
|
||||
self.first_chunk_timestamp_checked = true;
|
||||
if !self.ctx.first_chunk_timestamp_checked {
|
||||
self.ctx.first_chunk_timestamp_checked = true;
|
||||
if seg.timestamp_sec.abs() > 1e-6 {
|
||||
warn!(
|
||||
"StreamingFlacSink: first chunk timestamp is {:.6}s (expected 0.0)",
|
||||
@@ -794,21 +223,44 @@ impl NodeLogic for StreamingFlacSinkLogic {
|
||||
}
|
||||
|
||||
// Detect sample rate from first chunk and initialize encoder
|
||||
if self.sample_rate.is_none() {
|
||||
if self.ctx.sample_rate.is_none() {
|
||||
let sample_rate = chunk.sample_rate();
|
||||
self.sample_rate = Some(sample_rate);
|
||||
self.ctx.sample_rate = Some(sample_rate);
|
||||
debug!("Detected sample rate: {} Hz", sample_rate);
|
||||
|
||||
// Initialize the FLAC encoder now (first track starts at 0.0)
|
||||
self.initialize_encoder(sample_rate, 0.0).await?;
|
||||
self.ctx
|
||||
.initialize_encoder(
|
||||
sample_rate,
|
||||
0.0,
|
||||
|flac_stream,
|
||||
broadcast,
|
||||
header,
|
||||
current_timestamp,
|
||||
current_duration,
|
||||
max_lead,
|
||||
sample_rate,
|
||||
timestamp_offset_sec| {
|
||||
broadcast_flac_stream(
|
||||
flac_stream,
|
||||
broadcast,
|
||||
header,
|
||||
current_timestamp,
|
||||
current_duration,
|
||||
max_lead,
|
||||
sample_rate,
|
||||
timestamp_offset_sec,
|
||||
)
|
||||
},
|
||||
)
|
||||
.await?;
|
||||
}
|
||||
|
||||
// Convert chunk to PCM bytes
|
||||
let pcm_bytes = chunk_to_pcm_bytes(&chunk, self.bits_per_sample)?;
|
||||
let pcm_bytes = chunk_to_pcm_bytes(&chunk, self.ctx.bits_per_sample)?;
|
||||
|
||||
// Calculate exact duration from samples and sample rate
|
||||
let sample_rate = self
|
||||
.sample_rate
|
||||
let sample_rate = self.ctx.sample_rate
|
||||
.expect("sample_rate should be initialized");
|
||||
let duration_sec = chunk.len() as f64 / sample_rate as f64;
|
||||
|
||||
@@ -829,7 +281,7 @@ impl NodeLogic for StreamingFlacSinkLogic {
|
||||
let send_start = std::time::Instant::now();
|
||||
|
||||
// Get the sender (it should always be Some after initialization)
|
||||
let pcm_tx = match &self.pcm_tx {
|
||||
let pcm_tx = match &self.ctx.pcm_tx {
|
||||
Some(tx) => tx,
|
||||
None => {
|
||||
error!("PCM sender not initialized");
|
||||
@@ -855,9 +307,33 @@ impl NodeLogic for StreamingFlacSinkLogic {
|
||||
match marker.as_ref() {
|
||||
SyncMarker::TrackBoundary { metadata } => {
|
||||
// Only restart encoder if it's already initialized (not the first track)
|
||||
if self.sample_rate.is_some() && self.encoder_state.is_some() {
|
||||
if self.ctx.sample_rate.is_some() && self.ctx.encoder_state.is_some() {
|
||||
// Restart encoder to emit new header and reset timestamps
|
||||
if let Err(e) = self.restart_encoder_for_new_track().await {
|
||||
if let Err(e) = self
|
||||
.ctx
|
||||
.restart_encoder_for_new_track(
|
||||
|flac_stream,
|
||||
broadcast,
|
||||
header,
|
||||
current_timestamp,
|
||||
current_duration,
|
||||
max_lead,
|
||||
sample_rate,
|
||||
timestamp_offset_sec| {
|
||||
broadcast_flac_stream(
|
||||
flac_stream,
|
||||
broadcast,
|
||||
header,
|
||||
current_timestamp,
|
||||
current_duration,
|
||||
max_lead,
|
||||
sample_rate,
|
||||
timestamp_offset_sec,
|
||||
)
|
||||
},
|
||||
)
|
||||
.await
|
||||
{
|
||||
error!("Failed to restart encoder for new track: {}", e);
|
||||
break;
|
||||
}
|
||||
@@ -866,7 +342,7 @@ impl NodeLogic for StreamingFlacSinkLogic {
|
||||
}
|
||||
|
||||
// Update metadata for the new track
|
||||
if let Err(e) = self.update_metadata(metadata, seg.timestamp_sec).await {
|
||||
if let Err(e) = self.ctx.update_metadata(metadata, seg.timestamp_sec).await {
|
||||
error!("Failed to update metadata: {}", e);
|
||||
}
|
||||
}
|
||||
@@ -903,7 +379,6 @@ impl NodeLogic for StreamingFlacSinkLogic {
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
/// Streaming FLAC sink for multi-client HTTP streaming.
|
||||
pub struct StreamingFlacSink {
|
||||
inner: Node<StreamingFlacSinkLogic>,
|
||||
@@ -922,10 +397,7 @@ impl StreamingFlacSink {
|
||||
/// A tuple of `(sink, handle)` where:
|
||||
/// - `sink` is added to the audio pipeline
|
||||
/// - `handle` is used by HTTP handlers to serve streams
|
||||
pub fn new(
|
||||
encoder_options: EncoderOptions,
|
||||
bits_per_sample: u8,
|
||||
) -> (Self, StreamHandle) {
|
||||
pub fn new(encoder_options: EncoderOptions, bits_per_sample: u8) -> (Self, StreamHandle) {
|
||||
Self::with_max_broadcast_lead(
|
||||
encoder_options,
|
||||
bits_per_sample,
|
||||
@@ -954,8 +426,7 @@ impl StreamingFlacSink {
|
||||
let broadcast_capacity = calculate_broadcast_capacity(broadcast_max_lead_time);
|
||||
debug!(
|
||||
"Streaming Sink: using broadcast capacity of {} items (max_lead_time={:.1}s)",
|
||||
broadcast_capacity,
|
||||
broadcast_max_lead_time
|
||||
broadcast_capacity, broadcast_max_lead_time
|
||||
);
|
||||
|
||||
// Broadcast channel for FLAC bytes
|
||||
@@ -966,31 +437,34 @@ impl StreamingFlacSink {
|
||||
|
||||
// Stop token and client counter
|
||||
let stop_token = CancellationToken::new();
|
||||
let active_clients = Arc::new(AtomicUsize::new(0));
|
||||
let auto_stop = Arc::new(AtomicBool::new(true));
|
||||
|
||||
let handle = StreamHandle {
|
||||
broadcast: broadcast.clone(),
|
||||
metadata: metadata.clone(),
|
||||
active_clients,
|
||||
stop_token: stop_token.clone(),
|
||||
header: header.clone(),
|
||||
auto_stop: Arc::new(AtomicBool::new(true)),
|
||||
};
|
||||
let shared_handle = Arc::new(SharedStreamHandleInner::new(
|
||||
broadcast.clone(),
|
||||
metadata.clone(),
|
||||
stop_token.clone(),
|
||||
header.clone(),
|
||||
auto_stop.clone(),
|
||||
));
|
||||
|
||||
let handle = StreamHandle::new(shared_handle.clone());
|
||||
|
||||
let logic = StreamingFlacSinkLogic {
|
||||
encoder_options,
|
||||
bits_per_sample,
|
||||
pcm_tx: Some(pcm_tx),
|
||||
pcm_rx: Some(pcm_rx),
|
||||
metadata,
|
||||
broadcast,
|
||||
header,
|
||||
encoder_state: None,
|
||||
sample_rate: None,
|
||||
broadcast_max_lead_time: broadcast_max_lead_time.max(0.0),
|
||||
first_chunk_timestamp_checked: false,
|
||||
timestamp_offset_sec: 0.0,
|
||||
current_timestamp: Arc::new(RwLock::new(0.0)),
|
||||
ctx: SharedSinkContext {
|
||||
encoder_options,
|
||||
bits_per_sample,
|
||||
pcm_tx: Some(pcm_tx),
|
||||
pcm_rx: Some(pcm_rx),
|
||||
metadata,
|
||||
broadcast,
|
||||
header,
|
||||
encoder_state: None,
|
||||
sample_rate: None,
|
||||
broadcast_max_lead_time: broadcast_max_lead_time.max(0.0),
|
||||
first_chunk_timestamp_checked: false,
|
||||
timestamp_offset_sec: 0.0,
|
||||
current_timestamp: Arc::new(RwLock::new(0.0)),
|
||||
},
|
||||
};
|
||||
|
||||
let sink = Self {
|
||||
@@ -1030,7 +504,6 @@ impl TypedAudioNode for StreamingFlacSink {
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
/// Broadcaster task: reads FLAC bytes from encoder and broadcasts to all clients.
|
||||
/// Implements precise real-time pacing based on audio timestamps.
|
||||
/// Ensures data is sent at FLAC frame boundaries to prevent sync errors in strict decoders like FFPlay.
|
||||
@@ -1065,9 +538,8 @@ async fn broadcast_flac_stream(
|
||||
let mut read_count = 0u64;
|
||||
let sample_rate_f64 = sample_rate as f64;
|
||||
|
||||
// Sample counter for calculating accurate timestamps (reset on new headers)
|
||||
let mut encoded_samples = 0u64;
|
||||
|
||||
// Sample counter for calculating accurate timestamps (reset on new headers)
|
||||
let mut encoded_samples = 0u64;
|
||||
|
||||
loop {
|
||||
let read_start = std::time::Instant::now();
|
||||
@@ -1157,7 +629,8 @@ async fn broadcast_flac_stream(
|
||||
// Calculer le timestamp de cette FLAC frame (avec offset pour continuité entre tracks)
|
||||
let frame_start_samples = encoded_samples;
|
||||
encoded_samples = encoded_samples.saturating_add(total_samples);
|
||||
let audio_timestamp = timestamp_offset_sec + (frame_start_samples as f64 / sample_rate_f64);
|
||||
let audio_timestamp =
|
||||
timestamp_offset_sec + (frame_start_samples as f64 / sample_rate_f64);
|
||||
let segment_duration = total_samples as f64 / sample_rate_f64;
|
||||
|
||||
if stats_last_log.elapsed() >= Duration::from_secs(1) {
|
||||
@@ -1278,4 +751,3 @@ async fn broadcast_flac_stream(
|
||||
trace!("Broadcaster task completed successfully");
|
||||
Ok(())
|
||||
}
|
||||
|
||||
|
||||
239
pmoaudio-ext/src/sinks/streaming_icyflac_sink.rs
Normal file
239
pmoaudio-ext/src/sinks/streaming_icyflac_sink.rs
Normal file
@@ -0,0 +1,239 @@
|
||||
use std::{collections::VecDeque, pin::Pin, sync::Arc, task::{Context, Poll}};
|
||||
|
||||
use tokio::{io::{AsyncRead, ReadBuf}, sync::RwLock};
|
||||
|
||||
use crate::{MetadataSnapshot, sinks::{flac_frame_utils::FlacStreamState, streaming_sink_common::SharedStreamHandleInner, timed_broadcast::{self, TryRecvError}}};
|
||||
use bytes::Bytes;
|
||||
use std::io;
|
||||
|
||||
use tracing::{debug, error, info, trace, warn};
|
||||
|
||||
/// ICY-wrapped FLAC client stream (implements AsyncRead).
|
||||
///
|
||||
/// This stream injects ICY metadata blocks at regular intervals,
|
||||
/// allowing clients to display "Now Playing" information.
|
||||
/// As with [`FlacClientStream`], hitting [`TryRecvError::Lagged`]
|
||||
/// simply indicates that the timed broadcast discarded a stale chunk;
|
||||
/// the client resumes with fresh data to avoid wedging the HTTP response.
|
||||
pub struct IcyClientStream {
|
||||
rx: timed_broadcast::Receiver<Bytes>,
|
||||
metadata: Arc<RwLock<MetadataSnapshot>>,
|
||||
metaint: usize,
|
||||
byte_count: usize,
|
||||
buffer: VecDeque<u8>,
|
||||
current_metadata_version: u64,
|
||||
cached_icy_metadata: Bytes,
|
||||
finished: bool,
|
||||
handle: Arc<SharedStreamHandleInner>,
|
||||
state: FlacStreamState,
|
||||
current_epoch: u64,
|
||||
}
|
||||
|
||||
impl IcyClientStream {
|
||||
pub(crate) fn new(
|
||||
rx: timed_broadcast::Receiver<Bytes>,
|
||||
handle: Arc<SharedStreamHandleInner>,
|
||||
metaint: usize,
|
||||
) -> Self {
|
||||
Self {
|
||||
rx,
|
||||
metadata: handle.metadata.clone(),
|
||||
metaint,
|
||||
byte_count: 0,
|
||||
buffer: VecDeque::new(),
|
||||
current_metadata_version: 0,
|
||||
cached_icy_metadata: Bytes::new(),
|
||||
finished: false,
|
||||
handle,
|
||||
state: FlacStreamState::SendingHeader,
|
||||
current_epoch: 0,
|
||||
}
|
||||
}
|
||||
|
||||
pub fn current_epoch(&self) -> u64 {
|
||||
self.current_epoch
|
||||
}
|
||||
}
|
||||
|
||||
impl IcyClientStream {
|
||||
/// Format metadata as ICY metadata block.
|
||||
///
|
||||
/// ICY format: StreamTitle='Artist - Title';StreamUrl='url';
|
||||
/// Padded to multiple of 16 bytes, prefixed with length byte.
|
||||
///
|
||||
/// If cover_pk is available, constructs a URL for the cover image:
|
||||
/// - If pmoserver is initialized: http://server/covers/image/{pk}/256
|
||||
/// - Otherwise: relative URL /covers/image/{pk}/256
|
||||
fn format_icy_metadata(meta: &MetadataSnapshot) -> Bytes {
|
||||
let title = meta.title.as_deref().unwrap_or("Unknown");
|
||||
let artist = meta.artist.as_deref().unwrap_or("Unknown Artist");
|
||||
|
||||
// Build ICY metadata string with cover URL if available
|
||||
let mut metadata_str = format!("StreamTitle='{} - {}';", artist, title);
|
||||
|
||||
// Add cover URL if we have a cover_pk
|
||||
if let Some(pk) = &meta.cover_pk {
|
||||
// Use relative URL /covers/image/{pk}/256
|
||||
// This works when streaming from the same server that serves covers
|
||||
// VLC and other players will resolve relative URLs correctly
|
||||
metadata_str.push_str(&format!("StreamUrl='/covers/image/{}/256';", pk));
|
||||
} else if let Some(url) = &meta.cover_url {
|
||||
// Fallback to external cover URL if no local pk
|
||||
metadata_str.push_str(&format!("StreamUrl='{}';", url));
|
||||
}
|
||||
|
||||
// ICY metadata is padded to multiple of 16 bytes
|
||||
let metadata_bytes = metadata_str.as_bytes();
|
||||
let length = metadata_bytes.len();
|
||||
let padded_length = ((length + 15) / 16) * 16;
|
||||
let length_byte = (padded_length / 16) as u8;
|
||||
|
||||
let mut result = Vec::with_capacity(1 + padded_length);
|
||||
result.push(length_byte);
|
||||
result.extend_from_slice(metadata_bytes);
|
||||
result.resize(1 + padded_length, 0); // Pad with zeros
|
||||
|
||||
Bytes::from(result)
|
||||
}
|
||||
|
||||
/// Get metadata block if it needs to be inserted.
|
||||
async fn get_metadata_if_changed(&mut self) -> Option<Bytes> {
|
||||
let meta = self.metadata.read().await;
|
||||
if meta.version > self.current_metadata_version {
|
||||
self.current_metadata_version = meta.version;
|
||||
let icy_meta = Self::format_icy_metadata(&meta);
|
||||
self.cached_icy_metadata = icy_meta.clone();
|
||||
Some(icy_meta)
|
||||
} else if self.byte_count == 0 {
|
||||
// Always send metadata at the start
|
||||
Some(self.cached_icy_metadata.clone())
|
||||
} else {
|
||||
// No change, send empty metadata block
|
||||
Some(Bytes::from(vec![0u8]))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl AsyncRead for IcyClientStream {
|
||||
fn poll_read(
|
||||
mut self: Pin<&mut Self>,
|
||||
cx: &mut Context<'_>,
|
||||
buf: &mut ReadBuf<'_>,
|
||||
) -> Poll<io::Result<()>> {
|
||||
loop {
|
||||
// If in header state, send the header first
|
||||
if matches!(self.state, FlacStreamState::SendingHeader) {
|
||||
let header_opt = if let Ok(guard) = self.handle.header.try_read() {
|
||||
guard.clone()
|
||||
} else {
|
||||
None
|
||||
};
|
||||
|
||||
if let Some(header) = header_opt {
|
||||
self.buffer.extend(header.iter());
|
||||
debug!(
|
||||
"Sending cached FLAC header to new ICY client ({} bytes)",
|
||||
header.len()
|
||||
);
|
||||
self.state = FlacStreamState::Streaming;
|
||||
continue; // Now copy header to output buffer
|
||||
} else {
|
||||
// Header not yet captured - client will receive it via broadcast
|
||||
// Skip directly to streaming to avoid blocking
|
||||
debug!(
|
||||
"FLAC header not yet available, ICY client will receive it via broadcast"
|
||||
);
|
||||
self.state = FlacStreamState::Streaming;
|
||||
}
|
||||
}
|
||||
|
||||
// If we have buffered data, copy it
|
||||
if !self.buffer.is_empty() {
|
||||
let to_copy = self.buffer.len().min(buf.remaining());
|
||||
if to_copy == 0 {
|
||||
return Poll::Ready(Ok(()));
|
||||
}
|
||||
|
||||
let slice = self.buffer.make_contiguous();
|
||||
buf.put_slice(&slice[..to_copy]);
|
||||
self.buffer.drain(..to_copy);
|
||||
return Poll::Ready(Ok(()));
|
||||
}
|
||||
|
||||
if self.finished {
|
||||
return Poll::Ready(Ok(()));
|
||||
}
|
||||
|
||||
// Check if we need to insert metadata
|
||||
if self.byte_count % self.metaint == 0 && self.byte_count > 0 {
|
||||
// Time to insert ICY metadata
|
||||
// Use try_read to avoid blocking in poll context
|
||||
let update = {
|
||||
if let Ok(meta) = self.metadata.try_read() {
|
||||
if meta.version > self.current_metadata_version {
|
||||
Some((meta.version, Self::format_icy_metadata(&meta)))
|
||||
} else {
|
||||
None
|
||||
}
|
||||
} else {
|
||||
None
|
||||
}
|
||||
};
|
||||
|
||||
if let Some((new_version, new_metadata)) = update {
|
||||
self.current_metadata_version = new_version;
|
||||
self.cached_icy_metadata = new_metadata;
|
||||
}
|
||||
|
||||
let icy_data = self.cached_icy_metadata.clone();
|
||||
self.buffer.extend(icy_data.iter());
|
||||
self.byte_count = 0; // Reset counter after metadata
|
||||
continue;
|
||||
}
|
||||
|
||||
// Try to receive audio data
|
||||
match self.rx.try_recv() {
|
||||
Ok(packet) => {
|
||||
self.current_epoch = packet.epoch;
|
||||
// Calculate how many bytes until next metadata block
|
||||
let until_metadata = self.metaint - (self.byte_count % self.metaint);
|
||||
let to_buffer = packet.payload.len().min(until_metadata);
|
||||
|
||||
self.buffer.extend(packet.payload[..to_buffer].iter());
|
||||
self.byte_count += to_buffer;
|
||||
|
||||
// If we have more data, we'll process it in the next iteration
|
||||
if to_buffer < packet.payload.len() {
|
||||
// Save remaining for next iteration
|
||||
// For now, we'll just drop it and get it again
|
||||
// TODO: Improve this
|
||||
}
|
||||
}
|
||||
Err(TryRecvError::Empty) => {
|
||||
// No data available right now.
|
||||
// Schedule a wakeup after a small delay to avoid busy-loop polling.
|
||||
let waker = cx.waker().clone();
|
||||
tokio::spawn(async move {
|
||||
tokio::time::sleep(tokio::time::Duration::from_millis(10)).await;
|
||||
waker.wake();
|
||||
});
|
||||
return Poll::Pending;
|
||||
}
|
||||
Err(TryRecvError::Lagged(skipped)) => {
|
||||
warn!("ICY client lagged, skipped {} messages", skipped);
|
||||
}
|
||||
Err(TryRecvError::Closed) => {
|
||||
self.finished = true;
|
||||
return Poll::Ready(Ok(()));
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl Drop for IcyClientStream {
|
||||
fn drop(&mut self) {
|
||||
let remaining = self.handle.client_disconnected();
|
||||
debug!("ICY client disconnected (remaining: {})", remaining);
|
||||
}
|
||||
}
|
||||
@@ -49,7 +49,7 @@
|
||||
use std::collections::VecDeque;
|
||||
use std::io;
|
||||
use std::pin::Pin;
|
||||
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
|
||||
use std::sync::atomic::{AtomicBool, Ordering};
|
||||
use std::sync::Arc;
|
||||
use std::task::{Context, Poll};
|
||||
|
||||
@@ -62,99 +62,69 @@ use async_trait::async_trait;
|
||||
use bytes::Bytes;
|
||||
use pmoaudio::{
|
||||
pipeline::{AudioPipelineNode, Node, NodeLogic, PipelineHandle, StopReason},
|
||||
AudioError, AudioSegment, SyncMarker, TypeRequirement, TypedAudioNode,
|
||||
_AudioSegment,
|
||||
AudioError, AudioSegment, SyncMarker, TypeRequirement, TypedAudioNode, _AudioSegment,
|
||||
};
|
||||
use pmoflac::{encode_flac_stream, EncoderOptions, FlacEncodedStream, PcmFormat};
|
||||
use pmometadata::TrackMetadata;
|
||||
use pmoflac::{EncoderOptions, FlacEncodedStream};
|
||||
use tokio::io::{AsyncRead, AsyncReadExt, ReadBuf};
|
||||
use tokio::sync::{mpsc, RwLock};
|
||||
use tokio_util::sync::CancellationToken;
|
||||
use tracing::{debug, error, info, trace, warn};
|
||||
|
||||
use crate::byte_stream_reader::{PcmChunk,ByteStreamReader};
|
||||
use crate::byte_stream_reader::{PcmChunk};
|
||||
use crate::chunk_to_pcm::chunk_to_pcm_bytes;
|
||||
use crate::sinks::flac_frame_utils::{extract_sample_rate_from_streaminfo, read_flac_header};
|
||||
use crate::sinks::timed_broadcast::{DEFAULT_BROADCAST_MAX_LEAD_TIME, calculate_broadcast_capacity};
|
||||
|
||||
/// Snapshot of track metadata (reuse from streaming_flac_sink)
|
||||
pub use super::streaming_flac_sink::MetadataSnapshot;
|
||||
use crate::sinks::streaming_sink_common::{
|
||||
MetadataSnapshot, SharedClientStream, SharedSinkContext, SharedStreamHandleInner,
|
||||
};
|
||||
use crate::sinks::timed_broadcast::{
|
||||
calculate_broadcast_capacity, DEFAULT_BROADCAST_MAX_LEAD_TIME,
|
||||
};
|
||||
|
||||
/// Handle for accessing the OGG-FLAC stream and metadata from HTTP handlers.
|
||||
#[derive(Clone)]
|
||||
pub struct OggFlacStreamHandle {
|
||||
/// Broadcast sender for OGG-FLAC bytes
|
||||
|
||||
broadcast: timed_broadcast::Sender<Bytes>,
|
||||
|
||||
/// Current track metadata (read-only for consumers)
|
||||
metadata: Arc<RwLock<MetadataSnapshot>>,
|
||||
|
||||
/// Active client counter
|
||||
active_clients: Arc<AtomicUsize>,
|
||||
|
||||
/// Stop token to signal pipeline shutdown
|
||||
stop_token: CancellationToken,
|
||||
|
||||
/// Cached OGG-FLAC header (sent to new subscribers first)
|
||||
header: Arc<RwLock<Option<Bytes>>>,
|
||||
|
||||
auto_stop: Arc<AtomicBool>,
|
||||
inner: Arc<SharedStreamHandleInner>,
|
||||
}
|
||||
|
||||
impl OggFlacStreamHandle {
|
||||
/// Subscribe to the OGG-FLAC stream.
|
||||
///
|
||||
/// Returns an `AsyncRead` stream suitable for use with `tokio_util::io::ReaderStream`.
|
||||
pub fn new(inner: Arc<SharedStreamHandleInner>) -> Self {
|
||||
Self { inner }
|
||||
}
|
||||
|
||||
pub fn subscribe(&self) -> OggFlacClientStream {
|
||||
let count = self.active_clients.fetch_add(1, Ordering::SeqCst);
|
||||
debug!("New OGG-FLAC client subscribed (total: {})", count + 1);
|
||||
|
||||
OggFlacClientStream {
|
||||
rx: self.broadcast.subscribe(),
|
||||
buffer: VecDeque::new(),
|
||||
finished: false,
|
||||
handle: self.clone(),
|
||||
state: OggFlacStreamState::SendingHeader,
|
||||
current_epoch: 0,
|
||||
}
|
||||
let total = self.inner.client_connected();
|
||||
let rx = self.inner.register_client();
|
||||
debug!("New OGG-FLAC client subscribed (total: {})", total);
|
||||
OggFlacClientStream::new(rx, self.inner.clone())
|
||||
}
|
||||
|
||||
/// Get the current metadata snapshot.
|
||||
pub async fn get_metadata(&self) -> MetadataSnapshot {
|
||||
self.metadata.read().await.clone()
|
||||
self.inner.metadata.read().await.clone()
|
||||
}
|
||||
|
||||
/// Get the number of active clients.
|
||||
pub fn active_client_count(&self) -> usize {
|
||||
self.active_clients.load(Ordering::SeqCst)
|
||||
self.inner.active_clients.load(Ordering::SeqCst)
|
||||
}
|
||||
|
||||
pub fn set_auto_stop(&self, enabled: bool) {
|
||||
self.auto_stop.store(enabled, Ordering::SeqCst);
|
||||
self.inner.auto_stop.store(enabled, Ordering::SeqCst);
|
||||
}
|
||||
}
|
||||
|
||||
/// State for OGG-FLAC stream subscription.
|
||||
enum OggFlacStreamState {
|
||||
SendingHeader,
|
||||
Streaming,
|
||||
}
|
||||
|
||||
/// OGG-FLAC client stream (implements AsyncRead).
|
||||
pub struct OggFlacClientStream {
|
||||
rx: timed_broadcast::Receiver<Bytes>,
|
||||
buffer: VecDeque<u8>,
|
||||
finished: bool,
|
||||
handle: OggFlacStreamHandle,
|
||||
state: OggFlacStreamState,
|
||||
current_epoch: u64,
|
||||
inner: SharedClientStream,
|
||||
}
|
||||
|
||||
impl OggFlacClientStream {
|
||||
/// Dernier epoch observé (incrémenté à chaque TopZeroSync).
|
||||
fn new(rx: timed_broadcast::Receiver<Bytes>, handle: Arc<SharedStreamHandleInner>) -> Self {
|
||||
Self {
|
||||
inner: SharedClientStream::new(rx, handle),
|
||||
}
|
||||
}
|
||||
|
||||
pub fn current_epoch(&self) -> u64 {
|
||||
self.current_epoch
|
||||
self.inner.current_epoch()
|
||||
}
|
||||
}
|
||||
|
||||
@@ -164,272 +134,19 @@ impl AsyncRead for OggFlacClientStream {
|
||||
cx: &mut Context<'_>,
|
||||
buf: &mut ReadBuf<'_>,
|
||||
) -> Poll<io::Result<()>> {
|
||||
loop {
|
||||
// If in header state, send the header first
|
||||
if matches!(self.state, OggFlacStreamState::SendingHeader) {
|
||||
let header_opt = if let Ok(guard) = self.handle.header.try_read() {
|
||||
guard.clone()
|
||||
} else {
|
||||
None
|
||||
};
|
||||
|
||||
if let Some(header) = header_opt {
|
||||
self.buffer.extend(header.iter());
|
||||
debug!(
|
||||
"Sending cached OGG-FLAC header to new client ({} bytes)",
|
||||
header.len()
|
||||
);
|
||||
self.state = OggFlacStreamState::Streaming;
|
||||
continue; // Now copy header to output buffer
|
||||
} else {
|
||||
// Header not yet captured, skip to streaming
|
||||
self.state = OggFlacStreamState::Streaming;
|
||||
}
|
||||
}
|
||||
|
||||
// If we have buffered data, copy it
|
||||
if !self.buffer.is_empty() {
|
||||
let to_copy = self.buffer.len().min(buf.remaining());
|
||||
if to_copy == 0 {
|
||||
return Poll::Ready(Ok(()));
|
||||
}
|
||||
|
||||
let slice = self.buffer.make_contiguous();
|
||||
buf.put_slice(&slice[..to_copy]);
|
||||
self.buffer.drain(..to_copy);
|
||||
return Poll::Ready(Ok(()));
|
||||
}
|
||||
|
||||
if self.finished {
|
||||
return Poll::Ready(Ok(()));
|
||||
}
|
||||
|
||||
// Try to receive more data
|
||||
match self.rx.try_recv() {
|
||||
Ok(packet) => {
|
||||
self.current_epoch = packet.epoch;
|
||||
self.buffer.extend(packet.payload.iter());
|
||||
}
|
||||
Err(TryRecvError::Empty) => {
|
||||
// No data available right now.
|
||||
// Schedule a wakeup after a small delay to avoid busy-loop polling.
|
||||
let waker = cx.waker().clone();
|
||||
tokio::spawn(async move {
|
||||
tokio::time::sleep(tokio::time::Duration::from_micros(100)).await;
|
||||
waker.wake();
|
||||
});
|
||||
return Poll::Pending;
|
||||
}
|
||||
Err(TryRecvError::Lagged(skipped)) => {
|
||||
warn!("OGG-FLAC client lagged, skipped {} messages", skipped);
|
||||
}
|
||||
Err(TryRecvError::Closed) => {
|
||||
self.finished = true;
|
||||
return Poll::Ready(Ok(()));
|
||||
}
|
||||
}
|
||||
}
|
||||
Pin::new(&mut self.inner).poll_read(cx, buf)
|
||||
}
|
||||
}
|
||||
|
||||
impl Drop for OggFlacClientStream {
|
||||
fn drop(&mut self) {
|
||||
let count = self.handle.active_clients.fetch_sub(1, Ordering::SeqCst);
|
||||
debug!("OGG-FLAC client disconnected (remaining: {})", count - 1);
|
||||
|
||||
if count == 1 {
|
||||
if self.handle.auto_stop.load(Ordering::SeqCst) {
|
||||
debug!("Last OGG-FLAC client disconnected, signaling pipeline stop");
|
||||
self.handle.stop_token.cancel();
|
||||
} else {
|
||||
debug!("Last OGG-FLAC client disconnected, keeping pipeline alive");
|
||||
}
|
||||
}
|
||||
let remaining = self.inner.handle().client_disconnected();
|
||||
debug!("OGG-FLAC client disconnected (remaining: {})", remaining);
|
||||
}
|
||||
}
|
||||
|
||||
/// Internal state for encoder initialization.
|
||||
struct EncoderState {
|
||||
broadcaster_task: tokio::task::JoinHandle<()>,
|
||||
}
|
||||
|
||||
/// Logic for the streaming OGG-FLAC sink.
|
||||
struct StreamingOggFlacSinkLogic {
|
||||
encoder_options: EncoderOptions,
|
||||
bits_per_sample: u8,
|
||||
pcm_tx: Option<mpsc::Sender<PcmChunk>>,
|
||||
pcm_rx: Option<mpsc::Receiver<PcmChunk>>,
|
||||
metadata: Arc<RwLock<MetadataSnapshot>>,
|
||||
broadcast: timed_broadcast::Sender<Bytes>,
|
||||
header: Arc<RwLock<Option<Bytes>>>,
|
||||
encoder_state: Option<EncoderState>,
|
||||
sample_rate: Option<u32>,
|
||||
broadcast_max_lead_time: f64,
|
||||
/// Accumulated timestamp offset for maintaining continuity across encoder restarts
|
||||
timestamp_offset_sec: f64,
|
||||
/// Shared timestamp reference to read the last broadcast timestamp
|
||||
current_timestamp: Arc<RwLock<f64>>,
|
||||
}
|
||||
|
||||
impl StreamingOggFlacSinkLogic {
|
||||
/// Initialize the FLAC encoder once we know the sample rate.
|
||||
async fn initialize_encoder(&mut self, sample_rate: u32, timestamp_offset_sec: f64) -> Result<(), AudioError> {
|
||||
if self.encoder_state.is_some() {
|
||||
return Ok(()); // Already initialized
|
||||
}
|
||||
|
||||
debug!(
|
||||
"Initializing OGG-FLAC encoder with sample rate: {} Hz",
|
||||
sample_rate
|
||||
);
|
||||
|
||||
// Take the PCM receiver (we only initialize once)
|
||||
let pcm_rx = self
|
||||
.pcm_rx
|
||||
.take()
|
||||
.ok_or_else(|| AudioError::ProcessingError("PCM receiver already consumed".into()))?;
|
||||
|
||||
// Create shared timestamp and duration for pacing
|
||||
// Reuse existing current_timestamp Arc to maintain reference for reading later
|
||||
let current_timestamp = self.current_timestamp.clone();
|
||||
let current_duration = Arc::new(RwLock::new(0.0f64));
|
||||
|
||||
// Create ByteStreamReader for the encoder
|
||||
let pcm_reader =
|
||||
ByteStreamReader::new(pcm_rx, current_timestamp.clone(), current_duration.clone());
|
||||
|
||||
// Create PCM format
|
||||
let pcm_format = PcmFormat {
|
||||
sample_rate,
|
||||
channels: 2,
|
||||
bits_per_sample: self.bits_per_sample,
|
||||
};
|
||||
|
||||
// Start the FLAC encoder
|
||||
let flac_stream = encode_flac_stream(pcm_reader, pcm_format, self.encoder_options.clone())
|
||||
.await
|
||||
.map_err(|e| {
|
||||
AudioError::ProcessingError(format!("Failed to start FLAC encoder: {}", e))
|
||||
})?;
|
||||
|
||||
debug!("OGG-FLAC encoder initialized successfully");
|
||||
|
||||
// Spawn OGG wrapper + broadcaster task with timestamp and duration for pacing
|
||||
let broadcast = self.broadcast.clone();
|
||||
let header = self.header.clone();
|
||||
let max_lead = self.broadcast_max_lead_time;
|
||||
let broadcaster_task = tokio::spawn(async move {
|
||||
if let Err(e) = broadcast_ogg_flac_stream(
|
||||
flac_stream,
|
||||
broadcast,
|
||||
header,
|
||||
current_timestamp,
|
||||
current_duration,
|
||||
max_lead,
|
||||
timestamp_offset_sec,
|
||||
)
|
||||
.await
|
||||
{
|
||||
error!("OGG broadcaster task error: {}", e);
|
||||
}
|
||||
});
|
||||
|
||||
self.encoder_state = Some(EncoderState { broadcaster_task });
|
||||
|
||||
debug!("OGG broadcaster task spawned");
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Update metadata from a TrackBoundary marker.
|
||||
async fn update_metadata(
|
||||
&mut self,
|
||||
metadata_lock: &Arc<RwLock<dyn TrackMetadata>>,
|
||||
timestamp_sec: f64,
|
||||
) -> Result<(), AudioError> {
|
||||
let metadata = metadata_lock.read().await;
|
||||
|
||||
let mut snapshot = self.metadata.write().await;
|
||||
|
||||
// Extract all metadata fields
|
||||
snapshot.title = metadata.get_title().await.ok().flatten();
|
||||
snapshot.artist = metadata.get_artist().await.ok().flatten();
|
||||
snapshot.album = metadata.get_album().await.ok().flatten();
|
||||
snapshot.duration = metadata.get_duration().await.ok().flatten();
|
||||
snapshot.cover_url = metadata.get_cover_url().await.ok().flatten();
|
||||
snapshot.cover_pk = metadata.get_cover_pk().await.ok().flatten();
|
||||
snapshot.year = metadata.get_year().await.ok().flatten();
|
||||
|
||||
// Extract extra fields
|
||||
if let Ok(Some(extra)) = metadata.get_extra().await {
|
||||
snapshot.genre = extra.get("genre").cloned();
|
||||
snapshot.track_number = extra
|
||||
.get("track_number")
|
||||
.and_then(|s| s.parse::<u32>().ok());
|
||||
}
|
||||
|
||||
snapshot.audio_timestamp_sec = timestamp_sec;
|
||||
snapshot.version += 1;
|
||||
|
||||
debug!(
|
||||
"Metadata updated: v{} @ {:.2}s - {} - {} (cover_pk: {:?})",
|
||||
snapshot.version,
|
||||
timestamp_sec,
|
||||
snapshot.artist.as_deref().unwrap_or("?"),
|
||||
snapshot.title.as_deref().unwrap_or("?"),
|
||||
snapshot.cover_pk
|
||||
);
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Restart the FLAC encoder for a new track.
|
||||
/// This causes a new OGG stream header to be emitted and timestamps to reset to 0.
|
||||
async fn restart_encoder_for_new_track(&mut self) -> Result<(), AudioError> {
|
||||
let sample_rate = self
|
||||
.sample_rate
|
||||
.ok_or_else(|| AudioError::ProcessingError("Sample rate not initialized".into()))?;
|
||||
|
||||
debug!("Restarting OGG-FLAC encoder for new track");
|
||||
|
||||
// 1. Read current timestamp to maintain continuity
|
||||
let last_timestamp = *self.current_timestamp.read().await;
|
||||
debug!("Last timestamp before restart: {:.3}s", last_timestamp);
|
||||
|
||||
// 2. Close current PCM sender to signal encoder to finish
|
||||
if let Some(tx) = self.pcm_tx.take() {
|
||||
drop(tx);
|
||||
trace!("Dropped PCM sender to signal encoder finish");
|
||||
}
|
||||
|
||||
// 3. Wait for current broadcaster task to complete
|
||||
if let Some(state) = self.encoder_state.take() {
|
||||
trace!("Waiting for OGG broadcaster task to finish...");
|
||||
match state.broadcaster_task.await {
|
||||
Ok(_) => {
|
||||
trace!("OGG broadcaster task finished successfully");
|
||||
}
|
||||
Err(e) => {
|
||||
warn!("OGG broadcaster task error during restart: {:?}", e);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// 4. Update timestamp offset to maintain continuity across encoder restart
|
||||
self.timestamp_offset_sec += last_timestamp;
|
||||
debug!("New timestamp offset: {:.3}s", self.timestamp_offset_sec);
|
||||
|
||||
// 5. Create new PCM channel
|
||||
let (pcm_tx, pcm_rx) = mpsc::channel::<PcmChunk>(16);
|
||||
self.pcm_tx = Some(pcm_tx);
|
||||
self.pcm_rx = Some(pcm_rx);
|
||||
|
||||
// 6. Initialize new encoder with timestamp offset (will naturally emit new OGG stream header)
|
||||
self.initialize_encoder(sample_rate, self.timestamp_offset_sec).await?;
|
||||
|
||||
debug!("OGG-FLAC encoder restarted successfully for new track");
|
||||
Ok(())
|
||||
}
|
||||
ctx: SharedSinkContext,
|
||||
}
|
||||
|
||||
#[async_trait]
|
||||
@@ -462,20 +179,43 @@ impl NodeLogic for StreamingOggFlacSinkLogic {
|
||||
match &seg.segment {
|
||||
_AudioSegment::Chunk(chunk) => {
|
||||
// Detect sample rate from first chunk and initialize encoder
|
||||
if self.sample_rate.is_none() {
|
||||
if self.ctx.sample_rate.is_none() {
|
||||
let sample_rate = chunk.sample_rate();
|
||||
self.sample_rate = Some(sample_rate);
|
||||
self.ctx.sample_rate = Some(sample_rate);
|
||||
debug!("Detected sample rate: {} Hz", sample_rate);
|
||||
|
||||
// Initialize the FLAC encoder now (first track starts at 0.0)
|
||||
self.initialize_encoder(sample_rate, 0.0).await?;
|
||||
self.ctx
|
||||
.initialize_encoder(
|
||||
sample_rate,
|
||||
0.0,
|
||||
|flac_stream,
|
||||
broadcast,
|
||||
header,
|
||||
current_timestamp,
|
||||
current_duration,
|
||||
max_lead,
|
||||
sample_rate,
|
||||
timestamp_offset_sec| {
|
||||
broadcast_ogg_flac_stream(
|
||||
flac_stream,
|
||||
broadcast,
|
||||
header,
|
||||
current_timestamp,
|
||||
current_duration,
|
||||
max_lead,
|
||||
timestamp_offset_sec,
|
||||
)
|
||||
},
|
||||
)
|
||||
.await?;
|
||||
}
|
||||
|
||||
// Convert chunk to PCM bytes
|
||||
let pcm_bytes = chunk_to_pcm_bytes(&chunk, self.bits_per_sample)?;
|
||||
let pcm_bytes = chunk_to_pcm_bytes(&chunk, self.ctx.bits_per_sample)?;
|
||||
|
||||
// Calculate exact duration from samples and sample rate
|
||||
let sample_rate = self.sample_rate.expect("sample_rate should be initialized");
|
||||
let sample_rate = self.ctx.sample_rate.expect("sample_rate should be initialized");
|
||||
let duration_sec = chunk.len() as f64 / sample_rate as f64;
|
||||
|
||||
trace!(
|
||||
@@ -494,7 +234,7 @@ impl NodeLogic for StreamingOggFlacSinkLogic {
|
||||
};
|
||||
|
||||
// Get the sender (it should always be Some after initialization)
|
||||
let pcm_tx = match &self.pcm_tx {
|
||||
let pcm_tx = match &self.ctx.pcm_tx {
|
||||
Some(tx) => tx,
|
||||
None => {
|
||||
error!("OGG PCM sender not initialized");
|
||||
@@ -512,9 +252,32 @@ impl NodeLogic for StreamingOggFlacSinkLogic {
|
||||
match marker.as_ref() {
|
||||
SyncMarker::TrackBoundary { metadata } => {
|
||||
// Only restart encoder if it's already initialized (not the first track)
|
||||
if self.sample_rate.is_some() && self.encoder_state.is_some() {
|
||||
if self.ctx.sample_rate.is_some() && self.ctx.encoder_state.is_some() {
|
||||
// Restart encoder to emit new OGG stream header and reset timestamps
|
||||
if let Err(e) = self.restart_encoder_for_new_track().await {
|
||||
if let Err(e) = self
|
||||
.ctx
|
||||
.restart_encoder_for_new_track(
|
||||
|flac_stream,
|
||||
broadcast,
|
||||
header,
|
||||
current_timestamp,
|
||||
current_duration,
|
||||
max_lead,
|
||||
sample_rate,
|
||||
timestamp_offset_sec| {
|
||||
broadcast_ogg_flac_stream(
|
||||
flac_stream,
|
||||
broadcast,
|
||||
header,
|
||||
current_timestamp,
|
||||
current_duration,
|
||||
max_lead,
|
||||
timestamp_offset_sec,
|
||||
)
|
||||
},
|
||||
)
|
||||
.await
|
||||
{
|
||||
error!("Failed to restart OGG encoder for new track: {}", e);
|
||||
break;
|
||||
}
|
||||
@@ -523,7 +286,7 @@ impl NodeLogic for StreamingOggFlacSinkLogic {
|
||||
}
|
||||
|
||||
// Update metadata for the new track
|
||||
if let Err(e) = self.update_metadata(metadata, seg.timestamp_sec).await {
|
||||
if let Err(e) = self.ctx.update_metadata(metadata, seg.timestamp_sec).await {
|
||||
error!("Failed to update metadata: {}", e);
|
||||
}
|
||||
}
|
||||
@@ -610,40 +373,43 @@ impl StreamingOggFlacSink {
|
||||
let broadcast_capacity = calculate_broadcast_capacity(broadcast_max_lead_time);
|
||||
debug!(
|
||||
"Streaming Sink: using broadcast capacity of {} items (max_lead_time={:.1}s)",
|
||||
broadcast_capacity,
|
||||
broadcast_max_lead_time
|
||||
broadcast_capacity, broadcast_max_lead_time
|
||||
);
|
||||
let (broadcast, _) = timed_broadcast::channel(broadcast_capacity);
|
||||
|
||||
// OGG-FLAC header cache
|
||||
let header = Arc::new(RwLock::new(None));
|
||||
|
||||
// Stop token and client counter
|
||||
// Stop token and client control
|
||||
let stop_token = CancellationToken::new();
|
||||
let active_clients = Arc::new(AtomicUsize::new(0));
|
||||
let auto_stop = Arc::new(AtomicBool::new(true));
|
||||
|
||||
let handle = OggFlacStreamHandle {
|
||||
broadcast: broadcast.clone(),
|
||||
metadata: metadata.clone(),
|
||||
active_clients,
|
||||
stop_token: stop_token.clone(),
|
||||
header: header.clone(),
|
||||
auto_stop: Arc::new(AtomicBool::new(true)),
|
||||
};
|
||||
let shared_handle = Arc::new(SharedStreamHandleInner::new(
|
||||
broadcast.clone(),
|
||||
metadata.clone(),
|
||||
stop_token.clone(),
|
||||
header.clone(),
|
||||
auto_stop.clone(),
|
||||
));
|
||||
|
||||
let handle = OggFlacStreamHandle::new(shared_handle.clone());
|
||||
|
||||
let logic = StreamingOggFlacSinkLogic {
|
||||
encoder_options,
|
||||
bits_per_sample,
|
||||
pcm_tx: Some(pcm_tx),
|
||||
pcm_rx: Some(pcm_rx),
|
||||
metadata,
|
||||
broadcast,
|
||||
header,
|
||||
encoder_state: None,
|
||||
sample_rate: None,
|
||||
broadcast_max_lead_time: broadcast_max_lead_time.max(0.0),
|
||||
timestamp_offset_sec: 0.0,
|
||||
current_timestamp: Arc::new(RwLock::new(0.0)),
|
||||
ctx: SharedSinkContext {
|
||||
encoder_options,
|
||||
bits_per_sample,
|
||||
pcm_tx: Some(pcm_tx),
|
||||
pcm_rx: Some(pcm_rx),
|
||||
metadata,
|
||||
broadcast,
|
||||
header,
|
||||
encoder_state: None,
|
||||
sample_rate: None,
|
||||
broadcast_max_lead_time: broadcast_max_lead_time.max(0.0),
|
||||
first_chunk_timestamp_checked: false,
|
||||
timestamp_offset_sec: 0.0,
|
||||
current_timestamp: Arc::new(RwLock::new(0.0)),
|
||||
},
|
||||
};
|
||||
|
||||
let sink = Self {
|
||||
@@ -683,8 +449,6 @@ impl TypedAudioNode for StreamingOggFlacSink {
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
|
||||
/// OGG wrapper + broadcaster task: reads FLAC bytes from encoder, wraps in OGG pages, and broadcasts.
|
||||
/// Implements precise real-time pacing based on audio timestamps.
|
||||
/// Ensures FLAC frames are only sent at frame boundaries to prevent sync errors in strict decoders like FFPlay.
|
||||
@@ -828,10 +592,7 @@ async fn broadcast_ogg_flac_stream(
|
||||
trace!("Sent empty EOS page");
|
||||
}
|
||||
|
||||
trace!(
|
||||
"OGG-FLAC stream ended, total OGG bytes: {}",
|
||||
total_bytes
|
||||
);
|
||||
trace!("OGG-FLAC stream ended, total OGG bytes: {}", total_bytes);
|
||||
break;
|
||||
}
|
||||
Ok(n) => {
|
||||
@@ -906,8 +667,7 @@ async fn broadcast_ogg_flac_stream(
|
||||
}
|
||||
|
||||
// Extract just the first frame
|
||||
let first_frame: Vec<u8> =
|
||||
accumulator.drain(0..second_frame_start).collect();
|
||||
let first_frame: Vec<u8> = accumulator.drain(0..second_frame_start).collect();
|
||||
|
||||
// ╔═══════════════════════════════════════════════════════════════╗
|
||||
// ║ BACKPRESSURE INTELLIGENTE BASÉE SUR LE TIMING ║
|
||||
@@ -934,7 +694,8 @@ async fn broadcast_ogg_flac_stream(
|
||||
// Calculer le timestamp de cette FLAC frame (avec offset pour continuité entre tracks)
|
||||
let frame_start_samples = encoded_samples;
|
||||
encoded_samples = encoded_samples.saturating_add(first_frame_samples as u64);
|
||||
let audio_timestamp = timestamp_offset_sec + (frame_start_samples as f64 / sample_rate_f64);
|
||||
let audio_timestamp =
|
||||
timestamp_offset_sec + (frame_start_samples as f64 / sample_rate_f64);
|
||||
let segment_duration = first_frame_samples as f64 / sample_rate_f64;
|
||||
|
||||
// Check timing et apply pacing (skip si en retard)
|
||||
@@ -1023,7 +784,6 @@ async fn broadcast_ogg_flac_stream(
|
||||
Ok(())
|
||||
}
|
||||
|
||||
|
||||
/// Create OGG-FLAC identification packet (first packet in BOS page)
|
||||
/// Format: https://xiph.org/flac/ogg_mapping.html
|
||||
fn create_ogg_flac_identification(flac_header: &[u8]) -> Result<Vec<u8>, AudioError> {
|
||||
|
||||
394
pmoaudio-ext/src/sinks/streaming_sink_common.rs
Normal file
394
pmoaudio-ext/src/sinks/streaming_sink_common.rs
Normal file
@@ -0,0 +1,394 @@
|
||||
use std::collections::VecDeque;
|
||||
use std::future::Future;
|
||||
use std::io;
|
||||
use std::pin::Pin;
|
||||
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
|
||||
use std::sync::Arc;
|
||||
use std::task::{Context, Poll};
|
||||
use std::time::Duration;
|
||||
|
||||
use bytes::Bytes;
|
||||
use pmoaudio::AudioError;
|
||||
use pmoflac::{encode_flac_stream, EncoderOptions, FlacEncodedStream, PcmFormat};
|
||||
use pmometadata::TrackMetadata;
|
||||
use serde::{Deserialize, Serialize};
|
||||
use tokio::io::{AsyncRead, ReadBuf};
|
||||
use tokio::sync::{mpsc, RwLock};
|
||||
use tokio::task::JoinHandle;
|
||||
use tokio_util::sync::CancellationToken;
|
||||
use tracing::{debug, error, trace, warn};
|
||||
|
||||
use crate::byte_stream_reader::{ByteStreamReader, PcmChunk};
|
||||
use crate::sinks::timed_broadcast::{self, TryRecvError};
|
||||
|
||||
/// Snapshot of track metadata shared across streaming sinks.
|
||||
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
|
||||
pub struct MetadataSnapshot {
|
||||
pub title: Option<String>,
|
||||
pub artist: Option<String>,
|
||||
pub album: Option<String>,
|
||||
#[serde(skip_serializing_if = "Option::is_none")]
|
||||
pub duration: Option<Duration>,
|
||||
#[serde(skip_serializing_if = "Option::is_none")]
|
||||
pub cover_url: Option<String>,
|
||||
#[serde(skip_serializing_if = "Option::is_none")]
|
||||
pub cover_pk: Option<String>,
|
||||
#[serde(skip_serializing_if = "Option::is_none")]
|
||||
pub track_number: Option<u32>,
|
||||
#[serde(skip_serializing_if = "Option::is_none")]
|
||||
pub album_artist: Option<String>,
|
||||
#[serde(skip_serializing_if = "Option::is_none")]
|
||||
pub genre: Option<String>,
|
||||
#[serde(skip_serializing_if = "Option::is_none")]
|
||||
pub year: Option<u32>,
|
||||
pub audio_timestamp_sec: f64,
|
||||
pub version: u64,
|
||||
}
|
||||
|
||||
/// Shared handle state for streaming sinks.
|
||||
pub struct SharedStreamHandleInner {
|
||||
pub broadcast: timed_broadcast::Sender<Bytes>,
|
||||
pub metadata: Arc<RwLock<MetadataSnapshot>>,
|
||||
pub active_clients: Arc<AtomicUsize>,
|
||||
pub stop_token: CancellationToken,
|
||||
pub header: Arc<RwLock<Option<Bytes>>>,
|
||||
pub auto_stop: Arc<AtomicBool>,
|
||||
}
|
||||
|
||||
impl SharedStreamHandleInner {
|
||||
pub fn new(
|
||||
broadcast: timed_broadcast::Sender<Bytes>,
|
||||
metadata: Arc<RwLock<MetadataSnapshot>>,
|
||||
stop_token: CancellationToken,
|
||||
header: Arc<RwLock<Option<Bytes>>>,
|
||||
auto_stop: Arc<AtomicBool>,
|
||||
) -> Self {
|
||||
Self {
|
||||
broadcast,
|
||||
metadata,
|
||||
active_clients: Arc::new(AtomicUsize::new(0)),
|
||||
stop_token,
|
||||
header,
|
||||
auto_stop,
|
||||
}
|
||||
}
|
||||
|
||||
pub fn register_client(&self) -> timed_broadcast::Receiver<Bytes> {
|
||||
self.broadcast.subscribe()
|
||||
}
|
||||
|
||||
pub fn client_connected(&self) -> usize {
|
||||
self.active_clients.fetch_add(1, Ordering::SeqCst) + 1
|
||||
}
|
||||
|
||||
pub fn client_disconnected(&self) -> usize {
|
||||
let prev = self.active_clients.fetch_sub(1, Ordering::SeqCst);
|
||||
let remaining = prev.saturating_sub(1);
|
||||
if prev == 1 && self.auto_stop.load(Ordering::SeqCst) {
|
||||
trace!("Last client disconnected, signaling pipeline stop (shared handle)");
|
||||
self.stop_token.cancel();
|
||||
}
|
||||
remaining
|
||||
}
|
||||
}
|
||||
|
||||
enum StreamState {
|
||||
SendingHeader,
|
||||
Streaming,
|
||||
}
|
||||
|
||||
pub struct SharedClientStream {
|
||||
rx: timed_broadcast::Receiver<Bytes>,
|
||||
buffer: VecDeque<u8>,
|
||||
finished: bool,
|
||||
handle: Arc<SharedStreamHandleInner>,
|
||||
state: StreamState,
|
||||
current_epoch: u64,
|
||||
}
|
||||
|
||||
impl SharedClientStream {
|
||||
pub fn new(rx: timed_broadcast::Receiver<Bytes>, handle: Arc<SharedStreamHandleInner>) -> Self {
|
||||
Self {
|
||||
rx,
|
||||
buffer: VecDeque::new(),
|
||||
finished: false,
|
||||
handle,
|
||||
state: StreamState::SendingHeader,
|
||||
current_epoch: 0,
|
||||
}
|
||||
}
|
||||
|
||||
pub fn current_epoch(&self) -> u64 {
|
||||
self.current_epoch
|
||||
}
|
||||
|
||||
pub fn handle(&self) -> &Arc<SharedStreamHandleInner> {
|
||||
&self.handle
|
||||
}
|
||||
}
|
||||
|
||||
impl AsyncRead for SharedClientStream {
|
||||
fn poll_read(
|
||||
mut self: Pin<&mut Self>,
|
||||
cx: &mut Context<'_>,
|
||||
buf: &mut ReadBuf<'_>,
|
||||
) -> Poll<io::Result<()>> {
|
||||
loop {
|
||||
if matches!(self.state, StreamState::SendingHeader) {
|
||||
let header_opt = if let Ok(guard) = self.handle.header.try_read() {
|
||||
guard.clone()
|
||||
} else {
|
||||
None
|
||||
};
|
||||
|
||||
if let Some(header) = header_opt {
|
||||
self.buffer.extend(header.iter());
|
||||
trace!(
|
||||
"Sending cached header to new client ({} bytes)",
|
||||
header.len()
|
||||
);
|
||||
self.state = StreamState::Streaming;
|
||||
continue;
|
||||
} else {
|
||||
self.state = StreamState::Streaming;
|
||||
}
|
||||
}
|
||||
|
||||
if !self.buffer.is_empty() {
|
||||
let to_copy = self.buffer.len().min(buf.remaining());
|
||||
if to_copy == 0 {
|
||||
return Poll::Ready(Ok(()));
|
||||
}
|
||||
|
||||
let slice = self.buffer.make_contiguous();
|
||||
buf.put_slice(&slice[..to_copy]);
|
||||
self.buffer.drain(..to_copy);
|
||||
return Poll::Ready(Ok(()));
|
||||
}
|
||||
|
||||
if self.finished {
|
||||
return Poll::Ready(Ok(()));
|
||||
}
|
||||
|
||||
match self.rx.try_recv() {
|
||||
Ok(packet) => {
|
||||
self.current_epoch = packet.epoch;
|
||||
self.buffer.extend(packet.payload.iter());
|
||||
}
|
||||
Err(TryRecvError::Empty) => {
|
||||
let waker = cx.waker().clone();
|
||||
tokio::spawn(async move {
|
||||
tokio::time::sleep(Duration::from_millis(10)).await;
|
||||
waker.wake();
|
||||
});
|
||||
return Poll::Pending;
|
||||
}
|
||||
Err(TryRecvError::Lagged(skipped)) => {
|
||||
warn!("Client lagged, skipped {} messages", skipped);
|
||||
}
|
||||
Err(TryRecvError::Closed) => {
|
||||
self.finished = true;
|
||||
return Poll::Ready(Ok(()));
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
pub struct EncoderState {
|
||||
pub broadcaster_task: JoinHandle<()>,
|
||||
}
|
||||
|
||||
pub struct SharedSinkContext {
|
||||
pub encoder_options: EncoderOptions,
|
||||
pub bits_per_sample: u8,
|
||||
pub pcm_tx: Option<mpsc::Sender<PcmChunk>>,
|
||||
pub pcm_rx: Option<mpsc::Receiver<PcmChunk>>,
|
||||
pub metadata: Arc<RwLock<MetadataSnapshot>>,
|
||||
pub broadcast: timed_broadcast::Sender<Bytes>,
|
||||
pub header: Arc<RwLock<Option<Bytes>>>,
|
||||
pub encoder_state: Option<EncoderState>,
|
||||
pub sample_rate: Option<u32>,
|
||||
pub broadcast_max_lead_time: f64,
|
||||
pub first_chunk_timestamp_checked: bool,
|
||||
pub timestamp_offset_sec: f64,
|
||||
pub current_timestamp: Arc<RwLock<f64>>,
|
||||
}
|
||||
|
||||
impl SharedSinkContext {
|
||||
pub async fn initialize_encoder<Fut, F>(
|
||||
&mut self,
|
||||
sample_rate: u32,
|
||||
timestamp_offset_sec: f64,
|
||||
broadcaster: F,
|
||||
) -> Result<(), AudioError>
|
||||
where
|
||||
F: FnOnce(
|
||||
FlacEncodedStream,
|
||||
timed_broadcast::Sender<Bytes>,
|
||||
Arc<RwLock<Option<Bytes>>>,
|
||||
Arc<RwLock<f64>>,
|
||||
Arc<RwLock<f64>>,
|
||||
f64,
|
||||
u32,
|
||||
f64,
|
||||
) -> Fut
|
||||
+ Send
|
||||
+ 'static,
|
||||
Fut: Future<Output = Result<(), AudioError>> + Send + 'static,
|
||||
{
|
||||
if self.encoder_state.is_some() {
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
debug!(
|
||||
"Initializing FLAC encoder with sample rate: {} Hz",
|
||||
sample_rate
|
||||
);
|
||||
|
||||
let pcm_rx = self
|
||||
.pcm_rx
|
||||
.take()
|
||||
.ok_or_else(|| AudioError::ProcessingError("PCM receiver already consumed".into()))?;
|
||||
|
||||
let current_timestamp = self.current_timestamp.clone();
|
||||
let current_duration = Arc::new(RwLock::new(0.0f64));
|
||||
|
||||
let pcm_reader =
|
||||
ByteStreamReader::new(pcm_rx, current_timestamp.clone(), current_duration.clone());
|
||||
|
||||
let pcm_format = PcmFormat {
|
||||
sample_rate,
|
||||
channels: 2,
|
||||
bits_per_sample: self.bits_per_sample,
|
||||
};
|
||||
|
||||
let flac_stream = encode_flac_stream(pcm_reader, pcm_format, self.encoder_options.clone())
|
||||
.await
|
||||
.map_err(|e| {
|
||||
AudioError::ProcessingError(format!("Failed to start FLAC encoder: {}", e))
|
||||
})?;
|
||||
|
||||
debug!("FLAC encoder initialized successfully");
|
||||
|
||||
let broadcast = self.broadcast.clone();
|
||||
let header = self.header.clone();
|
||||
let max_lead = self.broadcast_max_lead_time;
|
||||
let current_timestamp_clone = current_timestamp.clone();
|
||||
let current_duration_clone = current_duration.clone();
|
||||
|
||||
let broadcaster_task = tokio::spawn(async move {
|
||||
if let Err(e) = broadcaster(
|
||||
flac_stream,
|
||||
broadcast,
|
||||
header,
|
||||
current_timestamp_clone,
|
||||
current_duration_clone,
|
||||
max_lead,
|
||||
sample_rate,
|
||||
timestamp_offset_sec,
|
||||
)
|
||||
.await
|
||||
{
|
||||
error!("Broadcaster task error: {}", e);
|
||||
}
|
||||
});
|
||||
|
||||
self.encoder_state = Some(EncoderState { broadcaster_task });
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub async fn restart_encoder_for_new_track<Fut, F>(
|
||||
&mut self,
|
||||
broadcaster: F,
|
||||
) -> Result<(), AudioError>
|
||||
where
|
||||
F: FnOnce(
|
||||
FlacEncodedStream,
|
||||
timed_broadcast::Sender<Bytes>,
|
||||
Arc<RwLock<Option<Bytes>>>,
|
||||
Arc<RwLock<f64>>,
|
||||
Arc<RwLock<f64>>,
|
||||
f64,
|
||||
u32,
|
||||
f64,
|
||||
) -> Fut
|
||||
+ Send
|
||||
+ 'static,
|
||||
Fut: Future<Output = Result<(), AudioError>> + Send + 'static,
|
||||
{
|
||||
let sample_rate = self
|
||||
.sample_rate
|
||||
.ok_or_else(|| AudioError::ProcessingError("Sample rate not initialized".into()))?;
|
||||
|
||||
debug!("Restarting FLAC encoder for new track");
|
||||
|
||||
let last_timestamp = *self.current_timestamp.read().await;
|
||||
debug!("Last timestamp before restart: {:.3}s", last_timestamp);
|
||||
|
||||
if let Some(tx) = self.pcm_tx.take() {
|
||||
drop(tx);
|
||||
trace!("Dropped PCM sender to signal encoder finish");
|
||||
}
|
||||
|
||||
if let Some(state) = self.encoder_state.take() {
|
||||
trace!("Waiting for broadcaster task to finish...");
|
||||
match state.broadcaster_task.await {
|
||||
Ok(_) => trace!("Broadcaster task finished successfully"),
|
||||
Err(e) => warn!("Broadcaster task error during restart: {:?}", e),
|
||||
}
|
||||
}
|
||||
|
||||
self.timestamp_offset_sec += last_timestamp;
|
||||
debug!("New timestamp offset: {:.3}s", self.timestamp_offset_sec);
|
||||
|
||||
let (pcm_tx, pcm_rx) = mpsc::channel::<PcmChunk>(16);
|
||||
self.pcm_tx = Some(pcm_tx);
|
||||
self.pcm_rx = Some(pcm_rx);
|
||||
|
||||
self.initialize_encoder(sample_rate, self.timestamp_offset_sec, broadcaster)
|
||||
.await?;
|
||||
|
||||
debug!("FLAC encoder restarted successfully for new track");
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub async fn update_metadata(
|
||||
&mut self,
|
||||
metadata_lock: &Arc<RwLock<dyn TrackMetadata>>,
|
||||
timestamp_sec: f64,
|
||||
) -> Result<(), AudioError> {
|
||||
let metadata = metadata_lock.read().await;
|
||||
let mut snapshot = self.metadata.write().await;
|
||||
|
||||
snapshot.title = metadata.get_title().await.ok().flatten();
|
||||
snapshot.artist = metadata.get_artist().await.ok().flatten();
|
||||
snapshot.album = metadata.get_album().await.ok().flatten();
|
||||
snapshot.duration = metadata.get_duration().await.ok().flatten();
|
||||
snapshot.cover_url = metadata.get_cover_url().await.ok().flatten();
|
||||
snapshot.cover_pk = metadata.get_cover_pk().await.ok().flatten();
|
||||
snapshot.year = metadata.get_year().await.ok().flatten();
|
||||
|
||||
if let Ok(Some(extra)) = metadata.get_extra().await {
|
||||
snapshot.genre = extra.get("genre").cloned();
|
||||
snapshot.track_number = extra
|
||||
.get("track_number")
|
||||
.and_then(|s| s.parse::<u32>().ok());
|
||||
}
|
||||
|
||||
snapshot.audio_timestamp_sec = timestamp_sec;
|
||||
snapshot.version += 1;
|
||||
|
||||
debug!(
|
||||
"Metadata updated: v{} @ {:.2}s - {} - {} (cover_pk: {:?})",
|
||||
snapshot.version,
|
||||
timestamp_sec,
|
||||
snapshot.artist.as_deref().unwrap_or("?"),
|
||||
snapshot.title.as_deref().unwrap_or("?"),
|
||||
snapshot.cover_pk
|
||||
);
|
||||
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
@@ -81,12 +81,12 @@ impl<'a> std::ops::DerefMut for ConnGuard<'a> {
|
||||
|
||||
impl DB {
|
||||
fn lock_conn(&self, ctx: &'static str) -> ConnGuard<'_> {
|
||||
trace!("DB mutex → acquiring ({ctx})");
|
||||
// trace!("DB mutex → acquiring ({ctx})");
|
||||
let start = Instant::now();
|
||||
let guard = self.conn.lock().unwrap();
|
||||
let waited = start.elapsed();
|
||||
|
||||
trace!("DB mutex → acquired ({ctx}) in {:?}", waited);
|
||||
// trace!("DB mutex → acquired ({ctx}) in {:?}", waited);
|
||||
if waited > std::time::Duration::from_millis(50) {
|
||||
warn!("DB mutex wait >50 ms ({}): {:?}", ctx, waited);
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user