kithara_queue/queue/selection/
control.rs1use kithara_bufpool::HasPool;
2use kithara_events::TrackId;
3use kithara_play::{SelectTransition, SelectionPlayback};
4use smallvec::SmallVec;
5
6use super::super::{
7 QueueControl,
8 types::{CrossfadeArm, PendingSelect, Transition},
9};
10use crate::{
11 attempts::LoadClass,
12 error::QueueError,
13 event::{AdvanceReason, QueueEvent, TrackStatus},
14};
15
16impl<S> QueueControl<S>
17where
18 S: HasPool<u8> + HasPool<f32> + Send + Sync + 'static,
19{
20 pub fn select(&self, id: TrackId, transition: Transition) -> Result<(), QueueError> {
29 self.with_open_result(|queue| {
30 queue.select_with(
31 id,
32 transition,
33 AdvanceReason::UserSelect,
34 SelectionPlayback::Play,
35 )
36 })
37 }
38
39 pub(in crate::queue) fn select_loaded_item(
40 &self,
41 index: usize,
42 id: TrackId,
43 crossfade: kithara_play::CrossfadeSettings,
44 reason: AdvanceReason,
45 playback: SelectionPlayback,
46 ) -> Result<(), QueueError> {
47 let was_playing = self.player.is_playing();
48 if was_playing && crossfade.duration > 0.0 {
49 self.bus.publish(QueueEvent::CrossfadeStarted {
50 settings: crossfade,
51 });
52 }
53 self.player.select_item_with_crossfade(
54 index,
55 SelectTransition {
56 playback,
57 crossfade,
58 },
59 )?;
60 let ids = self
61 .tracks()
62 .into_iter()
63 .map(|track| track.id)
64 .collect::<SmallVec<[_; 16]>>();
65 self.lock_navigation_mut().select(id, &ids);
66 self.bus.publish(QueueEvent::CurrentTrackAdvance {
67 reason,
68 id: Some(id),
69 });
70 self.set_status(id, TrackStatus::Consumed);
71 Ok(())
72 }
73
74 pub(in crate::queue) fn select_with(
78 &self,
79 id: TrackId,
80 transition: Transition,
81 reason: AdvanceReason,
82 playback: SelectionPlayback,
83 ) -> Result<(), QueueError> {
84 if matches!(
85 reason,
86 AdvanceReason::UserSelect
87 | AdvanceReason::UserNext
88 | AdvanceReason::UserPrev
89 | AdvanceReason::RemovedCurrent
90 ) {
91 self.autoplay_target.store(CrossfadeArm::Disarmed);
92 }
93 let default = *self
94 .crossfade_settings
95 .lock()
96 .unwrap_or_else(std::sync::PoisonError::into_inner);
97 let settings = transition.settings(default).validate()?;
98 let _apply = self.lock_select_apply();
99 self.select_with_reason_locked(id, settings, reason, playback)
100 }
101
102 pub(in crate::queue) fn select_with_reason(
103 &self,
104 id: TrackId,
105 transition: Transition,
106 reason: AdvanceReason,
107 ) -> Result<(), QueueError> {
108 let playback = if matches!(
109 reason,
110 AdvanceReason::NaturalEof
111 | AdvanceReason::TrackFailed
112 | AdvanceReason::CrossfadePreArm
113 | AdvanceReason::Repeat
114 ) || self.player.is_playing()
115 {
116 SelectionPlayback::Play
117 } else {
118 SelectionPlayback::Pause
119 };
120 self.select_with(id, transition, reason, playback)
121 }
122
123 pub(in crate::queue) fn select_with_reason_locked(
127 &self,
128 id: TrackId,
129 settings: kithara_play::CrossfadeSettings,
130 reason: AdvanceReason,
131 playback: SelectionPlayback,
132 ) -> Result<(), QueueError> {
133 let (index, status) = {
134 let guard = self.lock_tracks();
135 guard
136 .iter()
137 .enumerate()
138 .find(|(_, e)| e.id == id)
139 .map(|(i, e)| (i, e.status.clone()))
140 .ok_or(QueueError::UnknownTrackId(id))?
141 };
142
143 if self.player.current_index() == index
144 && matches!(status, TrackStatus::Consumed)
145 && matches!(
146 reason,
147 AdvanceReason::UserSelect
148 | AdvanceReason::UserNext
149 | AdvanceReason::UserPrev
150 | AdvanceReason::RemovedCurrent
151 | AdvanceReason::NaturalEof
152 )
153 {
154 self.cancel_stale_pending(id);
155 if self.player.is_playing() && reason != AdvanceReason::NaturalEof {
156 return Ok(());
157 }
158 let finished = self.player.playback_snapshot().is_some_and(|snapshot| {
159 snapshot.duration() > 0.0 && snapshot.position() >= snapshot.duration()
160 });
161 if reason == AdvanceReason::NaturalEof || finished {
162 self.player.seek_seconds(0.0)?;
163 }
164 if playback == SelectionPlayback::Play {
165 self.player.play();
166 } else {
167 self.player.pause();
168 }
169 let ids = self
170 .tracks()
171 .into_iter()
172 .map(|track| track.id)
173 .collect::<SmallVec<[_; 16]>>();
174 self.lock_navigation_mut().select(id, &ids);
175 self.bus.publish(QueueEvent::CurrentTrackAdvance {
176 reason,
177 id: Some(id),
178 });
179 return Ok(());
180 }
181
182 match status {
183 TrackStatus::Loaded => {
184 self.cancel_stale_pending(id);
185 self.select_loaded_item(index, id, settings, reason, playback)?;
186 Ok(())
187 }
188 TrackStatus::Pending | TrackStatus::Loading | TrackStatus::Slow => {
189 self.override_pending_select(PendingSelect {
190 reason,
191 settings,
192 playback,
193 id,
194 });
195 self.promote_pending_load(id);
196 Ok(())
197 }
198 TrackStatus::Cancelled | TrackStatus::Consumed | TrackStatus::Failed(_) => {
199 let source = self.tracks.source(id).ok_or(QueueError::NotReady(id))?;
200 self.override_pending_select(PendingSelect {
201 reason,
202 settings,
203 playback,
204 id,
205 });
206 self.set_status(id, TrackStatus::Pending);
207 self.spawn_apply_after_load(id, source, LoadClass::Interactive);
208 Ok(())
209 }
210 }
211 }
212}
213
214#[cfg(test)]
215mod tests {
216 use kithara_test_utils::kithara;
217
218 use super::{super::super::types::SelectPhase, *};
219 use crate::{
220 event::QueueEvent,
221 queue::state::tests::{make_queue, wait_for_queue_event},
222 };
223
224 fn append(queue: &crate::Queue<crate::test_pools::TestPools>, source: &str) -> TrackId {
225 queue
226 .append(source)
227 .expect("BUG: open queue must accept a track")
228 }
229
230 #[kithara::test(tokio)]
231 async fn select_unknown_id_errors() {
232 let queue = make_queue();
233 let err = queue
234 .select(TrackId(999), Transition::None)
235 .expect_err("unknown id should error");
236 assert!(matches!(err, QueueError::UnknownTrackId(_)));
237 }
238
239 #[kithara::test(tokio)]
240 async fn select_pending_track_stashes_pending_select() {
241 let queue = make_queue();
242 let id = append(&queue, "https://example.com/a.mp3");
243 let _ = queue.select(id, Transition::None);
244 let phase = *queue
245 .pending_select
246 .lock()
247 .expect("BUG: pending_select Mutex is not held across await");
248 match phase {
249 SelectPhase::Pending(pending) => {
250 assert_eq!(pending.id, id);
251 assert_eq!(pending.settings.duration, 0.0);
252 }
253 SelectPhase::Idle => panic!("BUG: select stashes pending entry"),
254 }
255 }
256
257 #[kithara::test(tokio)]
258 async fn advance_to_next_on_empty_emits_queue_ended() {
259 let queue = make_queue();
260 let mut rx = queue.subscribe();
261 assert!(
262 queue
263 .advance_to_next_inner(Transition::Crossfade, AdvanceReason::NaturalEof)
264 .expect("BUG: open queue advance must be admitted")
265 .is_none()
266 );
267 let saw_ended =
268 wait_for_queue_event(&mut rx, |ev| matches!(ev, QueueEvent::QueueEnded), 200).await;
269 assert!(saw_ended);
270 }
271
272 #[kithara::test(tokio)]
273 async fn manual_next_at_exhaustion_does_not_emit_queue_ended() {
274 let queue = make_queue();
275 let mut rx = queue.subscribe();
276 assert_eq!(queue.next(Transition::None).expect("manual next"), None);
277 assert!(
278 !wait_for_queue_event(&mut rx, |ev| matches!(ev, QueueEvent::QueueEnded), 50).await
279 );
280 }
281
282 #[kithara::test(tokio)]
283 async fn advance_to_next_cycles_then_emits_queue_ended() {
284 let queue = make_queue();
285 let a = append(&queue, "https://example.com/a.mp3");
286 let b = append(&queue, "https://example.com/b.mp3");
287 queue.lock_navigation_mut().select(b, &[a, b]);
288 let mut rx = queue.subscribe();
289
290 assert!(
291 queue
292 .advance_to_next_inner(Transition::Crossfade, AdvanceReason::NaturalEof)
293 .expect("BUG: open queue advance must be admitted")
294 .is_none()
295 );
296
297 let saw_ended =
298 wait_for_queue_event(&mut rx, |ev| matches!(ev, QueueEvent::QueueEnded), 400).await;
299 assert!(saw_ended, "QueueEnded should be broadcast at end-of-queue");
300 }
301
302 #[kithara::test(tokio)]
303 async fn admitted_pending_successor_becomes_navigation_authority() {
304 let queue = make_queue();
305 let first = append(&queue, "https://example.com/a.mp3");
306 let second = append(&queue, "https://example.com/b.mp3");
307 queue.lock_navigation_mut().select(first, &[first, second]);
308 queue.set_status(first, TrackStatus::Consumed);
309 queue.set_status(second, TrackStatus::Pending);
310
311 assert_eq!(
312 queue
313 .advance_to_next_inner(Transition::Crossfade, AdvanceReason::NaturalEof)
314 .expect("BUG: open queue advance must be admitted"),
315 Some(second)
316 );
317 assert_eq!(
318 queue.lock_navigation().current(),
319 Some(second),
320 "admitted automatic successor must remain authoritative while loading"
321 );
322 let SelectPhase::Pending(pending) = *queue.lock_pending_select_mut() else {
323 panic!("successor selection must remain pending")
324 };
325 assert_eq!(pending.playback, SelectionPlayback::Play);
326 assert_eq!(pending.reason, AdvanceReason::NaturalEof);
327 }
328
329 #[kithara::test(tokio)]
330 async fn pending_override_latches_profile_without_mutating_default() {
331 let queue = make_queue();
332 let id = append(&queue, "https://example.com/a.mp3");
333 let configured = kithara_play::CrossfadeSettings::new(
334 2.0,
335 kithara_play::CrossfadeCurve::EqualPower,
336 1.0,
337 0.5,
338 )
339 .expect("valid settings");
340 let override_settings = kithara_play::CrossfadeSettings::new(
341 4.0,
342 kithara_play::CrossfadeCurve::Linear,
343 0.25,
344 0.3,
345 )
346 .expect("valid settings");
347 queue
348 .set_crossfade_settings(configured)
349 .expect("valid settings");
350 queue
351 .select(
352 id,
353 Transition::CrossfadeWith {
354 settings: override_settings,
355 },
356 )
357 .expect("pending selection admitted");
358 queue
359 .set_crossfade_settings(kithara_play::CrossfadeSettings::default())
360 .expect("valid settings");
361
362 let SelectPhase::Pending(pending) = *queue.lock_pending_select_mut() else {
363 panic!("selection must remain pending")
364 };
365 assert_eq!(pending.settings, override_settings);
366 assert_eq!(
367 queue.crossfade_settings(),
368 kithara_play::CrossfadeSettings::default()
369 );
370 }
371}