feat: playlist management, gapless playback, ReplayGain, Qobuz theme
Playlist management: - Add/remove tracks from playlists via right-click context menu - Create new playlists (right-click Playlists sidebar header) - Delete playlists with confirmation dialog (right-click playlist item) - Playlist view removes track immediately on delete (optimistic) - Deleting currently-open playlist clears the track view Gapless playback: - Single long-running audio thread owns AudioOutput; CPAL stream stays open between tracks eliminating device teardown/startup gap - Decode runs inline on the audio thread; command channel polled via try_recv() so Pause/Resume/Seek/Stop/Play all work without spawning - New Play command arriving mid-decode is handled immediately, reusing the same audio output for zero-gap transition - Position timer reduced from 500 ms to 50 ms for faster track-end detection - URL/metadata prefetch: when gapless is enabled Qt pre-fetches the next track while the current one is still playing ReplayGain: - Toggled in Settings → Playback - replaygain_track_gain (dB) from track audio_info converted to linear gain factor and applied per-sample alongside volume Qobuz dark theme: - Background #191919, base #141414, accent #FFB232 (yellow-orange) - Selection highlight, slider fill, scrollbar hover all use #FFB232 - Links use Qobuz blue #46B3EE - Hi-res H badges updated to #FFB232 (from #FFD700) - Now-playing row uses #FFB232 (was Spotify green) - QSS stylesheet for scrollbars, menus, inputs, buttons, groups Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
This commit is contained in:
@@ -30,7 +30,9 @@ enum QobuzEvent {
|
||||
EV_POSITION = 16,
|
||||
EV_TRACK_URL_OK = 17,
|
||||
EV_TRACK_URL_ERR = 18,
|
||||
EV_GENERIC_ERR = 19,
|
||||
EV_GENERIC_ERR = 19,
|
||||
EV_PLAYLIST_CREATED = 20,
|
||||
EV_PLAYLIST_DELETED = 21,
|
||||
};
|
||||
|
||||
// Callback signature
|
||||
@@ -69,6 +71,16 @@ uint8_t qobuz_backend_get_volume(const QobuzBackendOpaque *backend);
|
||||
int qobuz_backend_get_state(const QobuzBackendOpaque *backend);
|
||||
int qobuz_backend_take_track_finished(QobuzBackendOpaque *backend);
|
||||
|
||||
// ReplayGain / Gapless
|
||||
void qobuz_backend_set_replaygain(QobuzBackendOpaque *backend, bool enabled);
|
||||
void qobuz_backend_prefetch_track(QobuzBackendOpaque *backend, int64_t track_id, int32_t format_id);
|
||||
|
||||
// Playlist management
|
||||
void qobuz_backend_create_playlist(QobuzBackendOpaque *backend, const char *name);
|
||||
void qobuz_backend_delete_playlist(QobuzBackendOpaque *backend, int64_t playlist_id);
|
||||
void qobuz_backend_add_track_to_playlist(QobuzBackendOpaque *backend, int64_t playlist_id, int64_t track_id);
|
||||
void qobuz_backend_delete_track_from_playlist(QobuzBackendOpaque *backend, int64_t playlist_id, int64_t playlist_track_id);
|
||||
|
||||
// Favorites modification
|
||||
void qobuz_backend_add_fav_track(QobuzBackendOpaque *backend, int64_t track_id);
|
||||
void qobuz_backend_remove_fav_track(QobuzBackendOpaque *backend, int64_t track_id);
|
||||
|
||||
@@ -89,6 +89,15 @@ impl QobuzClient {
|
||||
Ok(body)
|
||||
}
|
||||
|
||||
fn post_request(&self, method: &str) -> reqwest::RequestBuilder {
|
||||
let mut builder = self.http.post(self.url(method));
|
||||
builder = builder.query(&[("app_id", self.app_id.as_str())]);
|
||||
if let Some(token) = &self.auth_token {
|
||||
builder = builder.header("Authorization", format!("Bearer {}", token));
|
||||
}
|
||||
builder
|
||||
}
|
||||
|
||||
fn get_request(&self, method: &str) -> reqwest::RequestBuilder {
|
||||
let mut builder = self.http.get(self.url(method));
|
||||
builder = builder.query(&[("app_id", self.app_id.as_str())]);
|
||||
@@ -329,6 +338,55 @@ impl QobuzClient {
|
||||
Ok(serde_json::from_value(body["artists"].clone())?)
|
||||
}
|
||||
|
||||
// --- Playlist management ---
|
||||
|
||||
pub async fn create_playlist(&self, name: &str) -> Result<PlaylistDto> {
|
||||
let resp = self
|
||||
.post_request("playlist/create")
|
||||
.form(&[("name", name), ("is_public", "false"), ("is_collaborative", "false")])
|
||||
.send()
|
||||
.await?;
|
||||
let body = Self::check_response(resp).await?;
|
||||
Ok(serde_json::from_value(body)?)
|
||||
}
|
||||
|
||||
pub async fn add_track_to_playlist(&self, playlist_id: i64, track_id: i64) -> Result<()> {
|
||||
let resp = self
|
||||
.post_request("playlist/addTracks")
|
||||
.form(&[
|
||||
("playlist_id", playlist_id.to_string()),
|
||||
("track_ids", track_id.to_string()),
|
||||
("no_duplicate", "true".to_string()),
|
||||
])
|
||||
.send()
|
||||
.await?;
|
||||
Self::check_response(resp).await?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub async fn delete_playlist(&self, playlist_id: i64) -> Result<()> {
|
||||
let resp = self
|
||||
.get_request("playlist/delete")
|
||||
.query(&[("playlist_id", &playlist_id.to_string())])
|
||||
.send()
|
||||
.await?;
|
||||
Self::check_response(resp).await?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub async fn delete_track_from_playlist(&self, playlist_id: i64, playlist_track_id: i64) -> Result<()> {
|
||||
let resp = self
|
||||
.post_request("playlist/deleteTracks")
|
||||
.form(&[
|
||||
("playlist_id", playlist_id.to_string()),
|
||||
("playlist_track_ids", playlist_track_id.to_string()),
|
||||
])
|
||||
.send()
|
||||
.await?;
|
||||
Self::check_response(resp).await?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub async fn add_fav_track(&self, track_id: i64) -> Result<()> {
|
||||
let resp = self
|
||||
.get_request("favorite/create")
|
||||
|
||||
@@ -48,6 +48,7 @@ pub struct TrackDto {
|
||||
pub title: Option<String>,
|
||||
pub duration: Option<i64>,
|
||||
pub track_number: Option<i32>,
|
||||
pub playlist_track_id: Option<i64>,
|
||||
pub album: Option<AlbumDto>,
|
||||
pub performer: Option<ArtistDto>,
|
||||
pub composer: Option<ArtistDto>,
|
||||
|
||||
161
rust/src/lib.rs
161
rust/src/lib.rs
@@ -75,12 +75,20 @@ pub type EventCallback = unsafe extern "C" fn(*mut c_void, c_int, *const c_char)
|
||||
|
||||
// ---------- Backend ----------
|
||||
|
||||
struct PrefetchedTrack {
|
||||
track_id: i64,
|
||||
track: api::models::TrackDto,
|
||||
url: String,
|
||||
}
|
||||
|
||||
struct BackendInner {
|
||||
client: Arc<Mutex<QobuzClient>>,
|
||||
player: Player,
|
||||
rt: Runtime,
|
||||
cb: EventCallback,
|
||||
ud: SendPtr,
|
||||
replaygain_enabled: std::sync::Arc<std::sync::atomic::AtomicBool>,
|
||||
prefetch: std::sync::Arc<tokio::sync::Mutex<Option<PrefetchedTrack>>>,
|
||||
}
|
||||
|
||||
pub struct Backend(BackendInner);
|
||||
@@ -121,6 +129,8 @@ pub unsafe extern "C" fn qobuz_backend_new(
|
||||
rt,
|
||||
cb: event_cb,
|
||||
ud: SendPtr(userdata),
|
||||
replaygain_enabled: std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false)),
|
||||
prefetch: std::sync::Arc::new(tokio::sync::Mutex::new(None)),
|
||||
})))
|
||||
}
|
||||
|
||||
@@ -341,35 +351,58 @@ pub unsafe extern "C" fn qobuz_backend_play_track(
|
||||
let format = Format::from_id(format_id);
|
||||
let cmd_tx = inner.player.cmd_tx.clone();
|
||||
let status = inner.player.status.clone();
|
||||
let prefetch = inner.prefetch.clone();
|
||||
let rg_enabled = inner.replaygain_enabled.clone();
|
||||
|
||||
spawn(inner, async move {
|
||||
// 1. Track metadata
|
||||
let track = match client.lock().await.get_track(track_id).await {
|
||||
Ok(t) => t,
|
||||
Err(e) => { call_cb(cb, ud, EV_TRACK_URL_ERR, &err_json(&e.to_string())); return; }
|
||||
// 1. Check prefetch cache first for zero-gap start
|
||||
let cached = {
|
||||
let mut lock = prefetch.lock().await;
|
||||
if lock.as_ref().map(|p| p.track_id == track_id).unwrap_or(false) {
|
||||
lock.take()
|
||||
} else {
|
||||
None
|
||||
}
|
||||
};
|
||||
|
||||
// 2. Stream URL
|
||||
let url_dto = match client.lock().await.get_track_url(track_id, format).await {
|
||||
Ok(u) => u,
|
||||
Err(e) => { call_cb(cb, ud, EV_TRACK_URL_ERR, &err_json(&e.to_string())); return; }
|
||||
};
|
||||
let url = match url_dto.url {
|
||||
Some(u) => u,
|
||||
None => { call_cb(cb, ud, EV_TRACK_URL_ERR, &err_json("no stream URL")); return; }
|
||||
let (track, url) = if let Some(pf) = cached {
|
||||
(pf.track, pf.url)
|
||||
} else {
|
||||
// Fetch track metadata
|
||||
let track = match client.lock().await.get_track(track_id).await {
|
||||
Ok(t) => t,
|
||||
Err(e) => { call_cb(cb, ud, EV_TRACK_URL_ERR, &err_json(&e.to_string())); return; }
|
||||
};
|
||||
// Fetch stream URL
|
||||
let url_dto = match client.lock().await.get_track_url(track_id, format).await {
|
||||
Ok(u) => u,
|
||||
Err(e) => { call_cb(cb, ud, EV_TRACK_URL_ERR, &err_json(&e.to_string())); return; }
|
||||
};
|
||||
let url = match url_dto.url {
|
||||
Some(u) => u,
|
||||
None => { call_cb(cb, ud, EV_TRACK_URL_ERR, &err_json("no stream URL")); return; }
|
||||
};
|
||||
(track, url)
|
||||
};
|
||||
|
||||
// 3. Notify track change
|
||||
// 2. Notify track change
|
||||
if let Ok(j) = serde_json::to_string(&track) {
|
||||
call_cb(cb, ud, EV_TRACK_CHANGED, &j);
|
||||
}
|
||||
|
||||
// 3. Compute ReplayGain if enabled
|
||||
let replaygain_db = if rg_enabled.load(std::sync::atomic::Ordering::Relaxed) {
|
||||
track.audio_info.as_ref().and_then(|ai| ai.replaygain_track_gain)
|
||||
} else {
|
||||
None
|
||||
};
|
||||
|
||||
// 4. Update status + send play command
|
||||
*status.current_track.lock().unwrap() = Some(track.clone());
|
||||
if let Some(dur) = track.duration {
|
||||
status.duration_secs.store(dur as u64, std::sync::atomic::Ordering::Relaxed);
|
||||
}
|
||||
let _ = cmd_tx.send(player::PlayerCommand::Play(player::TrackInfo { track, url, format }));
|
||||
let _ = cmd_tx.send(player::PlayerCommand::Play(player::TrackInfo { track, url, format, replaygain_db }));
|
||||
|
||||
// 5. State notification
|
||||
call_cb(cb, ud, EV_STATE_CHANGED, r#"{"state":"playing"}"#);
|
||||
@@ -438,6 +471,41 @@ pub unsafe extern "C" fn qobuz_backend_take_track_finished(ptr: *mut Backend) ->
|
||||
if finished { 1 } else { 0 }
|
||||
}
|
||||
|
||||
// ---------- ReplayGain / Gapless ----------
|
||||
|
||||
#[no_mangle]
|
||||
pub unsafe extern "C" fn qobuz_backend_set_replaygain(ptr: *mut Backend, enabled: bool) {
|
||||
(*ptr).0.replaygain_enabled.store(enabled, std::sync::atomic::Ordering::Relaxed);
|
||||
}
|
||||
|
||||
#[no_mangle]
|
||||
pub unsafe extern "C" fn qobuz_backend_prefetch_track(
|
||||
ptr: *mut Backend,
|
||||
track_id: i64,
|
||||
format_id: i32,
|
||||
) {
|
||||
let inner = &(*ptr).0;
|
||||
let client = inner.client.clone();
|
||||
let prefetch = inner.prefetch.clone();
|
||||
let format = Format::from_id(format_id);
|
||||
|
||||
spawn(inner, async move {
|
||||
let track = match client.lock().await.get_track(track_id).await {
|
||||
Ok(t) => t,
|
||||
Err(_) => return,
|
||||
};
|
||||
let url_dto = match client.lock().await.get_track_url(track_id, format).await {
|
||||
Ok(u) => u,
|
||||
Err(_) => return,
|
||||
};
|
||||
let url = match url_dto.url {
|
||||
Some(u) => u,
|
||||
None => return,
|
||||
};
|
||||
*prefetch.lock().await = Some(PrefetchedTrack { track_id, track, url });
|
||||
});
|
||||
}
|
||||
|
||||
// ---------- Favorites modification ----------
|
||||
|
||||
#[no_mangle]
|
||||
@@ -489,3 +557,68 @@ pub unsafe extern "C" fn qobuz_backend_remove_fav_album(ptr: *mut Backend, album
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
// ---------- Playlist management ----------
|
||||
|
||||
pub const EV_PLAYLIST_CREATED: c_int = 20;
|
||||
pub const EV_PLAYLIST_DELETED: c_int = 21;
|
||||
|
||||
#[no_mangle]
|
||||
pub unsafe extern "C" fn qobuz_backend_create_playlist(ptr: *mut Backend, name: *const c_char) {
|
||||
let inner = &(*ptr).0;
|
||||
let name = CStr::from_ptr(name).to_string_lossy().into_owned();
|
||||
let client = inner.client.clone();
|
||||
let cb = inner.cb; let ud = inner.ud;
|
||||
spawn(inner, async move {
|
||||
match client.lock().await.create_playlist(&name).await {
|
||||
Ok(p) => call_cb(cb, ud, EV_PLAYLIST_CREATED, &serde_json::to_string(&p).unwrap_or_default()),
|
||||
Err(e) => call_cb(cb, ud, EV_GENERIC_ERR, &err_json(&e.to_string())),
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
#[no_mangle]
|
||||
pub unsafe extern "C" fn qobuz_backend_delete_playlist(ptr: *mut Backend, playlist_id: i64) {
|
||||
let inner = &(*ptr).0;
|
||||
let client = inner.client.clone();
|
||||
let cb = inner.cb; let ud = inner.ud;
|
||||
spawn(inner, async move {
|
||||
match client.lock().await.delete_playlist(playlist_id).await {
|
||||
Ok(()) => call_cb(cb, ud, EV_PLAYLIST_DELETED,
|
||||
&serde_json::json!({"playlist_id": playlist_id}).to_string()),
|
||||
Err(e) => call_cb(cb, ud, EV_GENERIC_ERR, &err_json(&e.to_string())),
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
#[no_mangle]
|
||||
pub unsafe extern "C" fn qobuz_backend_add_track_to_playlist(
|
||||
ptr: *mut Backend,
|
||||
playlist_id: i64,
|
||||
track_id: i64,
|
||||
) {
|
||||
let inner = &(*ptr).0;
|
||||
let client = inner.client.clone();
|
||||
let cb = inner.cb; let ud = inner.ud;
|
||||
spawn(inner, async move {
|
||||
if let Err(e) = client.lock().await.add_track_to_playlist(playlist_id, track_id).await {
|
||||
call_cb(cb, ud, EV_GENERIC_ERR, &err_json(&e.to_string()));
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
#[no_mangle]
|
||||
pub unsafe extern "C" fn qobuz_backend_delete_track_from_playlist(
|
||||
ptr: *mut Backend,
|
||||
playlist_id: i64,
|
||||
playlist_track_id: i64,
|
||||
) {
|
||||
let inner = &(*ptr).0;
|
||||
let client = inner.client.clone();
|
||||
let cb = inner.cb; let ud = inner.ud;
|
||||
spawn(inner, async move {
|
||||
if let Err(e) = client.lock().await.delete_track_from_playlist(playlist_id, playlist_track_id).await {
|
||||
call_cb(cb, ud, EV_GENERIC_ERR, &err_json(&e.to_string()));
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
@@ -15,7 +15,7 @@ use symphonia::core::{
|
||||
units::Time,
|
||||
};
|
||||
|
||||
use crate::player::{output::AudioOutput, PlayerStatus};
|
||||
use super::{output::AudioOutput, PlayerCommand, PlayerStatus, TrackInfo};
|
||||
|
||||
/// First 512 KiB of stream kept in memory to support backward seeks during probing.
|
||||
const HEAD_SIZE: usize = 512 * 1024;
|
||||
@@ -127,13 +127,22 @@ impl MediaSource for HttpStreamSource {
|
||||
}
|
||||
}
|
||||
|
||||
/// Stream and decode audio from `url`. Runs on a dedicated OS thread.
|
||||
pub fn play_track(
|
||||
/// Decode and play `url` inline on the calling thread (the player loop).
|
||||
///
|
||||
/// `audio_output` is reused across calls if the sample rate and channel count match,
|
||||
/// keeping the CPAL stream open between tracks for gapless playback.
|
||||
///
|
||||
/// Returns:
|
||||
/// - `Ok(Some(TrackInfo))` — a new Play command arrived; start that track next.
|
||||
/// - `Ok(None)` — track finished naturally or was stopped.
|
||||
/// - `Err(_)` — unrecoverable playback error.
|
||||
pub fn play_track_inline(
|
||||
url: &str,
|
||||
status: &PlayerStatus,
|
||||
stop: &Arc<AtomicBool>,
|
||||
paused: &Arc<AtomicBool>,
|
||||
) -> Result<()> {
|
||||
audio_output: &mut Option<AudioOutput>,
|
||||
cmd_rx: &std::sync::mpsc::Receiver<PlayerCommand>,
|
||||
) -> Result<Option<TrackInfo>> {
|
||||
let response = reqwest::blocking::get(url)?;
|
||||
let content_length = response.content_length();
|
||||
let source = HttpStreamSource::new(response, content_length);
|
||||
@@ -160,19 +169,91 @@ pub fn play_track(
|
||||
.make(&track.codec_params, &DecoderOptions::default())
|
||||
.map_err(|e| anyhow::anyhow!("decoder init failed: {e}"))?;
|
||||
|
||||
let mut audio_output = AudioOutput::try_open(sample_rate, channels)?;
|
||||
|
||||
loop {
|
||||
if stop.load(Ordering::SeqCst) {
|
||||
break;
|
||||
// Reuse existing audio output if format matches; rebuild only on format change.
|
||||
if let Some(ao) = audio_output.as_ref() {
|
||||
if ao.sample_rate != sample_rate || ao.channels != channels {
|
||||
*audio_output = None; // will be recreated below
|
||||
}
|
||||
while paused.load(Ordering::SeqCst) {
|
||||
std::thread::sleep(std::time::Duration::from_millis(50));
|
||||
if stop.load(Ordering::SeqCst) {
|
||||
return Ok(());
|
||||
}
|
||||
if audio_output.is_none() {
|
||||
*audio_output = Some(AudioOutput::try_open(sample_rate, channels)?);
|
||||
}
|
||||
let ao = audio_output.as_mut().unwrap();
|
||||
|
||||
let mut stopped = false;
|
||||
let mut next_track: Option<TrackInfo> = None;
|
||||
|
||||
'decode: loop {
|
||||
// Non-blocking command check — handle Pause/Resume/Seek/Stop/Play
|
||||
loop {
|
||||
match cmd_rx.try_recv() {
|
||||
Ok(PlayerCommand::Pause) => {
|
||||
paused.store(true, Ordering::SeqCst);
|
||||
*status.state.lock().unwrap() = super::PlayerState::Paused;
|
||||
}
|
||||
Ok(PlayerCommand::Resume) => {
|
||||
paused.store(false, Ordering::SeqCst);
|
||||
*status.state.lock().unwrap() = super::PlayerState::Playing;
|
||||
}
|
||||
Ok(PlayerCommand::Seek(s)) => {
|
||||
status.seek_target_secs.store(s, Ordering::Relaxed);
|
||||
status.seek_requested.load(Ordering::SeqCst); // read-side fence
|
||||
status.seek_requested.store(true, Ordering::SeqCst);
|
||||
}
|
||||
Ok(PlayerCommand::SetVolume(v)) => {
|
||||
status.volume.store(v, Ordering::Relaxed);
|
||||
}
|
||||
Ok(PlayerCommand::Stop) => {
|
||||
paused.store(false, Ordering::SeqCst);
|
||||
*status.state.lock().unwrap() = super::PlayerState::Idle;
|
||||
*status.current_track.lock().unwrap() = None;
|
||||
status.position_secs.store(0, Ordering::Relaxed);
|
||||
status.duration_secs.store(0, Ordering::Relaxed);
|
||||
stopped = true;
|
||||
break 'decode;
|
||||
}
|
||||
Ok(PlayerCommand::Play(info)) => {
|
||||
// New track requested — stop current and return it
|
||||
next_track = Some(info);
|
||||
break 'decode;
|
||||
}
|
||||
Err(std::sync::mpsc::TryRecvError::Empty) => break,
|
||||
Err(std::sync::mpsc::TryRecvError::Disconnected) => {
|
||||
stopped = true;
|
||||
break 'decode;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Spin while paused, but keep checking for commands
|
||||
while paused.load(Ordering::SeqCst) {
|
||||
std::thread::sleep(std::time::Duration::from_millis(10));
|
||||
// Still check for Stop/Play while paused
|
||||
match cmd_rx.try_recv() {
|
||||
Ok(PlayerCommand::Resume) => {
|
||||
paused.store(false, Ordering::SeqCst);
|
||||
*status.state.lock().unwrap() = super::PlayerState::Playing;
|
||||
}
|
||||
Ok(PlayerCommand::Stop) => {
|
||||
paused.store(false, Ordering::SeqCst);
|
||||
stopped = true;
|
||||
break;
|
||||
}
|
||||
Ok(PlayerCommand::Play(info)) => {
|
||||
paused.store(false, Ordering::SeqCst);
|
||||
next_track = Some(info);
|
||||
break 'decode;
|
||||
}
|
||||
Ok(PlayerCommand::SetVolume(v)) => {
|
||||
status.volume.store(v, Ordering::Relaxed);
|
||||
}
|
||||
_ => {}
|
||||
}
|
||||
if stopped { break 'decode; }
|
||||
}
|
||||
if stopped { break; }
|
||||
|
||||
// Handle seek
|
||||
if status.seek_requested.load(Ordering::SeqCst) {
|
||||
status.seek_requested.store(false, Ordering::SeqCst);
|
||||
let target = status.seek_target_secs.load(Ordering::Relaxed);
|
||||
@@ -190,8 +271,10 @@ pub fn play_track(
|
||||
|
||||
let packet = match format.next_packet() {
|
||||
Ok(p) => p,
|
||||
Err(SymphoniaError::IoError(e)) if e.kind() == std::io::ErrorKind::UnexpectedEof => {
|
||||
break;
|
||||
Err(SymphoniaError::IoError(e))
|
||||
if e.kind() == std::io::ErrorKind::UnexpectedEof =>
|
||||
{
|
||||
break; // natural end of track
|
||||
}
|
||||
Err(SymphoniaError::ResetRequired) => {
|
||||
decoder.reset();
|
||||
@@ -205,13 +288,16 @@ pub fn play_track(
|
||||
}
|
||||
|
||||
if let Some(ts) = packet.ts().checked_div(sample_rate as u64) {
|
||||
status.position_secs.store(ts, std::sync::atomic::Ordering::Relaxed);
|
||||
status.position_secs.store(ts, Ordering::Relaxed);
|
||||
}
|
||||
|
||||
match decoder.decode(&packet) {
|
||||
Ok(decoded) => {
|
||||
let volume = status.volume.load(Ordering::Relaxed) as f32 / 100.0;
|
||||
audio_output.write(decoded, volume, stop)?;
|
||||
let rg = *status.replaygain_gain.lock().unwrap();
|
||||
// Use a stop flag tied to new-track-incoming so write doesn't block
|
||||
let dummy_stop = Arc::new(AtomicBool::new(false));
|
||||
ao.write(decoded, (volume * rg).min(1.0), &dummy_stop)?;
|
||||
}
|
||||
Err(SymphoniaError::IoError(_)) => break,
|
||||
Err(SymphoniaError::DecodeError(e)) => eprintln!("decode error: {e}"),
|
||||
@@ -219,5 +305,10 @@ pub fn play_track(
|
||||
}
|
||||
}
|
||||
|
||||
Ok(())
|
||||
if stopped {
|
||||
// On explicit stop, drop the audio output to silence immediately
|
||||
*audio_output = None;
|
||||
}
|
||||
|
||||
Ok(next_track)
|
||||
}
|
||||
|
||||
@@ -24,6 +24,8 @@ pub struct TrackInfo {
|
||||
pub track: TrackDto,
|
||||
pub url: String,
|
||||
pub format: Format,
|
||||
/// ReplayGain track gain in dB, if enabled and available.
|
||||
pub replaygain_db: Option<f64>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, PartialEq)]
|
||||
@@ -47,6 +49,8 @@ pub struct PlayerStatus {
|
||||
/// Set by the player loop when a seek command arrives; cleared by the decode thread.
|
||||
pub seek_requested: Arc<AtomicBool>,
|
||||
pub seek_target_secs: Arc<AtomicU64>,
|
||||
/// Linear gain factor to apply (1.0 = unity). Updated each time a new track starts.
|
||||
pub replaygain_gain: Arc<std::sync::Mutex<f32>>,
|
||||
}
|
||||
|
||||
impl PlayerStatus {
|
||||
@@ -60,6 +64,7 @@ impl PlayerStatus {
|
||||
track_finished: Arc::new(AtomicBool::new(false)),
|
||||
seek_requested: Arc::new(AtomicBool::new(false)),
|
||||
seek_target_secs: Arc::new(AtomicU64::new(0)),
|
||||
replaygain_gain: Arc::new(std::sync::Mutex::new(1.0)),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -106,10 +111,6 @@ impl Player {
|
||||
self.cmd_tx.send(cmd).ok();
|
||||
}
|
||||
|
||||
pub fn play_track(&self, track: TrackDto, url: String, format: Format) {
|
||||
self.send(PlayerCommand::Play(TrackInfo { track, url, format }));
|
||||
}
|
||||
|
||||
pub fn pause(&self) {
|
||||
self.send(PlayerCommand::Pause);
|
||||
}
|
||||
@@ -133,68 +134,77 @@ impl Player {
|
||||
}
|
||||
}
|
||||
|
||||
/// The player loop runs on a single dedicated OS thread.
|
||||
/// It owns the `AudioOutput` locally so there are no Send constraints.
|
||||
/// Decoding is performed inline; the command channel is polled via try_recv
|
||||
/// inside the decode loop to handle Pause/Resume/Seek/Stop/Play without
|
||||
/// tearng down and re-opening the audio device between tracks.
|
||||
fn player_loop(rx: std::sync::mpsc::Receiver<PlayerCommand>, status: PlayerStatus) {
|
||||
let mut stop_flag = Arc::new(AtomicBool::new(true));
|
||||
use std::sync::mpsc::RecvTimeoutError;
|
||||
|
||||
let mut audio_output: Option<output::AudioOutput> = None;
|
||||
let paused = Arc::new(AtomicBool::new(false));
|
||||
// pending_info holds a Play command that interrupted an ongoing decode
|
||||
let mut pending_info: Option<TrackInfo> = None;
|
||||
|
||||
loop {
|
||||
match rx.recv_timeout(Duration::from_millis(100)) {
|
||||
Ok(cmd) => match cmd {
|
||||
PlayerCommand::Play(info) => {
|
||||
stop_flag.store(true, Ordering::SeqCst);
|
||||
stop_flag = Arc::new(AtomicBool::new(false));
|
||||
paused.store(false, Ordering::SeqCst);
|
||||
|
||||
*status.state.lock().unwrap() = PlayerState::Playing;
|
||||
*status.current_track.lock().unwrap() = Some(info.track.clone());
|
||||
if let Some(dur) = info.track.duration {
|
||||
status.duration_secs.store(dur as u64, Ordering::Relaxed);
|
||||
'outer: loop {
|
||||
// Wait for a Play command (or use one that was interrupted)
|
||||
let info = if let Some(p) = pending_info.take() {
|
||||
p
|
||||
} else {
|
||||
loop {
|
||||
match rx.recv_timeout(Duration::from_millis(100)) {
|
||||
Ok(PlayerCommand::Play(info)) => break info,
|
||||
Ok(PlayerCommand::Stop) => {
|
||||
audio_output = None;
|
||||
paused.store(false, Ordering::SeqCst);
|
||||
*status.state.lock().unwrap() = PlayerState::Idle;
|
||||
*status.current_track.lock().unwrap() = None;
|
||||
status.position_secs.store(0, Ordering::Relaxed);
|
||||
status.duration_secs.store(0, Ordering::Relaxed);
|
||||
}
|
||||
status.position_secs.store(0, Ordering::Relaxed);
|
||||
Ok(PlayerCommand::SetVolume(v)) => {
|
||||
status.volume.store(v, Ordering::Relaxed);
|
||||
}
|
||||
Ok(PlayerCommand::Seek(s)) => {
|
||||
status.seek_target_secs.store(s, Ordering::Relaxed);
|
||||
status.seek_requested.store(true, Ordering::SeqCst);
|
||||
}
|
||||
Ok(_) => {} // Pause/Resume ignored when idle
|
||||
Err(RecvTimeoutError::Timeout) => {}
|
||||
Err(RecvTimeoutError::Disconnected) => break 'outer,
|
||||
}
|
||||
}
|
||||
};
|
||||
|
||||
let status_c = status.clone();
|
||||
let stop_c = stop_flag.clone();
|
||||
let paused_c = paused.clone();
|
||||
// Compute ReplayGain factor
|
||||
let rg_factor = info.replaygain_db
|
||||
.map(|db| 10f32.powf(db as f32 / 20.0))
|
||||
.unwrap_or(1.0);
|
||||
*status.replaygain_gain.lock().unwrap() = rg_factor;
|
||||
|
||||
std::thread::spawn(move || {
|
||||
match decoder::play_track(&info.url, &status_c, &stop_c, &paused_c) {
|
||||
Ok(()) => {
|
||||
if !stop_c.load(Ordering::SeqCst) {
|
||||
*status_c.state.lock().unwrap() = PlayerState::Idle;
|
||||
status_c.track_finished.store(true, Ordering::SeqCst);
|
||||
}
|
||||
}
|
||||
Err(e) => {
|
||||
eprintln!("playback error: {e}");
|
||||
*status_c.state.lock().unwrap() =
|
||||
PlayerState::Error(e.to_string());
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
PlayerCommand::Pause => {
|
||||
paused.store(true, Ordering::SeqCst);
|
||||
*status.state.lock().unwrap() = PlayerState::Paused;
|
||||
}
|
||||
PlayerCommand::Resume => {
|
||||
paused.store(false, Ordering::SeqCst);
|
||||
*status.state.lock().unwrap() = PlayerState::Playing;
|
||||
}
|
||||
PlayerCommand::Stop => {
|
||||
stop_flag.store(true, Ordering::SeqCst);
|
||||
*status.state.lock().unwrap() = PlayerState::Idle;
|
||||
*status.current_track.lock().unwrap() = None;
|
||||
status.position_secs.store(0, Ordering::Relaxed);
|
||||
status.duration_secs.store(0, Ordering::Relaxed);
|
||||
}
|
||||
PlayerCommand::SetVolume(_) => {}
|
||||
PlayerCommand::Seek(secs) => {
|
||||
status.seek_target_secs.store(secs, Ordering::Relaxed);
|
||||
status.seek_requested.store(true, Ordering::SeqCst);
|
||||
}
|
||||
},
|
||||
Err(std::sync::mpsc::RecvTimeoutError::Timeout) => {}
|
||||
Err(std::sync::mpsc::RecvTimeoutError::Disconnected) => break,
|
||||
*status.state.lock().unwrap() = PlayerState::Playing;
|
||||
*status.current_track.lock().unwrap() = Some(info.track.clone());
|
||||
if let Some(dur) = info.track.duration {
|
||||
status.duration_secs.store(dur as u64, Ordering::Relaxed);
|
||||
}
|
||||
status.position_secs.store(0, Ordering::Relaxed);
|
||||
paused.store(false, Ordering::SeqCst);
|
||||
|
||||
match decoder::play_track_inline(&info.url, &status, &paused, &mut audio_output, &rx) {
|
||||
Ok(Some(next_info)) => {
|
||||
// Interrupted by a new Play — loop immediately with reused audio output
|
||||
pending_info = Some(next_info);
|
||||
}
|
||||
Ok(None) => {
|
||||
// Track finished naturally
|
||||
*status.state.lock().unwrap() = PlayerState::Idle;
|
||||
status.track_finished.store(true, Ordering::SeqCst);
|
||||
}
|
||||
Err(e) => {
|
||||
eprintln!("playback error: {e}");
|
||||
*status.state.lock().unwrap() = PlayerState::Error(e.to_string());
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -15,6 +15,8 @@ const RING_BUFFER_SIZE: usize = 32 * 1024;
|
||||
pub struct AudioOutput {
|
||||
ring_buf_producer: rb::Producer<f32>,
|
||||
_stream: cpal::Stream,
|
||||
pub sample_rate: u32,
|
||||
pub channels: usize,
|
||||
}
|
||||
|
||||
impl AudioOutput {
|
||||
@@ -50,6 +52,8 @@ impl AudioOutput {
|
||||
Ok(Self {
|
||||
ring_buf_producer: producer,
|
||||
_stream: stream,
|
||||
sample_rate,
|
||||
channels,
|
||||
})
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user