use kithara_bufpool::HasPool;
use kithara_play::{
InterruptionKind, PlayError, SeekOutcome, SelectionPlayback, SessionDuckingMode,
};
use smallvec::SmallVec;
use super::{
QueueControl,
types::{CachedPosition, PendingSelect, PlaybackView, SelectPhase, Transition},
};
use crate::{
attempts::LoadClass,
error::QueueError,
event::{AdvanceReason, TrackStatus},
};
impl<S> QueueControl<S>
where
S: HasPool<u8> + HasPool<f32> + Send + Sync + 'static,
{
fn freeze_cached_position(&self) {
if let Some(t) = self.player.position_seconds() {
self.write_cached_position(CachedPosition::known(t));
}
}
pub(super) fn is_paused(&self) -> bool {
self.player.is_paused()
}
fn maybe_arm_crossfade(&self) {
if self.is_paused() {
return;
}
let crossfade = self.player.crossfade_duration();
let view = self.playback_view();
let (Some(dur), Some(pos), Some(entry)) = (view.duration, view.position, self.current())
else {
return;
};
let armed_for = self.read_armed_for();
let time = super::types::PlaybackTime { dur, pos };
if !super::types::should_arm_crossfade(time, crossfade, entry.id, armed_for) {
return;
}
let transition = if crossfade > 0.0 {
Transition::Crossfade
} else {
Transition::None
};
self.advance_loaded_successor(entry.id, transition);
}
pub fn notify_audio_route_changed(&self, reason: &str) -> Result<(), QueueError> {
self.with_open_result(|queue| queue.player.invalidate_audio_route(reason))?;
Ok(())
}
pub fn notify_interruption(&self, kind: InterruptionKind) {
self.command(|queue| queue.player.notify_interruption(kind));
}
pub fn pause(&self) {
self.command(|queue| {
queue.player.pause();
let mut phase = queue.lock_pending_select_mut();
if let SelectPhase::Pending(mut pending) = *phase {
pending.playback = SelectionPlayback::Pause;
*phase = SelectPhase::Pending(pending);
}
drop(phase);
queue.freeze_cached_position();
});
}
pub fn play(&self) {
self.command(Self::play_inner);
}
fn play_inner(&self) {
let mut phase = self.lock_pending_select_mut();
if let SelectPhase::Pending(mut pending) = *phase {
pending.playback = SelectionPlayback::Play;
*phase = SelectPhase::Pending(pending);
}
drop(phase);
self.player.play();
let _apply = self.lock_select_apply();
let pending = match *self.lock_pending_select_mut() {
SelectPhase::Pending(pending) => Some(pending),
SelectPhase::Idle => None,
};
let index = self.player.current_index();
if pending.is_none() && self.player.item_has_resource(index) {
return;
}
let current = {
let guard = self.lock_tracks();
pending
.map_or_else(
|| guard.get(index),
|pending| guard.iter().find(|entry| entry.id == pending.id),
)
.map(|entry| (entry.id, entry.status.clone()))
};
let Some((id, status)) = current else {
return;
};
match status {
TrackStatus::Loaded => self.set_status(id, TrackStatus::Consumed),
TrackStatus::Pending | TrackStatus::Loading | TrackStatus::Slow => {
self.override_pending_select(pending.unwrap_or_else(|| PendingSelect {
id,
settings: Transition::None.settings(self.crossfade_settings()),
playback: SelectionPlayback::Play,
reason: AdvanceReason::UserSelect,
}));
self.promote_pending_load(id);
}
TrackStatus::Failed(_) => {
let Some(source) = self.tracks.source(id) else {
return;
};
self.override_pending_select(PendingSelect {
id,
settings: Transition::None.settings(self.crossfade_settings()),
playback: SelectionPlayback::Play,
reason: AdvanceReason::UserSelect,
});
self.set_status(id, TrackStatus::Pending);
self.spawn_apply_after_load(id, source, LoadClass::Interactive);
}
TrackStatus::Consumed | TrackStatus::Cancelled => {}
}
}
#[must_use]
pub fn playback_view(&self) -> PlaybackView {
let mut view = self
.player
.playback_snapshot()
.map(PlaybackView::from)
.unwrap_or_default();
view.position = self.position_seconds();
view
}
pub(super) fn seek_player(&self, seconds: f64) -> Result<SeekOutcome, PlayError> {
self.with_open_result(|queue| queue.seek_player_inner(seconds))
}
fn seek_player_inner(&self, seconds: f64) -> Result<SeekOutcome, PlayError> {
if self.current().is_none() {
let id = { self.lock_navigation().last_selected() };
if let Some(id) = id {
let ids = self
.tracks()
.into_iter()
.map(|track| track.id)
.collect::<SmallVec<[_; 16]>>();
self.lock_navigation_mut().select(id, &ids);
self.handle_current_item_changed();
}
}
let outcome = self.player.seek_seconds(seconds)?;
if let SeekOutcome::Landed { landed_at, .. } = outcome {
self.write_cached_position(CachedPosition::known(landed_at.as_secs_f64()));
}
Ok(outcome)
}
pub fn set_session_ducking(&self, mode: SessionDuckingMode) -> Result<(), QueueError> {
self.with_open_result(|queue| queue.player.set_session_ducking(mode))?;
Ok(())
}
pub(super) fn tick_player(&self) -> Result<(), PlayError> {
self.with_open_result(Self::tick_player_inner)
}
fn tick_player_inner(&self) -> Result<(), PlayError> {
self.player.tick()?;
self.player.process_notifications();
self.drain_player_events();
self.update_cached_position();
self.maybe_arm_crossfade();
Ok(())
}
fn update_cached_position(&self) {
const MIN_STABLE_POSITION_SECS: f64 = 0.5;
if self.is_paused() {
return;
}
let Some(t) = self.player.position_seconds() else {
return;
};
let prev = Option::<f64>::from(self.read_cached_position());
if t == 0.0 && prev.is_some_and(|p| p > MIN_STABLE_POSITION_SECS) {
return;
}
self.write_cached_position(CachedPosition::known(t));
}
delegate::delegate! {
to self {
#[must_use]
#[into]
#[call(read_cached_position)]
pub fn position_seconds(&self) -> Option<f64>;
#[expr($.map_err(QueueError::from))]
#[call(seek_player)]
pub fn seek(&self, seconds: f64) -> Result<SeekOutcome, QueueError>;
#[expr($.map_err(QueueError::from))]
#[call(tick_player)]
pub fn tick(&self) -> Result<(), QueueError>;
}
}
}
#[cfg(test)]
mod tests {
use kithara_events::{SlotId, TrackId};
use kithara_platform::sync::Arc;
use kithara_play::{ItemRole, PlayerEvent, TrackRef};
use kithara_test_utils::kithara;
use crate::{
event::{QueueEvent, TrackStatus},
queue::{
state::tests::make_queue,
types::{CrossfadeArm, PlaybackTime, SelectPhase, should_arm_crossfade},
},
track::{TrackRecord, TrackSource},
};
#[kithara::test(tokio)]
async fn spurious_item_did_play_to_end_is_filtered() {
let queue = make_queue();
let _a = queue.append("https://example.com/a.mp3");
let _b = queue.append("https://example.com/b.mp3");
queue.player.bus().publish(PlayerEvent::ItemDidPlayToEnd {
item: ItemRole::Leading(TrackRef::new(
TrackId::allocate(),
SlotId::new(0),
Arc::from(""),
)),
});
queue
.tick()
.expect("BUG: tick returned error in test setup");
assert_eq!(
queue.lock_navigation().current(),
None,
"navigation must not have advanced"
);
}
#[kithara::test(tokio)]
async fn eof_after_queue_end_does_not_restart_from_first_track() {
let queue = make_queue();
let a = TrackId::allocate();
let b = TrackId::allocate();
queue.tracks.lock().extend([
TrackRecord::new(a, "a".into(), TrackSource::from("a")),
TrackRecord::new(b, "b".into(), TrackSource::from("b")),
]);
queue.lock_navigation_mut().select(b, &[a, b]);
queue.lock_navigation_mut().finish();
let mut rx = queue.subscribe();
queue.player.bus().publish(PlayerEvent::ItemDidPlayToEnd {
item: ItemRole::Leading(TrackRef::new(
b,
SlotId::new(0),
Arc::from(format!("test://memory/{}", b.as_u64())),
)),
});
queue
.tick()
.expect("BUG: tick returned error in test setup");
assert_eq!(
queue.lock_navigation().current(),
None,
"stale EOF must not restart the queue"
);
let saw_ended = crate::queue::state::tests::wait_for_queue_event(
&mut rx,
|ev| matches!(ev, QueueEvent::QueueEnded),
200,
)
.await;
assert!(!saw_ended, "stale EOF must not duplicate QueueEnded");
}
#[kithara::test(tokio)]
async fn play_retries_the_current_track_after_its_prefetch_failed() {
let queue = make_queue();
let id = queue
.append("https://example.com/a.mp3")
.expect("open queue accepts a track");
queue.set_status(id, TrackStatus::Failed("network offline".into()));
queue.play_inner();
let SelectPhase::Pending(pending) = *queue.lock_pending_select_mut() else {
panic!("play must retain selection while retrying the failed track")
};
assert_eq!(pending.id, id);
assert_eq!(pending.playback, kithara_play::SelectionPlayback::Play);
}
#[kithara::test(tokio)]
#[case::append(false)]
#[case::insert(true)]
async fn play_promotes_the_initial_pending_prefetch(#[case] insert: bool) {
let queue = make_queue();
let id = if insert {
queue.insert("https://example.com/a.mp3", None)
} else {
queue.append("https://example.com/a.mp3")
}
.expect("open queue accepts a track");
assert!(!queue.tracks.attempt_selected(id));
queue.play();
assert!(queue.tracks.attempt_selected(id));
}
#[kithara::test]
#[case::remaining_equals_crossfade(157.0, 162.0, 5.0, TrackId(1), CrossfadeArm::Disarmed, true)]
#[case::remaining_below_crossfade(160.0, 162.0, 5.0, TrackId(1), CrossfadeArm::Disarmed, true)]
#[case::far_from_end(100.0, 162.0, 5.0, TrackId(1), CrossfadeArm::Disarmed, false)]
#[case::already_armed_for_same_track(
160.0,
162.0,
5.0,
TrackId(1),
CrossfadeArm::armed(TrackId(1)),
false
)]
#[case::armed_for_different_track_still_arms(
160.0,
162.0,
5.0,
TrackId(1),
CrossfadeArm::armed(TrackId(0)),
true
)]
#[case::crossfade_zero_at_tail_no_pre_arm(
161.9,
162.0,
0.0,
TrackId(1),
CrossfadeArm::Disarmed,
false
)]
#[case::crossfade_zero_quiet_middle(
161.0,
162.0,
0.0,
TrackId(1),
CrossfadeArm::Disarmed,
false
)]
#[case::zero_position_rejected(0.0, 162.0, 5.0, TrackId(1), CrossfadeArm::Disarmed, false)]
#[case::zero_duration_rejected(10.0, 0.0, 5.0, TrackId(1), CrossfadeArm::Disarmed, false)]
fn should_arm_crossfade_cases(
#[case] pos: f64,
#[case] dur: f64,
#[case] crossfade: f32,
#[case] current_id: TrackId,
#[case] armed_for: CrossfadeArm,
#[case] expected: bool,
) {
assert_eq!(
should_arm_crossfade(PlaybackTime { dur, pos }, crossfade, current_id, armed_for),
expected
);
}
}