Skip to main content

kithara_queue/queue/selection/
control.rs

1use 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    /// Select a track by id, applying the given [`Transition`]. If the
21    /// track is still loading or pending, both the id and the
22    /// transition are stashed and applied when loading finishes.
23    ///
24    /// # Errors
25    /// Returns [`QueueError::UnknownTrackId`] if `id` is not in the queue,
26    /// [`QueueError::NotReady`] if the track is in a terminal failed state,
27    /// or [`QueueError::Play`] if the underlying `select_item` call fails.
28    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    /// Serializes the whole select against a concurrent `spawn_apply_after_load` completion so
75    /// marking the prior pending attempt `Cancelled` and a loading track's apply never interleave,
76    /// which would let the superseded track barge in.
77    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    /// `is_playing` is a session flag, not a verdict on the current item: the render thread queues
124    /// the natural end but clears the flag only at the next `process`, so a repeat-one advance must
125    /// re-select the item that just ended despite the flag.
126    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}