diff --git a/crates/player/src/engine.rs b/crates/player/src/engine.rs index 1148dec..a30f83b 100644 --- a/crates/player/src/engine.rs +++ b/crates/player/src/engine.rs @@ -85,6 +85,23 @@ struct EngineState { silence_skip_enabled: bool, /// Сколько подряд идёт тишина прямо сейчас — сбрасывается любым не-тихим пакетом. silent_run: Duration, + /// Фоновая загрузка сетевого трека (SMB/http(s)://) в процессе, если есть. Открытие + /// такого трека может занять секунды (таймаут соединения при недоступном сервере) — это + /// НЕ должно блокировать всю очередь команд движка целиком, иначе, например, попытка + /// запустить недоступный подкаст не давала бы вообще ничего другого включить/переключить + /// (в том числе аудиокнигу), пока не истечёт таймаут — реальный баг, так и было раньше. + pending_load: Option, + /// Растёт при каждой новой попытке загрузки трека — если результат фоновой загрузки + /// пришёл, но поколение уже не совпадает с текущим (пока грузили — пользователь сменил + /// трек ещё раз), результат просто отбрасывается как устаревший. + load_generation: u64, +} + +struct PendingLoad { + generation: u64, + item: QueueItem, + queue_index: Option, + rx: std::sync::mpsc::Receiver)>>, } fn run(cmd_rx: Receiver, evt_tx: Sender) -> anyhow::Result<()> { @@ -167,6 +184,8 @@ fn run(cmd_rx: Receiver, evt_tx: Sender) -> anyhow:: pitch_resampler: None, silence_skip_enabled: false, silent_run: Duration::ZERO, + pending_load: None, + load_generation: 0, }; let mut current_track: Option = None; @@ -193,6 +212,12 @@ fn run(cmd_rx: Receiver, evt_tx: Sender) -> anyhow:: } } + // Опрашиваем фоновую загрузку сетевого трека (если есть) — неблокирующе, тем же + // способом, что и команды выше. Пока результата нет, движок просто идёт дальше как + // обычно (продолжает декодировать текущий трек, если есть, готов принять новые + // команды на следующей итерации) — сетевое подключение никого не задерживает. + poll_pending_load(&mut state, &mut current_track, &mut resampler, &paused_flag, &evt_tx); + if state.status == PlaybackStatus::Playing && current_track.is_some() { let need_more = producer.free_len() > 4096; if need_more { @@ -315,6 +340,7 @@ fn convert_channels(samples: &[f32], from_channels: usize, to_channels: usize) - out } + #[cfg(test)] mod network_error_tests { use super::*; @@ -560,6 +586,144 @@ fn fade_gain(state: &EngineState) -> Option { fade_gain_at(state.current_track_end, absolute_position(state)) } +#[cfg(test)] +mod nonblocking_load_tests { + use super::*; + use std::sync::atomic::AtomicBool; + + fn test_state() -> EngineState { + EngineState { + queue: Queue::default(), + speed: 1.0, + eq: None, + eq_gains: [0.0; EQ_BANDS], + status: PlaybackStatus::Stopped, + device_rate: 44100, + device_channels: 2, + frames_decoded_native: 0, + native_rate: 44100, + current_track_start: Duration::ZERO, + current_track_end: None, + is_radio: false, + is_network_source: false, + spectrum_buffer: Vec::new(), + last_spectrum_send: Instant::now(), + fade_out_enabled: false, + loudness_norm_enabled: false, + current_track_replay_gain: None, + boost_gain: 1.0, + stretcher: None, + pitch: 1.0, + pitch_resampler: None, + silence_skip_enabled: false, + silent_run: Duration::ZERO, + pending_load: None, + load_generation: 0, + } + } + + /// Реальный баг: попытка запустить недоступный (зависший, не отвечающий) сетевой + /// трек не должна мешать сразу следом запустить что-то совсем другое (например, + /// локальную аудиокнигу). Нужен запущенный вручную hang_server.py на 127.0.0.1:8949 + /// (сервер, который принимает соединение, но никогда не отвечает — честная имитация + /// зависшего/недоступного подкаст-сервера) и /tmp/local_audiobook.mp3. + /// Запускать явно: `cargo test -- --ignored --nocapture`. + #[test] + #[ignore] + fn switching_to_local_track_is_not_blocked_by_hung_network_load() { + let mut state = test_state(); + let mut current_track: Option = None; + let mut resampler: Option = None; + let paused_flag = Arc::new(AtomicBool::new(false)); + let (evt_tx, evt_rx) = std::sync::mpsc::channel(); + + // Шаг 1: запрашиваем недоступный (зависший) сетевой трек — ровно как SetQueue + // на подкаст, который не отвечает. + state.queue.set_tracks( + vec![QueueItem::whole_file(PathBuf::from("http://127.0.0.1:8949/hobbytalks.mp3"))], + 0, + ); + let start = Instant::now(); + load_current(&mut state, &mut current_track, &mut resampler, &paused_flag, &evt_tx); + let load_current_call_duration = start.elapsed(); + + // КЛЮЧЕВАЯ ПРОВЕРКА: сам вызов load_current для сетевого трека должен вернуться + // почти мгновенно (не ждать таймаута соединения) — открытие ушло в фон. + assert!( + load_current_call_duration < Duration::from_millis(500), + "load_current for a network path took {load_current_call_duration:?} — it must return immediately, not block" + ); + assert!(state.pending_load.is_some(), "network load should be pending in the background"); + assert!(current_track.is_none(), "no track installed yet — background load hasn't finished"); + + // Шаг 2: РОВНО ТА ЖЕ СИТУАЦИЯ, ЧТО У ПОЛЬЗОВАТЕЛЯ — пока подкаст всё ещё висит на + // фоновом подключении, пробуем запустить совсем другой (локальный) трек. + state.queue.set_tracks( + vec![QueueItem::whole_file(PathBuf::from("/tmp/local_audiobook.mp3"))], + 0, + ); + let start2 = Instant::now(); + load_current(&mut state, &mut current_track, &mut resampler, &paused_flag, &evt_tx); + let switch_duration = start2.elapsed(); + + assert!( + switch_duration < Duration::from_millis(500), + "switching to a local file took {switch_duration:?} — must not be blocked by the still-hanging network load" + ); + assert!(current_track.is_some(), "local audiobook should be loaded and playing right away"); + assert!(state.pending_load.is_none(), "local file load doesn't need a pending background load"); + + // И убеждаемся, что реально пришло событие TrackChanged на локальный файл (не завис + // тихо где-то) — читаем события из канала, которые должны были прийти. + let mut saw_local_track_changed = false; + while let Ok(evt) = evt_rx.try_recv() { + if let PlayerEvent::TrackChanged(info) = evt { + if info.path == PathBuf::from("/tmp/local_audiobook.mp3") { + saw_local_track_changed = true; + } + } + } + assert!(saw_local_track_changed, "should have received TrackChanged for the local audiobook"); + + println!( + "load_current(hung network)={load_current_call_duration:?}, switch_to_local={switch_duration:?} — both non-blocking, as expected" + ); + } + + /// Когда фоновая (зависшая) загрузка ВСЁ ЖЕ рано или поздно завершится (или провалится + /// таймаутом) уже ПОСЛЕ того, как пользователь переключился на другой трек — её + /// результат должен быть тихо отброшен (не должен неожиданно подменить уже играющий + /// новый трек). Проверяем через generation напрямую, без реального ожидания таймаута. + #[test] + fn stale_background_load_result_is_discarded_after_switching_tracks() { + let mut state = test_state(); + let mut current_track: Option = None; + let mut resampler: Option = None; + let paused_flag = Arc::new(AtomicBool::new(false)); + let (evt_tx, evt_rx) = std::sync::mpsc::channel(); + + // Симулируем "старую" фоновую загрузку с generation=1, которая вот-вот (успешно!) + // завершится — но пользователь уже переключился на что-то другое. + let (tx, rx) = std::sync::mpsc::channel(); + state.pending_load = Some(PendingLoad { + generation: 1, + item: QueueItem::whole_file(PathBuf::from("http://example.com/stale.mp3")), + queue_index: Some(0), + rx, + }); + state.load_generation = 2; // уже "новее" — пользователь успел переключиться + + // Старая загрузка наконец отвечает (успешно!) — но она больше не актуальна. + let _ = tx.send(Err(anyhow::anyhow!("irrelevant — should never be processed"))); + + poll_pending_load(&mut state, &mut current_track, &mut resampler, &paused_flag, &evt_tx); + + assert!(state.pending_load.is_none(), "stale pending load should be cleared after polling"); + assert!(current_track.is_none(), "stale result must not install anything"); + assert!(evt_rx.try_recv().is_err(), "no event should be sent for a discarded stale result"); + } +} + #[cfg(test)] mod fade_tests { use super::*; @@ -913,6 +1077,44 @@ fn install_track( let _ = evt_tx.send(PlayerEvent::StatusChanged(state.status)); } +/// Неблокирующе проверяет, не подоспел ли результат фоновой загрузки сетевого трека +/// (см. `load_current`) — вызывать каждую итерацию главного цикла движка, тем же способом, +/// что и разбор команд. Устаревшие результаты (пользователь успел переключить трек, пока +/// грузилось — несовпадение generation) молча отбрасываются, ничего не делая. +fn poll_pending_load( + state: &mut EngineState, + current_track: &mut Option, + resampler: &mut Option, + paused_flag: &Arc, + evt_tx: &Sender, +) { + let Some(poll_result) = state.pending_load.as_ref().map(|p| p.rx.try_recv()) else { + return; + }; + match poll_result { + Ok(result) => { + let pending = state.pending_load.take().unwrap(); + if pending.generation == state.load_generation { + finish_load( + state, + current_track, + resampler, + paused_flag, + evt_tx, + pending.item, + pending.queue_index, + result, + ); + } + // иначе — устарело, просто отбрасываем результат, ничего не делаем. + } + Err(TryRecvError::Empty) => {} // ещё не готово, подождём следующей итерации + Err(TryRecvError::Disconnected) => { + state.pending_load = None; + } + } +} + fn load_current( state: &mut EngineState, current_track: &mut Option, @@ -920,21 +1122,66 @@ fn load_current( paused_flag: &Arc, evt_tx: &Sender, ) { + // Любая предыдущая фоновая загрузка теперь неактуальна (её результат, если придёт + // позже, будет отброшен по несовпадению поколения — см. опрос pending_load в run()). + state.load_generation += 1; + state.pending_load = None; + let Some(item) = state.queue.current_item().cloned() else { *current_track = None; return; }; let queue_index = state.queue.current_index(); - match open_track(&item.path) { - Ok(mut track) => { - // Стриминг по сети (SMB-шара или http(s):// — например, эпизод подкаста без - // скачивания) — сетевые сбои для таких треков не должны валить весь движок, - // см. комментарий у поля is_network_source. - let path_str = item.path.to_string_lossy(); - state.is_network_source = smbfs::is_smb_path(&item.path) - || path_str.starts_with("http://") - || path_str.starts_with("https://"); + if is_network_path(&item.path) { + // Не блокируем движок сетевым подключением — открываем в отдельном потоке и + // возвращаемся немедленно; результат подхватит опрос pending_load в главном цикле, + // как только он будет готов. До этого момента движок как ни в чём не бывало + // обрабатывает остальные команды (в т.ч. запуск СОВСЕМ ДРУГОГО трека). + let generation = state.load_generation; + let path = item.path.clone(); + let (tx, rx) = std::sync::mpsc::channel(); + std::thread::spawn(move || { + let result = open_track(&path).map(|track| { + let tags = read_tags(&path).ok(); + (track, tags) + }); + let _ = tx.send(result); + }); + state.pending_load = Some(PendingLoad { generation, item, queue_index, rx }); + return; + } + + // Локальный файл — открытие практически мгновенное, фоновый поток тут не нужен и + // добавил бы только лишнюю задержку в один тик цикла. + let result = open_track(&item.path).map(|t| { + let tags = read_tags(&item.path).ok(); + (t, tags) + }); + finish_load(state, current_track, resampler, paused_flag, evt_tx, item, queue_index, result); +} + +fn is_network_path(path: &std::path::Path) -> bool { + let s = path.to_string_lossy(); + smbfs::is_smb_path(path) || s.starts_with("http://") || s.starts_with("https://") +} + +/// Общая часть завершения загрузки трека — вызывается либо сразу (локальные файлы), либо +/// когда подоспел результат фоновой загрузки (сетевые треки), см. `load_current`. +#[allow(clippy::too_many_arguments)] +fn finish_load( + state: &mut EngineState, + current_track: &mut Option, + resampler: &mut Option, + paused_flag: &Arc, + evt_tx: &Sender, + item: QueueItem, + queue_index: Option, + result: anyhow::Result<(DecodedTrack, Option)>, +) { + match result { + Ok((mut track, tags)) => { + state.is_network_source = is_network_path(&item.path); // Виртуальный трек (например, кусок cue-листа) — сразу прыгаем на его начало // внутри физического файла, вместо проигрывания с самого начала файла. @@ -960,7 +1207,7 @@ fn load_current( (item.start.as_secs_f64() * state.native_rate as f64) as u64; } - let mut info = read_tags(&item.path).unwrap_or_else(|_| TrackInfo { + let mut info = tags.unwrap_or_else(|| TrackInfo { path: item.path.clone(), title: item .path