This commit is contained in:
Magnus Root 2026-08-21 11:18:32 +03:00
parent c5c202384b
commit d450f20339

View file

@ -85,6 +85,23 @@ struct EngineState {
silence_skip_enabled: bool,
/// Сколько подряд идёт тишина прямо сейчас — сбрасывается любым не-тихим пакетом.
silent_run: Duration,
/// Фоновая загрузка сетевого трека (SMB/http(s)://) в процессе, если есть. Открытие
/// такого трека может занять секунды (таймаут соединения при недоступном сервере) — это
/// НЕ должно блокировать всю очередь команд движка целиком, иначе, например, попытка
/// запустить недоступный подкаст не давала бы вообще ничего другого включить/переключить
/// (в том числе аудиокнигу), пока не истечёт таймаут — реальный баг, так и было раньше.
pending_load: Option<PendingLoad>,
/// Растёт при каждой новой попытке загрузки трека — если результат фоновой загрузки
/// пришёл, но поколение уже не совпадает с текущим (пока грузили — пользователь сменил
/// трек ещё раз), результат просто отбрасывается как устаревший.
load_generation: u64,
}
struct PendingLoad {
generation: u64,
item: QueueItem,
queue_index: Option<usize>,
rx: std::sync::mpsc::Receiver<anyhow::Result<(DecodedTrack, Option<TrackInfo>)>>,
}
fn run(cmd_rx: Receiver<PlayerCommand>, evt_tx: Sender<PlayerEvent>) -> anyhow::Result<()> {
@ -167,6 +184,8 @@ fn run(cmd_rx: Receiver<PlayerCommand>, evt_tx: Sender<PlayerEvent>) -> anyhow::
pitch_resampler: None,
silence_skip_enabled: false,
silent_run: Duration::ZERO,
pending_load: None,
load_generation: 0,
};
let mut current_track: Option<DecodedTrack> = None;
@ -193,6 +212,12 @@ fn run(cmd_rx: Receiver<PlayerCommand>, evt_tx: Sender<PlayerEvent>) -> 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<f32> {
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<DecodedTrack> = None;
let mut resampler: Option<SpeedResampler> = 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<DecodedTrack> = None;
let mut resampler: Option<SpeedResampler> = 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<DecodedTrack>,
resampler: &mut Option<SpeedResampler>,
paused_flag: &Arc<AtomicBool>,
evt_tx: &Sender<PlayerEvent>,
) {
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<DecodedTrack>,
@ -920,21 +1122,66 @@ fn load_current(
paused_flag: &Arc<AtomicBool>,
evt_tx: &Sender<PlayerEvent>,
) {
// Любая предыдущая фоновая загрузка теперь неактуальна (её результат, если придёт
// позже, будет отброшен по несовпадению поколения — см. опрос 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<DecodedTrack>,
resampler: &mut Option<SpeedResampler>,
paused_flag: &Arc<AtomicBool>,
evt_tx: &Sender<PlayerEvent>,
item: QueueItem,
queue_index: Option<usize>,
result: anyhow::Result<(DecodedTrack, Option<TrackInfo>)>,
) {
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