/
magnusroot
/
nm
Обзор
Документация
Войти
/
magnusroot
/
nm
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
master
crates/player/src/engine.rs
717 строк
31 KB
Magnus Root
First version
13 июл 2026, 09:24
13 июл 2026, 09:24
50f8b8c
Код
Авторство
О чём код?
use std::path::PathBuf; use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; use std::sync::mpsc::{Receiver, Sender, TryRecvError}; use std::sync::{Arc, Mutex}; use std::time::{Duration, Instant}; use cpal::traits::{DeviceTrait, HostTrait, StreamTrait}; use ringbuf::{HeapConsumer, HeapProducer, HeapRb}; use crate::decode::{next_packet_samples, open_track, open_url, read_tags, seek as decode_seek, DecodedTrack}; use crate::eq::Equalizer; use crate::leveler::Leveler; use crate::queue::Queue; use crate::resample::SpeedResampler; use crate::spectrum::{SpectrumAnalyzer, FFT_SIZE}; use crate::types::{OpenMode, PlaybackStatus, PlayerCommand, PlayerEvent, QueueItem, TrackInfo, EQ_BANDS}; /// Не чаще раза в этот интервал шлём событие со спектром — иначе можно засыпать UI-канал /// событиями чаще, чем это осмысленно для визуализации на экране. const SPECTRUM_SEND_INTERVAL: Duration = Duration::from_millis(60); /// Точка входа: запускает движок в отдельном потоке и возвращает канал команд + канал событий. pub fn spawn() -> (Sender<PlayerCommand>, Receiver<PlayerEvent>) { let (cmd_tx, cmd_rx) = std::sync::mpsc::channel::<PlayerCommand>(); let (evt_tx, evt_rx) = std::sync::mpsc::channel::<PlayerEvent>(); std::thread::spawn(move || { if let Err(e) = run(cmd_rx, evt_tx.clone()) { let _ = evt_tx.send(PlayerEvent::Error(format!("engine crashed: {e}"))); } }); (cmd_tx, evt_rx) } struct EngineState { queue: Queue, speed: f32, eq: Option<Equalizer>, eq_gains: [f32; EQ_BANDS], status: PlaybackStatus, device_rate: u32, device_channels: usize, frames_decoded_native: u64, // абсолютная позиция в физическом файле, в native-сэмплах native_rate: u32, /// Абсолютное (от начала физического файла) начало текущего логического трека. /// Для обычных файлов — всегда ноль; для cue-треков — начало соответствующего диапазона. current_track_start: Duration, /// Абсолютный (от начала физического файла) конец текущего логического трека, если задан /// (cue-треки, кроме последнего в файле). При достижении — переход к следующему в очереди. current_track_end: Option<Duration>, /// true, пока играет интернет-радио: очередь/next/prev/repeat/seek не применяются, /// а обрыв потока просто останавливает воспроизведение вместо перехода дальше. is_radio: bool, /// Накопитель моно-сэмплов (пост-EQ, то есть то, что реально слышит пользователь) /// для спектрального анализатора. spectrum_buffer: Vec<f32>, last_spectrum_send: Instant, /// Плавное затухание громкости в последние секунды трека. fade_out_enabled: bool, /// Выравнивание громкости между треками (тег ReplayGain либо AGC-левеллер). loudness_norm_enabled: bool, /// Линейный коэффициент усиления из тега ReplayGain текущего трека, если есть. current_track_replay_gain: Option<f32>, } fn run(cmd_rx: Receiver<PlayerCommand>, evt_tx: Sender<PlayerEvent>) -> anyhow::Result<()> { let host = cpal::default_host(); let device = host .default_output_device() .ok_or_else(|| anyhow::anyhow!("no output audio device found"))?; let supported = device.default_output_config()?; let device_rate = supported.sample_rate().0; let device_channels = supported.channels() as usize; let config: cpal::StreamConfig = supported.into(); // Ring buffer между потоком декодирования и cpal callback. Consumer обёрнут в Mutex, // чтобы движок мог его очистить при seek (иначе доиграл бы старое аудио ещё пару секунд). let rb = HeapRb::<f32>::new(device_rate as usize * device_channels * 2); // ~2 сек буфера let (mut producer, consumer) = rb.split(); let consumer_shared: Arc<Mutex<HeapConsumer<f32>>> = Arc::new(Mutex::new(consumer)); let consumer_for_stream = consumer_shared.clone(); let underrun_frames = Arc::new(AtomicU64::new(0)); let paused_flag = Arc::new(AtomicBool::new(true)); let paused_for_stream = paused_flag.clone(); let muted_flag = Arc::new(AtomicBool::new(false)); let muted_for_stream = muted_flag.clone(); let stream = device.build_output_stream( &config, move |data: &mut [f32], _| { if paused_for_stream.load(Ordering::Relaxed) { data.fill(0.0); return; } let muted = muted_for_stream.load(Ordering::Relaxed); let Ok(mut consumer) = consumer_for_stream.lock() else { data.fill(0.0); return; }; let mut i = 0; while i < data.len() { match consumer.pop() { Some(s) => { data[i] = if muted { 0.0 } else { s }; i += 1; } None => { underrun_frames.fetch_add(1, Ordering::Relaxed); data[i] = 0.0; i += 1; } } } }, move |err| eprintln!("audio stream error: {err}"), None, )?; stream.play()?; let mut state = EngineState { queue: Queue::default(), speed: 1.0, eq: None, eq_gains: [0.0; EQ_BANDS], status: PlaybackStatus::Stopped, device_rate, device_channels, frames_decoded_native: 0, native_rate: device_rate, current_track_start: Duration::ZERO, current_track_end: None, is_radio: false, spectrum_buffer: Vec::with_capacity(FFT_SIZE * 2), last_spectrum_send: Instant::now(), fade_out_enabled: false, loudness_norm_enabled: false, current_track_replay_gain: None, }; let mut current_track: Option<DecodedTrack> = None; let mut resampler: Option<SpeedResampler> = None; let mut spectrum_analyzer = SpectrumAnalyzer::new(); let mut leveler = Leveler::new(); loop { // Обрабатываем все накопившиеся команды не блокируясь. loop { match cmd_rx.try_recv() { Ok(cmd) => handle_command( cmd, &mut state, &mut current_track, &mut resampler, &paused_flag, &muted_flag, &consumer_shared, &evt_tx, ), Err(TryRecvError::Empty) => break, Err(TryRecvError::Disconnected) => return Ok(()), } } if state.status == PlaybackStatus::Playing && current_track.is_some() { let need_more = producer.free_len() > 4096; if need_more { decode_and_push( &mut state, &mut current_track, &mut resampler, &mut producer, &paused_flag, &evt_tx, &mut spectrum_analyzer, &mut leveler, )?; } else { std::thread::sleep(Duration::from_millis(5)); } let _ = evt_tx.send(PlayerEvent::Position(current_position(&state))); } else { std::thread::sleep(Duration::from_millis(20)); } } } fn absolute_position(state: &EngineState) -> Duration { if state.native_rate == 0 { return Duration::ZERO; } Duration::from_secs_f64(state.frames_decoded_native as f64 / state.native_rate as f64) } /// Позиция относительно начала ТЕКУЩЕГО логического трека — то, что видит пользователь /// на прогресс-баре (для cue-трека это не то же самое, что позиция в физическом файле). fn current_position(state: &EngineState) -> Duration { absolute_position(state).saturating_sub(state.current_track_start) } /// Длительность затухания в конце трека. const FADE_OUT_DURATION: Duration = Duration::from_secs(4); /// Коэффициент усиления для плавного затухания в последние `FADE_OUT_DURATION` секунд /// трека. `None`, если конец трека неизвестен (например, интернет-радио — там затухать /// нечему, там нет конца). fn fade_gain_at(current_track_end: Option<Duration>, absolute_pos: Duration) -> Option<f32> { let end = current_track_end?; if absolute_pos >= end { return Some(0.0); } let remaining = end - absolute_pos; if remaining >= FADE_OUT_DURATION { return Some(1.0); } Some((remaining.as_secs_f32() / FADE_OUT_DURATION.as_secs_f32()).clamp(0.0, 1.0)) } fn fade_gain(state: &EngineState) -> Option<f32> { fade_gain_at(state.current_track_end, absolute_position(state)) } #[cfg(test)] mod fade_tests { use super::*; #[test] fn no_fade_far_from_end() { let end = Some(Duration::from_secs(200)); let pos = Duration::from_secs(10); assert_eq!(fade_gain_at(end, pos), Some(1.0)); } #[test] fn fades_linearly_in_last_window() { let end = Some(Duration::from_secs(200)); // ровно на середине окна затухания (4с) -> примерно половина громкости let pos = end.unwrap() - Duration::from_secs(2); let gain = fade_gain_at(end, pos).unwrap(); assert!((gain - 0.5).abs() < 0.01, "gain={gain}"); } #[test] fn silent_at_or_past_end() { let end = Some(Duration::from_secs(200)); assert_eq!(fade_gain_at(end, Duration::from_secs(200)), Some(0.0)); assert_eq!(fade_gain_at(end, Duration::from_secs(250)), Some(0.0)); } #[test] fn no_fade_when_end_unknown() { // например, интернет-радио — у него нет известного конца assert_eq!(fade_gain_at(None, Duration::from_secs(999)), None); } } fn handle_command( cmd: PlayerCommand, state: &mut EngineState, current_track: &mut Option<DecodedTrack>, resampler: &mut Option<SpeedResampler>, paused_flag: &Arc<AtomicBool>, muted_flag: &Arc<AtomicBool>, consumer_shared: &Arc<Mutex<HeapConsumer<f32>>>, evt_tx: &Sender<PlayerEvent>, ) { match cmd { PlayerCommand::OpenPath(path, mode) => { state.is_radio = false; let tracks = match mode { OpenMode::SingleTrack => vec![QueueItem::whole_file(path.clone())], OpenMode::WholeAlbum => match sibling_album_tracks(&path) { Ok(t) => t.into_iter().map(QueueItem::whole_file).collect(), Err(e) => { let _ = evt_tx.send(PlayerEvent::Error(e.to_string())); vec![QueueItem::whole_file(path.clone())] } }, }; let start_idx = tracks.iter().position(|t| t.path == path).unwrap_or(0); state.queue.set_tracks(tracks, start_idx); load_current(state, current_track, resampler, paused_flag, evt_tx); } PlayerCommand::SetQueue(tracks, start_idx) => { state.is_radio = false; state.queue.set_tracks(tracks, start_idx); load_current(state, current_track, resampler, paused_flag, evt_tx); } PlayerCommand::Enqueue(path) => { state.queue.push(QueueItem::whole_file(path)); } PlayerCommand::PlayRadio { name, url } => { state.is_radio = true; state.current_track_start = Duration::ZERO; state.current_track_end = None; match open_url(&url, evt_tx.clone()) { Ok(track) => { install_track(state, current_track, resampler, paused_flag, evt_tx, track); let _ = evt_tx.send(PlayerEvent::TrackChanged(TrackInfo { path: PathBuf::from(&url), title: name, artist: "Internet Radio".to_string(), album: url, duration: None, cover: None, })); } Err(e) => { let _ = evt_tx.send(PlayerEvent::Error(format!( "failed to open radio stream: {e}" ))); state.status = PlaybackStatus::Stopped; } } } PlayerCommand::PlayPause => { state.status = match state.status { PlaybackStatus::Playing => { paused_flag.store(true, Ordering::Relaxed); PlaybackStatus::Paused } PlaybackStatus::Paused | PlaybackStatus::Stopped => { if current_track.is_some() { paused_flag.store(false, Ordering::Relaxed); PlaybackStatus::Playing } else { PlaybackStatus::Stopped } } }; let _ = evt_tx.send(PlayerEvent::StatusChanged(state.status)); } PlayerCommand::Stop => { state.status = PlaybackStatus::Stopped; paused_flag.store(true, Ordering::Relaxed); *current_track = None; state.is_radio = false; state.spectrum_buffer.clear(); state.current_track_replay_gain = None; let _ = evt_tx.send(PlayerEvent::StatusChanged(state.status)); let _ = evt_tx.send(PlayerEvent::Spectrum(vec![0.0; crate::spectrum::SPECTRUM_BANDS])); } PlayerCommand::Next => { if state.is_radio { return; // радио: следующего трека не существует } if state.queue.advance() { load_current(state, current_track, resampler, paused_flag, evt_tx); } else { state.status = PlaybackStatus::Stopped; paused_flag.store(true, Ordering::Relaxed); let _ = evt_tx.send(PlayerEvent::QueueFinished); } } PlayerCommand::Previous => { if state.is_radio { return; } if state.queue.previous() { load_current(state, current_track, resampler, paused_flag, evt_tx); } } PlayerCommand::Seek(requested_relative) => { if state.is_radio { let _ = evt_tx.send(PlayerEvent::Error( "can't seek in a live radio stream".to_string(), )); return; } let Some(track) = current_track.as_mut() else { return; }; // requested_relative — позиция от начала ТЕКУЩЕГО логического трека; переводим // в абсолютную позицию внутри физического файла и не даём выйти за его границу. let mut absolute_target = state.current_track_start + requested_relative; if let Some(end) = state.current_track_end { absolute_target = absolute_target.min(end); } match decode_seek(track, absolute_target) { Ok(actual_absolute) => { state.frames_decoded_native = (actual_absolute.as_secs_f64() * state.native_rate as f64) as u64; // Сбрасываем и буфер вывода (иначе доиграет старое аудио перед прыжком), // и ресемплер (его внутренняя история сэмплов больше не валидна). if let Ok(mut c) = consumer_shared.lock() { c.clear(); } let base_ratio = state.device_rate as f64 / state.native_rate as f64; *resampler = SpeedResampler::new(state.device_channels, base_ratio, state.speed).ok(); let _ = evt_tx.send(PlayerEvent::Position(current_position(state))); } Err(e) => { let _ = evt_tx.send(PlayerEvent::Error(format!("seek failed: {e}"))); } } } PlayerCommand::SetSpeed(speed) => { state.speed = speed; if let (Some(track), Some(_)) = (current_track.as_ref(), resampler.as_ref()) { let base_ratio = state.device_rate as f64 / track.spec.rate as f64; if let Ok(r) = SpeedResampler::new(state.device_channels, base_ratio, speed) { *resampler = Some(r); } } let _ = evt_tx.send(PlayerEvent::SpeedChanged(speed)); } PlayerCommand::ToggleShuffle => { state.queue.set_shuffle(!state.queue.shuffle); let _ = evt_tx.send(PlayerEvent::ShuffleChanged(state.queue.shuffle)); } PlayerCommand::CycleRepeat => { state.queue.cycle_repeat(); let _ = evt_tx.send(PlayerEvent::RepeatChanged(state.queue.repeat)); } PlayerCommand::SetEqBand(band, gain) => { state.eq_gains[band] = gain; if let Some(eq) = state.eq.as_mut() { eq.set_band(band, gain); } let _ = evt_tx.send(PlayerEvent::EqChanged(state.eq_gains)); } PlayerCommand::ResetEq => { state.eq_gains = [0.0; EQ_BANDS]; if let Some(eq) = state.eq.as_mut() { eq.reset(); } let _ = evt_tx.send(PlayerEvent::EqChanged(state.eq_gains)); } PlayerCommand::ToggleMute => { let new_muted = !muted_flag.load(Ordering::Relaxed); muted_flag.store(new_muted, Ordering::Relaxed); let _ = evt_tx.send(PlayerEvent::MuteChanged(new_muted)); } PlayerCommand::ToggleFadeOut => { state.fade_out_enabled = !state.fade_out_enabled; let _ = evt_tx.send(PlayerEvent::FadeOutChanged(state.fade_out_enabled)); } PlayerCommand::ToggleLoudnessNorm => { state.loudness_norm_enabled = !state.loudness_norm_enabled; let _ = evt_tx.send(PlayerEvent::LoudnessNormChanged(state.loudness_norm_enabled)); } } } /// Общая часть загрузки декодированного трека в состояние движка (используется и файлами, и радио). fn install_track( state: &mut EngineState, current_track: &mut Option<DecodedTrack>, resampler: &mut Option<SpeedResampler>, paused_flag: &Arc<AtomicBool>, evt_tx: &Sender<PlayerEvent>, track: DecodedTrack, ) { state.native_rate = track.spec.rate; state.frames_decoded_native = 0; let channels = state.device_channels; state.eq = Some(Equalizer::new(track.spec.rate as f32, channels)); for (band, gain) in state.eq_gains.iter().enumerate() { state.eq.as_mut().unwrap().set_band(band, *gain); } let base_ratio = state.device_rate as f64 / track.spec.rate as f64; *resampler = SpeedResampler::new(channels, base_ratio, state.speed).ok(); *current_track = Some(track); state.status = PlaybackStatus::Playing; paused_flag.store(false, Ordering::Relaxed); let _ = evt_tx.send(PlayerEvent::StatusChanged(state.status)); } fn load_current( state: &mut EngineState, current_track: &mut Option<DecodedTrack>, resampler: &mut Option<SpeedResampler>, paused_flag: &Arc<AtomicBool>, evt_tx: &Sender<PlayerEvent>, ) { 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) => { // Виртуальный трек (например, кусок cue-листа) — сразу прыгаем на его начало // внутри физического файла, вместо проигрывания с самого начала файла. if item.start > Duration::ZERO { if let Err(e) = decode_seek(&mut track, item.start) { let _ = evt_tx.send(PlayerEvent::Error(format!( "seek to track start failed: {e}" ))); } } state.current_track_start = item.start; // Для обычных (не-cue) файлов item.end всегда None — используем длительность // самого файла (если известна), чтобы затухание в конце работало и для них, // а не только для cue-треков. state.current_track_end = item.end.or(track.duration); state.current_track_replay_gain = item .replay_gain_db .map(|db| 10f32.powf(db / 20.0)); install_track(state, current_track, resampler, paused_flag, evt_tx, track); if item.start > Duration::ZERO { state.frames_decoded_native = (item.start.as_secs_f64() * state.native_rate as f64) as u64; } let mut info = read_tags(&item.path).unwrap_or_else(|_| TrackInfo { path: item.path.clone(), title: item .path .file_stem() .and_then(|s| s.to_str()) .unwrap_or("Unknown") .to_string(), artist: "Unknown Artist".to_string(), album: "Unknown Album".to_string(), duration: None, cover: None, }); if let Some(title) = &item.title_override { info.title = title.clone(); } if let Some(artist) = &item.artist_override { info.artist = artist.clone(); } if let Some(end) = item.end { info.duration = Some(end.saturating_sub(item.start)); } else if item.start > Duration::ZERO { // последний трек cue-листа: длительность = длительность файла минус start info.duration = info.duration.map(|d| d.saturating_sub(item.start)); } let _ = evt_tx.send(PlayerEvent::TrackChanged(info)); if let Some(idx) = queue_index { let _ = evt_tx.send(PlayerEvent::QueueIndexChanged(idx)); } } Err(e) => { let _ = evt_tx.send(PlayerEvent::Error(format!( "failed to open {}: {e}", item.path.display() ))); } } } /// Общая логика "трек закончился (физически или по границе cue) — что делать дальше". fn advance_or_finish( state: &mut EngineState, current_track: &mut Option<DecodedTrack>, resampler: &mut Option<SpeedResampler>, paused_flag: &Arc<AtomicBool>, evt_tx: &Sender<PlayerEvent>, ) { if state.is_radio { state.status = PlaybackStatus::Stopped; paused_flag.store(true, Ordering::Relaxed); *current_track = None; let _ = evt_tx.send(PlayerEvent::QueueFinished); return; } if state.queue.advance() { load_current(state, current_track, resampler, paused_flag, evt_tx); } else { state.status = PlaybackStatus::Stopped; paused_flag.store(true, Ordering::Relaxed); *current_track = None; let _ = evt_tx.send(PlayerEvent::QueueFinished); } } fn decode_and_push( state: &mut EngineState, current_track: &mut Option<DecodedTrack>, resampler: &mut Option<SpeedResampler>, producer: &mut HeapProducer<f32>, paused_flag: &Arc<AtomicBool>, evt_tx: &Sender<PlayerEvent>, spectrum_analyzer: &mut SpectrumAnalyzer, leveler: &mut Leveler, ) -> anyhow::Result<()> { let Some(track) = current_track.as_mut() else { return Ok(()); }; let packet_result = next_packet_samples(track); // Для радио сетевая ошибка (обрыв соединения) приходит как Err, а не Ok(None) — // трактуем это так же, как конец потока. let packet_result = match packet_result { Ok(v) => Ok(v), Err(e) if state.is_radio => { let _ = evt_tx.send(PlayerEvent::Error(format!("radio stream error: {e}"))); Ok(None) } Err(e) => Err(e), }; match packet_result? { Some(mut samples) => { let frames = samples.len() / track.spec.channels.count(); state.frames_decoded_native += frames as u64; if let Some(eq) = state.eq.as_mut() { eq.process_interleaved(&mut samples); } let channels = track.spec.channels.count().max(1); if state.loudness_norm_enabled { if let Some(gain) = state.current_track_replay_gain { // Точное значение из тега ReplayGain — просто статическое усиление. for s in samples.iter_mut() { *s *= gain; } } else { // Тега нет — используем адаптивный левеллер (AGC) как fallback. leveler.process(&mut samples, channels, state.native_rate as f32); } } if state.fade_out_enabled { if let Some(fade) = fade_gain(state) { if fade < 1.0 { for s in samples.iter_mut() { *s *= fade; } } } } // Копим моно-микс пост-EQ сэмплов (то, что реально слышит пользователь) для // визуализатора спектра — считаем FFT не чаще, чем раз в SPECTRUM_SEND_INTERVAL, // чтобы не заваливать канал событий чаще, чем это осмысленно для отрисовки. state .spectrum_buffer .extend(samples.chunks(channels).map(|frame| { frame.iter().sum::<f32>() / channels as f32 })); if state.spectrum_buffer.len() >= FFT_SIZE && state.last_spectrum_send.elapsed() >= SPECTRUM_SEND_INTERVAL { let start = state.spectrum_buffer.len() - FFT_SIZE; let bands = spectrum_analyzer.compute(&state.spectrum_buffer[start..], state.native_rate as f32); let _ = evt_tx.send(PlayerEvent::Spectrum(bands.to_vec())); state.spectrum_buffer.clear(); state.last_spectrum_send = Instant::now(); } // не даём буферу расти бесконечно, если по какой-то причине долго не отправляли if state.spectrum_buffer.len() > FFT_SIZE * 8 { state.spectrum_buffer.clear(); } let out = if let Some(r) = resampler.as_mut() { r.push_interleaved(&samples)? } else { samples }; for s in out { // если буфер полон, немного подождём, чтобы не терять сэмплы while producer.push(s).is_err() { std::thread::sleep(Duration::from_millis(1)); } } // Достигли конца виртуального (cue) трека внутри файла — переходим дальше, // не дожидаясь физического EOF файла. if let Some(end) = state.current_track_end { if absolute_position(state) >= end { advance_or_finish(state, current_track, resampler, paused_flag, evt_tx); } } Ok(()) } None => { // физический конец файла — переходим дальше согласно repeat/shuffle (для файлов) // либо просто останавливаемся (для радио: следующего трека не существует) if let Some(r) = resampler.as_mut() { if let Ok(tail) = r.flush() { for s in tail { let _ = producer.push(s); } } } advance_or_finish(state, current_track, resampler, paused_flag, evt_tx); Ok(()) } } } /// Для WholeAlbum: находим все аудиофайлы в той же папке, что и path, сортируем по имени /// (простая эвристика для MVP; сортировку по номеру трека из тегов добавим на уровне library). fn sibling_album_tracks(path: &PathBuf) -> anyhow::Result<Vec<PathBuf>> { let dir = path .parent() .ok_or_else(|| anyhow::anyhow!("no parent directory"))?; let mut entries: Vec<PathBuf> = std::fs::read_dir(dir)? .filter_map(|e| e.ok()) .map(|e| e.path()) .filter(|p| { p.extension() .and_then(|e| e.to_str()) .map(|ext| { matches!( ext.to_lowercase().as_str(), "mp3" | "flac" | "ogg" | "wav" | "m4a" ) }) .unwrap_or(false) }) .collect(); entries.sort(); Ok(entries) }