use std::{
ops::Deref,
sync::{Mutex, PoisonError},
};
use kithara_assets::{AssetStore, StorageBackend};
use kithara_bufpool::HasPool;
use kithara_events::{EventBus, EventReceiver, TrackId};
use kithara_platform::{
CancelScope, CancelToken, sync::Arc, tokio::runtime::Handle as RuntimeHandle,
};
use kithara_play::{
CrossfadeSettings, PlayError, PlayerImpl,
player::{PlayerControl, PlayerControlSource},
};
use super::{
engine_events::PlayerBusEvent,
types::{AtomicCachedPosition, AtomicTrackId, CachedPosition, CrossfadeArm, SelectPhase},
};
use crate::{
config::QueueConfig,
loader::Loader,
navigation::{ActionAtItemEnd, NavigationState},
track::{TrackRecord, Tracks},
};
#[doc(hidden)]
pub struct QueueRuntime<S>
where
S: HasPool<u8> + HasPool<f32> + Send + Sync + 'static,
{
pub(super) loader: Arc<Loader<S>>,
pub(super) navigation: Arc<Mutex<NavigationState>>,
pub(super) pending_select: Arc<Mutex<SelectPhase>>,
pub(super) select_apply: Arc<Mutex<()>>,
pub(super) tracks: Arc<Tracks<S>>,
pub(super) cached_position: AtomicCachedPosition,
pub(super) autoplay_target: AtomicTrackId,
pub(super) crossfade_armed_for: AtomicTrackId,
pub(super) shutdown: CancelToken,
pub(super) bus: EventBus,
pub(super) action_at_item_end: Mutex<ActionAtItemEnd>,
pub(super) admission: Mutex<()>,
pub(super) crossfade_settings: Mutex<CrossfadeSettings>,
pub(super) player_rx: Mutex<EventReceiver<PlayerBusEvent>>,
pub(super) should_autoplay: bool,
}
#[derive_where::derive_where(Clone; S: HasPool<u8> + HasPool<f32> + Send + Sync + 'static)]
pub struct QueueControl<S>
where
S: HasPool<u8> + HasPool<f32> + Send + Sync + 'static,
{
pub(super) player: PlayerControl<S>,
runtime: Arc<QueueRuntime<S>>,
}
pub struct Queue<S>
where
S: HasPool<u8> + HasPool<f32> + Send + Sync + 'static,
{
pub(super) player: PlayerImpl<S>,
pub(super) control: QueueControl<S>,
}
impl<S> Deref for QueueControl<S>
where
S: HasPool<u8> + HasPool<f32> + Send + Sync + 'static,
{
type Target = QueueRuntime<S>;
fn deref(&self) -> &Self::Target {
&self.runtime
}
}
impl<S> Deref for Queue<S>
where
S: HasPool<u8> + HasPool<f32> + Send + Sync + 'static,
{
type Target = QueueControl<S>;
fn deref(&self) -> &Self::Target {
&self.control
}
}
impl<S> Queue<S>
where
S: HasPool<u8> + HasPool<f32> + Send + Sync + 'static,
{
#[must_use]
pub fn new(config: QueueConfig<S>) -> Self {
let QueueConfig {
player,
runtime,
store,
cancel: config_cancel,
max_concurrent_loads,
max_history_size,
prefetch_duration,
should_autoplay,
playback_order,
action_at_item_end,
crossfade_settings,
} = config;
let cancel = CancelScope::new(config_cancel).token();
let store = store.unwrap_or_else(|| {
AssetStore::builder(player.pools().clone())
.backend(StorageBackend::default())
.cancel(cancel.child())
.build()
});
player.set_auto_advance_enabled(false);
player.set_prefetch_duration(prefetch_duration);
player.set_crossfade_duration(crossfade_settings.duration);
let bus = player.bus().clone();
let player_control = player.control();
let tracks = Arc::new(Tracks::new(bus.clone()));
let loader = Arc::new(Loader::new(
player_control.clone(),
runtime.or_else(|| RuntimeHandle::try_current().ok()),
store,
max_concurrent_loads,
Arc::clone(&tracks),
cancel.child(),
));
let player_rx = player.subscribe();
let mut navigation = NavigationState::new(max_history_size);
navigation.set_playback_order(playback_order, &[]);
let runtime = Arc::new(QueueRuntime {
loader,
tracks,
bus,
should_autoplay,
admission: Mutex::new(()),
shutdown: cancel,
navigation: Arc::new(Mutex::new(navigation)),
action_at_item_end: Mutex::new(action_at_item_end),
crossfade_settings: Mutex::new(crossfade_settings),
pending_select: Arc::new(Mutex::new(SelectPhase::Idle)),
select_apply: Arc::new(Mutex::new(())),
player_rx: Mutex::new(player_rx),
crossfade_armed_for: AtomicTrackId::disarmed(),
autoplay_target: AtomicTrackId::disarmed(),
cached_position: AtomicCachedPosition::unknown(),
});
Self {
player,
control: QueueControl {
runtime,
player: player_control,
},
}
}
}
impl<S> QueueControl<S>
where
S: HasPool<u8> + HasPool<f32> + Send + Sync + 'static,
{
pub fn close(&self) -> Result<(), PlayError> {
let _admission = self.lock_admission();
self.player.close()?;
self.shutdown.cancel();
Ok(())
}
pub(in crate::queue) fn command(&self, operation: impl FnOnce(&Self)) {
let _ = self.with_open(operation);
}
fn ensure_open(&self) -> Result<(), PlayError> {
if self.is_closed() {
Err(PlayError::Closed)
} else {
Ok(())
}
}
pub(crate) fn invalidate(&self) {
self.shutdown.cancel();
}
#[must_use]
pub fn is_closed(&self) -> bool {
self.shutdown.is_cancelled() || self.player.is_closed()
}
pub(in crate::queue) fn lock_admission(&self) -> std::sync::MutexGuard<'_, ()> {
self.admission
.lock()
.unwrap_or_else(PoisonError::into_inner)
}
pub(super) fn lock_navigation(&self) -> std::sync::MutexGuard<'_, NavigationState> {
self.navigation
.lock()
.unwrap_or_else(PoisonError::into_inner)
}
pub(super) fn lock_navigation_mut(&self) -> std::sync::MutexGuard<'_, NavigationState> {
self.navigation
.lock()
.unwrap_or_else(PoisonError::into_inner)
}
pub(in crate::queue) fn lock_pending_select_mut(
&self,
) -> std::sync::MutexGuard<'_, SelectPhase> {
self.pending_select
.lock()
.unwrap_or_else(PoisonError::into_inner)
}
pub(in crate::queue) fn lock_select_apply(&self) -> std::sync::MutexGuard<'_, ()> {
self.select_apply
.lock()
.unwrap_or_else(PoisonError::into_inner)
}
pub(in crate::queue) fn with_open<T>(
&self,
operation: impl FnOnce(&Self) -> T,
) -> Result<T, PlayError> {
let _admission = self.lock_admission();
self.ensure_open()?;
Ok(operation(self))
}
pub(in crate::queue) fn with_open_result<T, E>(
&self,
operation: impl FnOnce(&Self) -> Result<T, E>,
) -> Result<T, E>
where
E: From<PlayError>,
{
let _admission = self.lock_admission();
self.ensure_open().map_err(E::from)?;
operation(self)
}
delegate::delegate! {
to self.tracks {
#[call(lock)]
pub(super) fn lock_tracks(&self) -> std::sync::MutexGuard<'_, Vec<TrackRecord<S>>>;
#[call(lock)]
pub(super) fn lock_tracks_mut(&self) -> std::sync::MutexGuard<'_, Vec<TrackRecord<S>>>;
pub(super) fn set_status(&self, id: TrackId, status: crate::event::TrackStatus);
}
to self.crossfade_armed_for {
#[call(load)]
pub(super) fn read_armed_for(&self) -> CrossfadeArm;
#[call(take_if_matches)]
pub(super) fn take_armed_for_if_matches(&self, id: TrackId) -> bool;
#[call(store)]
pub(super) fn write_armed_for(&self, arm: CrossfadeArm);
}
to self.cached_position {
#[call(load)]
pub(super) fn read_cached_position(&self) -> CachedPosition;
#[call(store)]
pub(super) fn write_cached_position(&self, pos: CachedPosition);
}
}
}
impl<S> Drop for Queue<S>
where
S: HasPool<u8> + HasPool<f32> + Send + Sync + 'static,
{
fn drop(&mut self) {
self.control.invalidate();
}
}
#[cfg(test)]
pub(crate) mod tests {
use core::sync::atomic::{AtomicU64, Ordering};
use std::{
num::NonZeroU32,
sync::mpsc::{self, RecvTimeoutError},
thread,
};
use kithara_audio::ConsumerWakeMode;
use kithara_events::{Envelope, EventReceiver};
use kithara_platform::{
sync::{Arc, Mutex},
time::{Duration, Instant, timeout},
};
use kithara_play::{
AllocatedSlot, BeatGrid, Cmd, NodeInputs, PlayError, PlayWorker, PlayWorkerConfig,
PlayerConfig, Reply, SessionBinding, SessionDispatcher, SessionSampleRate, SharedEq,
SlotId, bridge::slot_channels,
};
use kithara_test_utils::kithara;
use super::*;
use crate::{
event::QueueEvent,
test_pools::{TestPools, pools},
};
pub(crate) const TEST_SAMPLE_RATE: NonZeroU32 = match NonZeroU32::new(44_100) {
Some(sample_rate) => sample_rate,
None => unreachable!(),
};
pub(in crate::queue) fn make_store() -> AssetStore<TestPools> {
AssetStore::builder(pools())
.backend(StorageBackend::Memory)
.build()
}
pub(in crate::queue) fn make_queue() -> Queue<TestPools> {
Queue::new(queue_config())
}
struct TestSession {
next_slot: AtomicU64,
nodes: Mutex<Vec<NodeInputs>>,
}
impl SessionDispatcher<TestPools> for TestSession {
fn consumer_wake_mode(&self) -> ConsumerWakeMode {
ConsumerWakeMode::RealtimeDeferred
}
fn exec(&self, cmd: Cmd<TestPools>) -> Result<Reply, PlayError> {
let reply = match cmd {
Cmd::RegisterPlayer { .. } => {
Reply::PlayerRegistered(kithara_play::session::RegisteredPlayer {
id: 1,
eq: SharedEq::new(10),
})
}
Cmd::AllocateSlot { .. } => {
let slot = SlotId::new(self.next_slot.fetch_add(1, Ordering::Relaxed));
let (inputs, control) = slot_channels(SharedEq::new(10));
self.nodes.lock().push(inputs);
Reply::SlotAllocated(AllocatedSlot::new(control, slot))
}
Cmd::QuerySampleRate => Reply::SampleRate(SessionSampleRate::new(None, 44_100)),
Cmd::QueryStreamShape => Reply::StreamShape(None),
_ => Reply::Ok,
};
Ok(reply)
}
}
pub(crate) fn test_session() -> SessionBinding<TestPools> {
SessionBinding::new(
Arc::new(TestSession {
next_slot: AtomicU64::new(0),
nodes: Mutex::default(),
}),
TEST_SAMPLE_RATE,
)
}
fn queue_config() -> QueueConfig<TestPools> {
QueueConfig::builder()
.player(player())
.store(make_store())
.build()
}
fn player() -> PlayerImpl<TestPools> {
let worker = PlayWorker::new(PlayWorkerConfig::builder(pools()).build());
PlayerImpl::new(
PlayerConfig::builder()
.sample_rate(TEST_SAMPLE_RATE)
.worker(worker)
.session(test_session())
.build(),
)
}
pub(in crate::queue) async fn wait_for_queue_event<F>(
rx: &mut EventReceiver<QueueEvent>,
mut matches: F,
timeout_ms: u64,
) -> bool
where
F: FnMut(&QueueEvent) -> bool,
{
let deadline = Instant::now() + Duration::from_millis(timeout_ms);
loop {
let remaining = deadline.saturating_duration_since(Instant::now());
if remaining.is_zero() {
return false;
}
match timeout(remaining, rx.recv()).await {
Ok(Ok(Envelope { event: ev, .. })) if matches(&ev) => return true,
Ok(Ok(_)) => continue,
Ok(Err(_)) | Err(_) => return false,
}
}
}
#[kithara::test]
fn queue_new_constructs_without_panic() {
let _queue = make_queue();
}
#[kithara::test]
fn queue_preserves_the_resident_players_canonical_grid() {
let player = player();
let grid_id = player.id();
let snapshot = player.snapshot();
let queue = Queue::new(QueueConfig::builder().player(player).build());
assert_eq!(queue.id(), grid_id);
assert_eq!(queue.snapshot(), snapshot);
}
#[kithara::test]
fn queue_control_rejects_mutation_after_close() {
let queue = make_queue();
let control = queue.control.clone();
control.close().expect("unstarted fixture must close");
assert!(control.runtime.shutdown.is_cancelled());
assert!(matches!(
control.append("https://example.com/a.mp3"),
Err(crate::QueueError::Play(PlayError::Closed))
));
assert!(queue.is_empty());
}
#[kithara::test]
fn close_waits_for_an_admitted_queue_mutation() {
let queue = make_queue();
let mutation_control = queue.control.clone();
let close_control = queue.control.clone();
let (entered_tx, entered_rx) = mpsc::channel();
let (release_tx, release_rx) = mpsc::channel();
let (mutation_tx, mutation_rx) = mpsc::channel();
let mutation = thread::spawn(move || {
let result = mutation_control.with_open(|_| {
entered_tx.send(()).expect("test receiver remains alive");
release_rx.recv().expect("test sender releases mutation");
});
mutation_tx
.send(result)
.expect("test receiver remains alive");
});
entered_rx
.recv()
.expect("mutation must enter the queue admission gate");
let (close_tx, close_rx) = mpsc::channel();
let close = thread::spawn(move || {
close_tx
.send(close_control.close())
.expect("test receiver remains alive");
});
assert!(
matches!(
close_rx.recv_timeout(Duration::from_millis(50)),
Err(RecvTimeoutError::Timeout)
),
"close must not overtake an admitted queue mutation"
);
release_tx.send(()).expect("mutation thread remains alive");
mutation_rx
.recv()
.expect("mutation must complete after release")
.expect("admitted mutation remains open");
close_rx
.recv()
.expect("close must complete after the mutation")
.expect("unstarted fixture must close");
mutation.join().expect("mutation thread must not panic");
close.join().expect("close thread must not panic");
assert!(queue.is_closed());
}
#[kithara::test]
fn the_configured_prefetch_lead_reaches_the_player() {
let queue = Queue::new(
QueueConfig::builder()
.player(player())
.store(make_store())
.prefetch_duration(8.0)
.build(),
);
assert!((queue.player.prefetch_duration() - 8.0).abs() < f32::EPSILON);
}
#[kithara::test]
fn crossfade_arm_disarmed_after_construction() {
let queue = make_queue();
assert_eq!(queue.read_armed_for(), CrossfadeArm::Disarmed);
}
#[kithara::test]
fn crossfade_arm_take_only_disarms_matching_track() {
let queue = make_queue();
queue.write_armed_for(CrossfadeArm::armed(TrackId(9)));
assert!(!queue.take_armed_for_if_matches(TrackId(10)));
assert_eq!(
queue.read_armed_for(),
CrossfadeArm::Armed {
for_track: TrackId(9),
}
);
assert!(queue.take_armed_for_if_matches(TrackId(9)));
assert_eq!(queue.read_armed_for(), CrossfadeArm::Disarmed);
}
#[kithara::test]
fn cached_position_unknown_after_construction() {
let queue = make_queue();
assert_eq!(Option::<f64>::from(queue.read_cached_position()), None);
}
#[kithara::test]
fn cached_position_round_trips_through_queue() {
let queue = make_queue();
queue.write_cached_position(CachedPosition::known(12.5));
assert_eq!(
Option::<f64>::from(queue.read_cached_position()),
Some(12.5)
);
}
#[kithara::test]
fn select_phase_idle_after_construction() {
let queue = make_queue();
assert!(matches!(
*queue.lock_pending_select_mut(),
SelectPhase::Idle
));
}
}