pub mod action; mod cmdline; pub mod command; pub mod event; pub mod input; pub(crate) mod popup; pub mod state; pub mod update; use std::io::{self, Write as _}; use std::path::{Path, PathBuf}; use std::process::{Command, Stdio}; use std::sync::Arc; use std::time::Duration; use anyhow::Result; use crokey::KeyCombination; use crossterm::event::{Event as TermEvent, EventStream, KeyEvent, KeyEventKind}; use futures_util::StreamExt; use ratatui::DefaultTerminal; use tokio::sync::mpsc; use tokio::time::MissedTickBehavior; use crate::config::keymap::{KeyResolution, Keymap}; use crate::library::Library; use crate::player; use crate::ui; use event::AppEvent; use state::AppState; use update::{Effect, update}; const TICK_INTERVAL: Duration = Duration::from_millis(250); const VISUALIZER_TICK_INTERVAL: Duration = Duration::from_millis(50); const ACTIVE_IDLE_LEASE_MS: i64 = 5 * 60 * 1000; /// Handles shared by background tasks; AppState stays pure UI data. pub struct Runtime { pub event_tx: mpsc::UnboundedSender, pub library: Arc, pub devices: Arc, pub federation: Arc, /// When the last Federation-tab status snapshot was requested. pub fed_status_at: Option, /// Coalesces urgent personal-device syncs after remote playback commands. pub device_sync_running: Arc, pub device_sync_requested: Arc, /// Placeholder ids of federated tracks being downloaded right now. pub fed_resolving: std::sync::Mutex>, /// Caps concurrent artwork loads so they never starve the disk. pub art_semaphore: Arc, /// The terminal screen was externally disturbed and needs a full repaint. pub force_redraw: bool, /// Monotonic sequence for live search; stale responses are dropped. pub search_seq: Arc, pub player: player::Controller, pub player_start_pending: bool, pub media_tx: std::sync::mpsc::Sender, pub last_media_push: Option, } fn now_epoch_seconds() -> i64 { std::time::SystemTime::now() .duration_since(std::time::UNIX_EPOCH) .map(|d| d.as_secs() as i64) .unwrap_or(0) } fn err_string(err: anyhow::Error) -> String { format!("{err:#}") } pub async fn run( mut terminal: DefaultTerminal, mut keymap: Keymap, startup_warning: Option, event_tx: mpsc::UnboundedSender, mut event_rx: mpsc::UnboundedReceiver, media_tx: std::sync::mpsc::Sender, ) -> Result<()> { let db_path = crate::library::default_db_path()?; let library = Arc::new(Library::open(&db_path)?); tracing::info!(path = %db_path.display(), "library opened"); let (settings, settings_warning) = crate::config::settings::load(); let status_message = match (startup_warning, settings_warning) { (Some(left), Some(right)) => Some(format!("{left}; {right}")), (Some(message), None) | (None, Some(message)) => Some(message), (None, None) => None, }; let mut state = AppState { status_message, ..AppState::default() }; state.player.volume = settings.volume; state.global.filters = settings.library; if let Err(err) = state.visualizer.load_library() { state.status_message = Some(format!("visualizations disabled: {err:#}")); } let devices = crate::devices::DeviceSync::new(Arc::clone(&library))?; devices.set_event_tx(event_tx.clone()); let federation = crate::federation::Federation::new(Arc::clone(&library), Arc::clone(&devices)); state.federation.settings = federation.settings(); state.federation.devices = Some(devices.status()); if let Ok((device_id, device_name)) = devices.identity_summary() { state.device_playback.self_device_id = device_id.clone(); state.device_playback.self_device_name = device_name.clone(); state.device_playback.active_device_id = Some(device_id); state.device_playback.active_device_name = Some(device_name); } let player_events = event_tx.clone(); let mut runtime = Runtime { event_tx, library, devices, federation, fed_status_at: None, device_sync_running: Arc::new(std::sync::atomic::AtomicBool::new(false)), device_sync_requested: Arc::new(std::sync::atomic::AtomicBool::new(false)), fed_resolving: std::sync::Mutex::new(std::collections::HashSet::new()), art_semaphore: Arc::new(tokio::sync::Semaphore::new(4)), force_redraw: false, search_seq: Arc::new(std::sync::atomic::AtomicU64::new(0)), player: player::spawn(move |event| { let _ = player_events.send(AppEvent::Player(event)); }), player_start_pending: false, media_tx, last_media_push: None, }; { let fed = Arc::clone(&runtime.federation); let devices = Arc::clone(&runtime.devices); let tx = runtime.event_tx.clone(); tokio::spawn(async move { fed.start_if_enabled().await; let status = fed.status().await; let _ = tx.send(AppEvent::FederationStatus(status)); let _ = tx.send(AppEvent::DeviceSyncStatus(devices.status())); }); } let mut input = EventStream::new(); let mut tick = tokio::time::interval(TICK_INTERVAL); tick.set_missed_tick_behavior(MissedTickBehavior::Skip); let mut visual_tick = tokio::time::interval(VISUALIZER_TICK_INTERVAL); visual_tick.set_missed_tick_behavior(MissedTickBehavior::Skip); loop { if runtime.force_redraw { terminal.clear()?; runtime.force_redraw = false; } terminal.draw(|frame| ui::draw(frame, &state, &keymap))?; tokio::select! { maybe_event = input.next() => match maybe_event { Some(Ok(event)) => handle_terminal_event(&mut state, &mut keymap, &mut runtime, event), Some(Err(err)) => return Err(err.into()), None => state.should_quit = true, }, Some(app_event) = event_rx.recv() => handle_app_event(&mut state, &mut runtime, app_event), _ = tick.tick() => { expire_quit_confirmation(&mut state); sync_player_shared(&mut state, &runtime); maybe_prefetch_next(&mut state, &runtime); push_media_update(&state, &mut runtime, false); } _ = visual_tick.tick(), if state.visualizer.active => { sync_player_shared(&mut state, &runtime); } } if state.should_quit { runtime.federation.shutdown().await; return Ok(()); } maintenance(&mut state, &mut runtime); } } fn sync_player_shared(state: &mut AppState, runtime: &Runtime) { if state.device_playback.is_control() { extrapolate_control_position(state); } else if state.player.current.is_some() && !runtime.player_start_pending { state.player.position_secs = runtime.player.shared.position().as_secs_f64(); state.player.paused = runtime.player.shared.paused(); } state.player.audio_analysis = if state.device_playback.is_control() { player::AudioAnalysisSnapshot::default() } else { runtime.player.shared.audio_analysis() }; publish_playback_snapshot(state, runtime); } fn unix_time_ms() -> i64 { std::time::SystemTime::now() .duration_since(std::time::UNIX_EPOCH) .map(|d| d.as_millis() as i64) .unwrap_or(0) } fn playback_repeat_to_wire(mode: state::RepeatMode) -> crate::devices::PlaybackRepeat { match mode { state::RepeatMode::Off => crate::devices::PlaybackRepeat::Off, state::RepeatMode::One => crate::devices::PlaybackRepeat::One, state::RepeatMode::All => crate::devices::PlaybackRepeat::All, } } fn playback_repeat_from_wire(mode: crate::devices::PlaybackRepeat) -> state::RepeatMode { match mode { crate::devices::PlaybackRepeat::Off => state::RepeatMode::Off, crate::devices::PlaybackRepeat::One => state::RepeatMode::One, crate::devices::PlaybackRepeat::All => state::RepeatMode::All, } } fn playback_state_from_ui(state: &AppState) -> crate::devices::PlaybackStateWire { crate::devices::PlaybackStateWire { queue: state .player .queue .iter() .map(crate::devices::PlaybackTrack::from_track) .collect(), queue_pos: state.player.queue_pos, playing: state.player.playing, paused: state.player.paused, idle_since_ms: (!state.player.playing || state.player.paused) .then_some(state.device_playback.local_idle_since_ms) .flatten(), position_secs: state.player.position_secs, volume: state.player.volume, shuffle: state.player.shuffle, repeat: playback_repeat_to_wire(state.player.repeat), } } fn playback_track_to_ui( wire: &crate::devices::PlaybackTrack, library: Option<&Library>, ) -> crate::library::models::TrackItem { if let (Some(library), Some(content_id)) = (library, wire.content_id.as_deref()) && let Ok(Some(track)) = library.track_by_content_id(content_id) { return track; } wire.to_track_item() } fn apply_playback_state_to_ui( state: &mut AppState, wire: &crate::devices::PlaybackStateWire, library: Option<&Library>, ) { state.player.queue = wire .queue .iter() .map(|track| playback_track_to_ui(track, library)) .collect(); state.player.queue_pos = wire .queue_pos .min(state.player.queue.len().saturating_sub(1)); state.player.playing = wire.playing && !state.player.queue.is_empty(); state.player.paused = wire.paused; state.device_playback.local_idle_since_ms = if state.player.playing && !state.player.paused { None } else { wire.idle_since_ms.or_else(|| Some(unix_time_ms())) }; state.player.position_secs = wire.position_secs.max(0.0); state.player.volume = wire.volume.min(100); state.player.shuffle = wire.shuffle; state.player.repeat = playback_repeat_from_wire(wire.repeat); state.player.prefetched_pos = None; state.player.original_order = None; state.player.current = state .player .playing .then(|| state.player.queue.get(state.player.queue_pos).cloned()) .flatten(); if !state.player.playing { state.player.current = None; state.player.paused = false; } state.queue_tab.cursor = state .queue_tab .cursor .min(state.player.queue.len().saturating_sub(1)); } fn publish_playback_snapshot(state: &mut AppState, runtime: &Runtime) { let Ok((device_id, device_name)) = runtime.devices.identity_summary() else { return; }; update_local_idle_since(state); state.device_playback.self_device_id = device_id.clone(); state.device_playback.self_device_name = device_name.clone(); if state.device_playback.role == state::DevicePlaybackRole::Active { state.device_playback.active_device_id = Some(device_id.clone()); state.device_playback.active_device_name = Some(device_name.clone()); } let snapshot = crate::devices::PlaybackSnapshot { device_id, device_name, active: state.device_playback.role == state::DevicePlaybackRole::Active, updated_at_ms: unix_time_ms(), state: playback_state_from_ui(state), }; runtime.devices.publish_playback(snapshot); } fn update_local_idle_since(state: &mut AppState) { if !state.player.playing || state.player.paused { if state.device_playback.local_idle_since_ms.is_none() { state.device_playback.local_idle_since_ms = Some(unix_time_ms()); } } else { state.device_playback.local_idle_since_ms = None; } } fn active_snapshot_idle_since(snapshot: &crate::devices::PlaybackSnapshot) -> Option { if snapshot.state.playing && !snapshot.state.paused { None } else { snapshot .state .idle_since_ms .or(Some(snapshot.updated_at_ms)) } } fn active_idle_lease_expired(snapshot: &crate::devices::PlaybackSnapshot, now: i64) -> bool { active_snapshot_idle_since(snapshot) .is_some_and(|idle_since| now.saturating_sub(idle_since) >= ACTIVE_IDLE_LEASE_MS) } fn extrapolate_control_position(state: &mut AppState) { let Some(snapshot) = state.device_playback.last_remote_snapshot.as_ref() else { return; }; if !snapshot.state.playing { return; } let elapsed = if snapshot.state.paused { 0.0 } else { (unix_time_ms().saturating_sub(snapshot.updated_at_ms) as f64 / 1000.0).max(0.0) }; let duration = state .player .current .as_ref() .map(|track| track.duration_seconds) .unwrap_or(0.0); let position = snapshot.state.position_secs + elapsed; state.player.position_secs = if duration > 0.0 { position.min(duration) } else { position }; } pub(crate) fn become_control_device( state: &mut AppState, runtime: &Runtime, snapshot: crate::devices::PlaybackSnapshot, ) { if state.device_playback.role == state::DevicePlaybackRole::Active { runtime.player.stop(); } state.device_playback.role = state::DevicePlaybackRole::Control; state.device_playback.active_device_id = Some(snapshot.device_id.clone()); state.device_playback.active_device_name = Some(snapshot.device_name.clone()); state.device_playback.last_remote_snapshot = Some(snapshot.clone()); state .device_playback .remote .insert(snapshot.device_id.clone(), snapshot.clone()); apply_playback_state_to_ui(state, &snapshot.state, Some(runtime.library.as_ref())); extrapolate_control_position(state); state.status_message = Some(format!("controlling {}", snapshot.device_name)); } pub(crate) fn become_active_device(state: &mut AppState, runtime: &mut Runtime, start_audio: bool) { let was_control = state.device_playback.role == state::DevicePlaybackRole::Control; state.device_playback.role = state::DevicePlaybackRole::Active; let Ok((device_id, device_name)) = runtime.devices.identity_summary() else { return; }; state.device_playback.self_device_id = device_id.clone(); state.device_playback.self_device_name = device_name.clone(); state.device_playback.active_device_id = Some(device_id); state.device_playback.active_device_name = Some(device_name); state.device_playback.last_remote_snapshot = None; state.device_playback.local_idle_since_ms = None; if was_control && start_audio && state.player.playing { start_current_audio( state, runtime, state.player.position_secs, state.player.paused, ); } publish_playback_snapshot(state, runtime); } fn record_control_playback_state(state: &mut AppState, runtime: &Runtime) { if !state.device_playback.is_control() { return; } update_local_idle_since(state); let Some(target) = state.device_playback.active_device_id.clone() else { return; }; let command = crate::devices::PlaybackCommand::SetState { state: playback_state_from_ui(state), }; if let Err(err) = runtime.devices.record_playback_command(&target, command) { tracing::warn!(%err, target, "recording playback command failed"); state.status_message = Some(format!("device command failed: {err:#}")); } else { request_urgent_device_sync(runtime); } } fn request_urgent_device_sync(runtime: &Runtime) { use std::sync::atomic::Ordering; runtime.device_sync_requested.store(true, Ordering::SeqCst); if runtime .device_sync_running .compare_exchange(false, true, Ordering::SeqCst, Ordering::SeqCst) .is_ok() { spawn_urgent_device_sync( Arc::clone(&runtime.federation), Arc::clone(&runtime.devices), runtime.event_tx.clone(), Arc::clone(&runtime.device_sync_running), Arc::clone(&runtime.device_sync_requested), ); } } fn spawn_urgent_device_sync( federation: Arc, devices: Arc, tx: mpsc::UnboundedSender, running: Arc, requested: Arc, ) { use std::sync::atomic::Ordering; tokio::spawn(async move { loop { requested.store(false, Ordering::SeqCst); match federation.device_sync_now().await { Ok(()) => {} Err(err) => { tracing::debug!("urgent device sync failed: {err:#}"); } } let _ = tx.send(AppEvent::DeviceSyncStatus(devices.status())); if !requested.swap(false, Ordering::SeqCst) { break; } } running.store(false, Ordering::SeqCst); if requested.load(Ordering::SeqCst) && running .compare_exchange(false, true, Ordering::SeqCst, Ordering::SeqCst) .is_ok() { spawn_urgent_device_sync(federation, devices, tx, running, requested); } }); } const ARTISTS_PREFETCH_MARGIN: usize = 24; /// How many artist tiles one screen holds right now (grid geometry from the /// live terminal size), so the initial load always fills the viewport. fn artist_grid_capacity() -> usize { let (width, height) = crossterm::terminal::size().unwrap_or((80, 24)); let columns = usize::from((width.saturating_sub(2) / state::TILE_WIDTH).max(1)); let rows = usize::from((height.saturating_sub(5) / state::TILE_HEIGHT).max(1)); columns * rows } /// Runs after every event: kicks off whatever background work the current /// state needs — the first artists page, the next page when the selection /// nears the end, and artwork for loaded artists. fn maintenance(state: &mut AppState, runtime: &mut Runtime) { { let global = &mut state.global; // Keep at least a full screen plus a margin loaded, and stay ahead // of the cursor: a big terminal fills itself on startup without any // scrolling, page after page. let needed = artist_grid_capacity().max(global.selected + ARTISTS_PREFETCH_MARGIN) + ARTISTS_PREFETCH_MARGIN; if global.has_more && !global.loading && !global.reloading && global.error.is_none() && global.artists.len() < needed { global.loading = true; let page = global.next_page; let hide_featured_only = global.filters.hide_featured_only; let limit = *global .page_limit .get_or_insert_with(|| (needed as i64).clamp(48, 200)); let library = Arc::clone(&runtime.library); let tx = runtime.event_tx.clone(); tokio::task::spawn_blocking(move || { let result = library .artists(page, limit, hide_featured_only) .map_err(err_string); let _ = tx.send(AppEvent::ArtistsLoaded(result)); }); } } // Liked ids load once per session — markers are shown everywhere. if !state.likes_loaded { state.likes_loaded = true; let library = Arc::clone(&runtime.library); let tx = runtime.event_tx.clone(); tokio::task::spawn_blocking(move || { let result = library.likes().map_err(err_string); let _ = tx.send(AppEvent::LikesLoaded(result)); let result = library.fed_like_ids().map_err(err_string); let _ = tx.send(AppEvent::FedLikesLoaded(result)); }); } // Playlists tab data (also wanted while the add-to-playlist picker is // open from any tab). let picker_open = matches!(state.popup, Some(state::Popup::AddToPlaylist { .. })); if state.active_tab == state::Tab::Playlists || picker_open { if state.playlists.list.is_none() { state.playlists.list = Some(state::Loadable::Loading); let library = Arc::clone(&runtime.library); let tx = runtime.event_tx.clone(); tokio::task::spawn_blocking(move || { let result = library.playlists().map_err(err_string); let _ = tx.send(AppEvent::PlaylistsLoaded(result)); }); } if let Some(opened) = state.playlists.opened { let id = opened.id; if let std::collections::hash_map::Entry::Vacant(entry) = state.playlist_views.entry(id) { entry.insert(state::Loadable::Loading); let library = Arc::clone(&runtime.library); let tx = runtime.event_tx.clone(); tokio::task::spawn_blocking(move || { let result = library.playlist(id).map_err(err_string); let _ = tx.send(AppEvent::PlaylistViewLoaded { id, result }); }); } } } // Drill-down views pushed on the stack fetch their data on first sight. for view in state.global.stack.clone() { match view { state::GlobalView::Artist { id, .. } => { if let std::collections::hash_map::Entry::Vacant(entry) = state.artist_views.entry(id) { entry.insert(state::Loadable::Loading); let library = Arc::clone(&runtime.library); let tx = runtime.event_tx.clone(); tokio::task::spawn_blocking(move || { let result = library.artist(id).map_err(err_string); let _ = tx.send(AppEvent::ArtistViewLoaded { id, result }); }); } } state::GlobalView::Release { id, .. } => { if let std::collections::hash_map::Entry::Vacant(entry) = state.release_views.entry(id) { entry.insert(state::Loadable::Loading); let library = Arc::clone(&runtime.library); let tx = runtime.event_tx.clone(); tokio::task::spawn_blocking(move || { let result = library.release(id).map_err(err_string); let _ = tx.send(AppEvent::ReleaseViewLoaded { id, result }); }); } } state::GlobalView::Search { .. } | state::GlobalView::FedArtist { .. } | state::GlobalView::FedRelease { .. } => {} } } // Refresh the Federation tab status while it is visible. if state.active_tab == state::Tab::Federation { let due = runtime .fed_status_at .is_none_or(|at| at.elapsed() > Duration::from_secs(2)); if due { runtime.fed_status_at = Some(std::time::Instant::now()); fed_spawn_status(runtime); } } // Artwork wanted by everything currently loaded, at its display size. let mut wanted: Vec<(String, u16, u16)> = Vec::new(); let tile = (state::ART_CELL_WIDTH, state::ART_CELL_HEIGHT); let header = (state::ART_HEADER_WIDTH, state::ART_HEADER_HEIGHT); for artist in &state.global.artists { if let Some(path) = &artist.image_path { wanted.push((path.clone(), tile.0, tile.1)); } } for detail in state.artist_views.values() { if let state::Loadable::Ready(detail) = detail { if let Some(path) = &detail.image_path { wanted.push((path.clone(), header.0, header.1)); } for release in &detail.releases { if let Some(path) = &release.cover_path { wanted.push((path.clone(), tile.0, tile.1)); } } } } for detail in state.release_views.values() { if let state::Loadable::Ready(detail) = detail && let Some(path) = &detail.cover_path { wanted.push((path.clone(), header.0, header.1)); } } if let Some((_, state::Loadable::Ready(card))) = &state.fed_artist_view { if let Some(path) = &card.image_path { wanted.push((path.clone(), header.0, header.1)); } for release in &card.releases { if let Some(path) = &release.cover_path { wanted.push((path.clone(), tile.0, tile.1)); wanted.push((path.clone(), header.0, header.1)); } } } for (path, width, height) in wanted { let key = crate::art::cache_key(&path, width, height); if state.art.contains_key(&key) { continue; } state.art.insert(key.clone(), state::ArtState::Loading); spawn_art_fetch(runtime, key, path, width, height); } } /// Load and decode a local image file for the art cache. fn spawn_art_fetch(runtime: &Runtime, key: String, path: String, width: u16, height: u16) { let tx = runtime.event_tx.clone(); let semaphore = Arc::clone(&runtime.art_semaphore); tokio::spawn(async move { let Ok(_permit) = semaphore.acquire_owned().await else { return; }; let art = tokio::task::spawn_blocking(move || -> anyhow::Result { let bytes = std::fs::read(&path)?; crate::art::decode_to_cells(&bytes, width, height) }) .await .map_err(anyhow::Error::from) .and_then(|result| result) .map_err(|err| tracing::warn!(%err, "artwork load failed")) .ok() .map(Arc::new); let _ = tx.send(AppEvent::ArtLoaded { key, art }); }); } /// Execute a side effect requested by update(). fn perform_effect(state: &mut AppState, runtime: &mut Runtime, effect: Effect) { if state.device_playback.is_control() && is_controlled_playback_effect(&effect) { perform_control_playback_effect(state, runtime, effect); return; } match effect { Effect::PlayCurrent => { play_current(state, runtime); push_media_metadata(state, runtime); push_media_update(state, runtime, true); } Effect::TogglePause => { if state.player.paused { runtime.player.pause(); } else { runtime.player.resume(); } push_media_update(state, runtime, true); } Effect::StopPlayback => { runtime.player_start_pending = false; runtime.player.stop(); push_media_update(state, runtime, true); } Effect::SeekBy(delta) => { let target = (state.player.position_secs + delta as f64).max(0.0); state.player.position_secs = target; runtime .player .seek(std::time::Duration::from_secs_f64(target)); } Effect::SetVolume(volume) => { runtime.player.set_volume(player::amplitude(volume)); save_app_settings(state); } Effect::SetOptions => {} Effect::PlaybackQueueChanged => {} Effect::EnqueueRelease { id, next } => { let library = Arc::clone(&runtime.library); let tx = runtime.event_tx.clone(); tokio::task::spawn_blocking(move || match library.release(id) { Ok(detail) => { let _ = tx.send(AppEvent::EnqueueTracks { tracks: detail.tracks, next, }); } Err(err) => { tracing::warn!(%err, release = id, "queueing a release failed"); let _ = tx.send(AppEvent::StatusMessage(format!("queue failed: {err:#}"))); } }); } Effect::ToggleLikes { track_ids, fed_tracks, } => { let track_ids: Vec = track_ids.into_iter().filter(|id| *id >= 0).collect(); if track_ids.is_empty() && fed_tracks.is_empty() { return; } let library = Arc::clone(&runtime.library); let devices = Arc::clone(&runtime.devices); let tx = runtime.event_tx.clone(); tokio::task::spawn_blocking(move || { for track_id in track_ids { match library.toggle_like(track_id) { Ok(liked) => { if let Err(err) = devices.record_track_like(track_id, liked) { tracing::warn!(%err, track_id, "recording synced like failed"); } let _ = tx.send(AppEvent::LikeToggled { track_id, liked }); } Err(err) => { tracing::warn!(%err, track_id, "like toggle failed"); let _ = tx.send(AppEvent::StatusMessage(format!("like failed: {err:#}"))); break; } } } for fed in fed_tracks { match library.toggle_fed_like(&fed) { Ok(liked) => { if let Err(err) = devices.record_fed_like(&fed, liked) { tracing::warn!(%err, title = %fed.title, "recording synced federated like failed"); } let _ = tx.send(AppEvent::FedLikeToggled { item_id: fed.item_id.clone(), content_id: fed.content_id.clone(), liked, }); } Err(err) => { tracing::warn!(%err, title = %fed.title, "federated like toggle failed"); let _ = tx.send(AppEvent::StatusMessage(format!("like failed: {err:#}"))); break; } } } }); } Effect::RemoveFromPlaylist { playlist_id, track_ids, content_ids, } => { let library = Arc::clone(&runtime.library); let devices = Arc::clone(&runtime.devices); let tx = runtime.event_tx.clone(); tokio::task::spawn_blocking(move || { let event = match library .remove_tracks_from_playlist(playlist_id, &track_ids) .and_then(|()| { library.remove_content_ids_from_playlist(playlist_id, &content_ids) }) { Ok(()) => { if let Err(err) = devices.record_playlist_content_removed(playlist_id, &content_ids) { tracing::warn!(%err, playlist_id, "recording synced playlist removal failed"); } AppEvent::LibraryChanged { message: Some(format!("removed {} track(s)", track_ids.len())), } } Err(err) => AppEvent::StatusMessage(format!("remove failed: {err:#}")), }; let _ = tx.send(event); }); } Effect::FedApplySettings => fed_apply_settings(state, runtime), Effect::FedSyncNow => { let fed = Arc::clone(&runtime.federation); let tx = runtime.event_tx.clone(); tokio::spawn(async move { let message = match fed.sync_now().await { Ok(()) => "federation: library published".to_string(), Err(err) => format!("federation sync failed: {err:#}"), }; let _ = tx.send(AppEvent::FederationStatus(fed.status().await)); let _ = tx.send(AppEvent::StatusMessage(message)); }); } Effect::FedShowTicket => { let fed = Arc::clone(&runtime.federation); let tx = runtime.event_tx.clone(); tokio::spawn(async move { let result = fed.ticket().await.map_err(|err| format!("{err:#}")); let _ = tx.send(AppEvent::FedTicket(result)); }); } Effect::DeviceShowInvite | Effect::DeviceConnectInvite(_) | Effect::DeviceSyncNow | Effect::DeviceSetName(_) | Effect::DeviceRevoke(_) if !state.connected_devices_enabled() => { state.status_message = Some("enable federation before using connected devices".into()); } Effect::DeviceShowInvite => { let fed = Arc::clone(&runtime.federation); let devices = Arc::clone(&runtime.devices); let tx = runtime.event_tx.clone(); tokio::spawn(async move { let result = fed.device_invite().await.map_err(|err| format!("{err:#}")); let _ = tx.send(AppEvent::DeviceInvite(result)); let _ = tx.send(AppEvent::FederationStatus(fed.status().await)); let _ = tx.send(AppEvent::DeviceSyncStatus(devices.status())); }); } Effect::DeviceConnectInvite(invite) => device_connect(runtime, invite), Effect::DeviceSyncNow => { let fed = Arc::clone(&runtime.federation); let devices = Arc::clone(&runtime.devices); let tx = runtime.event_tx.clone(); tokio::spawn(async move { let message = match fed.device_sync_now().await { Ok(()) => "devices: sync complete".to_string(), Err(err) => format!("devices: {err:#}"), }; let _ = tx.send(AppEvent::DeviceSyncStatus(devices.status())); let _ = tx.send(AppEvent::StatusMessage(message)); }); } Effect::DeviceSetName(name) => { let fed = Arc::clone(&runtime.federation); let devices = Arc::clone(&runtime.devices); let tx = runtime.event_tx.clone(); tokio::spawn(async move { let ticket = fed.ticket().await.ok(); let message = match devices.set_device_name(&name, ticket.as_deref()) { Ok(()) => "device name saved".to_string(), Err(err) => format!("device name: {err:#}"), }; let _ = tx.send(AppEvent::DeviceSyncStatus(devices.status())); let _ = tx.send(AppEvent::StatusMessage(message)); }); } Effect::DeviceRevoke(device_id) => { let devices = Arc::clone(&runtime.devices); let tx = runtime.event_tx.clone(); tokio::task::spawn_blocking(move || { let message = match devices.revoke_device(&device_id) { Ok(()) => format!("device {} revoked", &device_id[..device_id.len().min(10)]), Err(err) => format!("revoke failed: {err:#}"), }; let _ = tx.send(AppEvent::DeviceSyncStatus(devices.status())); let _ = tx.send(AppEvent::StatusMessage(message)); }); } Effect::FedOpenArtist(name) => { let fed = Arc::clone(&runtime.federation); let tx = runtime.event_tx.clone(); tokio::spawn(async move { let result = fed .artist_card(&name) .await .map_err(|err| format!("{err:#}")); let card = result.as_ref().ok().cloned(); let _ = tx.send(AppEvent::FedArtistLoaded { name: name.clone(), result, }); // Stream the artwork in after the card is on screen: the // artist image first, then every release cover. let Some(card) = card else { return }; if let Some(path) = fed.card_image(&card.owners, &name, None).await { let _ = tx.send(AppEvent::FedCardArt { name: name.clone(), release: None, path, }); } for release in &card.releases { let Some(path) = fed .card_image(&release.owners, &name, Some(&release.title)) .await else { continue; }; let _ = tx.send(AppEvent::FedCardArt { name: name.clone(), release: Some(release.title.clone()), path, }); } }); } Effect::FedDownload { tracks } => fed_download_spawn(runtime, tracks, None), Effect::FedFetchTrackInfo { tracks } => { for (placeholder_id, fed_track) in tracks { let federation = Arc::clone(&runtime.federation); let tx = runtime.event_tx.clone(); tokio::spawn(async move { let mut preview = crate::federation::pending_track(&fed_track); preview.id = placeholder_id; let item_id = fed_track.item_id.clone(); let result = federation .track_info(preview) .await .map_err(|err| format!("{err:#}")); let _ = tx.send(AppEvent::FedTrackInfoLoaded { placeholder_id, item_id, result, }); }); } } Effect::OpenVisualizerEditor { path } => match open_visualizer_editor(&path) { Ok(()) => { runtime.force_redraw = true; match state.visualizer.load_library() { Ok(()) => { clamp_settings_cursor(state); state.status_message = Some(format!("visualization script saved: {}", path.display())); } Err(err) => { state.status_message = Some(format!("visualizations: {err:#}")); } } } Err(err) => { runtime.force_redraw = true; state.status_message = Some(format!("editor failed: {err:#}")); let _ = state.visualizer.load_library(); clamp_settings_cursor(state); } }, Effect::RemoveQueueIndices { restart_paused, stop, .. } => { if stop { runtime.player_start_pending = false; runtime.player.stop(); push_media_update(state, runtime, true); } else if let Some(paused) = restart_paused { start_current_audio(state, runtime, 0.0, paused); push_media_metadata(state, runtime); push_media_update(state, runtime, true); } else { push_media_update(state, runtime, true); } } } } fn is_controlled_playback_effect(effect: &Effect) -> bool { matches!( effect, Effect::PlayCurrent | Effect::TogglePause | Effect::StopPlayback | Effect::SeekBy(_) | Effect::SetVolume(_) | Effect::SetOptions | Effect::RemoveQueueIndices { .. } | Effect::PlaybackQueueChanged ) } fn perform_control_playback_effect(state: &mut AppState, runtime: &mut Runtime, effect: Effect) { match effect { Effect::PlayCurrent => { state.player.current = state.player.queue.get(state.player.queue_pos).cloned(); state.player.playing = state.player.current.is_some(); state.player.paused = false; state.player.position_secs = 0.0; } Effect::TogglePause => {} Effect::StopPlayback => { state.player.playing = false; state.player.current = None; state.player.paused = false; state.player.position_secs = 0.0; } Effect::SeekBy(delta) => { state.player.position_secs = (state.player.position_secs + delta as f64).max(0.0); } Effect::SetVolume(volume) => { state.player.volume = volume.min(100); save_app_settings(state); } Effect::SetOptions | Effect::RemoveQueueIndices { .. } | Effect::PlaybackQueueChanged => {} _ => {} } record_control_playback_state(state, runtime); } fn clamp_settings_cursor(state: &mut AppState) { let last = state::settings_rows(state).len().saturating_sub(1); state.settings_cursor = state.settings_cursor.min(last); } fn open_visualizer_editor(path: &Path) -> Result<()> { let editor = std::env::var("VISUAL") .or_else(|_| std::env::var("EDITOR")) .unwrap_or_else(|_| "vi".to_string()); let _ = crossterm::terminal::disable_raw_mode(); let _ = crossterm::execute!( io::stdout(), crossterm::terminal::LeaveAlternateScreen, crossterm::event::DisableBracketedPaste ); let status = if cfg!(windows) { Command::new("cmd") .args(["/C", &format!("{editor} {}", path.display())]) .status() } else { Command::new("sh") .arg("-c") .arg(format!("{editor} {}", shell_quote(path))) .status() }; let _ = crossterm::execute!( io::stdout(), crossterm::terminal::EnterAlternateScreen, crossterm::event::EnableBracketedPaste ); let _ = crossterm::terminal::enable_raw_mode(); match status { Ok(status) if status.success() => Ok(()), Ok(status) => anyhow::bail!("editor exited with {status}"), Err(err) => Err(err.into()), } } fn shell_quote(path: &Path) -> String { let value = path.to_string_lossy(); format!("'{}'", value.replace('\'', "'\\''")) } fn copy_text_to_clipboard(text: &str) -> bool { let command: &[&str] = if cfg!(target_os = "macos") { &["pbcopy"] } else if cfg!(target_os = "windows") { &["clip"] } else { &["wl-copy"] }; let Some((program, args)) = command.split_first() else { return false; }; let mut child = match Command::new(program) .args(args) .stdin(Stdio::piped()) .spawn() { Ok(child) => child, Err(_) if !cfg!(target_os = "macos") && !cfg!(target_os = "windows") => { match Command::new("xclip") .args(["-selection", "clipboard"]) .stdin(Stdio::piped()) .spawn() { Ok(child) => child, Err(_) => return false, } } Err(_) => return false, }; let Some(mut stdin) = child.stdin.take() else { return false; }; if stdin.write_all(text.as_bytes()).is_err() { return false; } drop(stdin); child.wait().is_ok_and(|status| status.success()) } /// Start playing `queue[queue_pos]`: open the local file in a background /// task and hand the reader to the audio thread. fn play_current(state: &mut AppState, runtime: &mut Runtime) { start_current_audio(state, runtime, 0.0, false); } fn start_current_audio( state: &mut AppState, runtime: &mut Runtime, position_secs: f64, paused: bool, ) { let Some(track) = state.player.queue.get(state.player.queue_pos).cloned() else { return; }; // The track that was playing until now was cut short by this switch. let previous_started_at = state.player.track_started_at; let same_track_started_at = if let Some(previous) = state.player.current.take() { let same_track = previous.id == track.id; if state.player.playing && previous.id != track.id { report_history( runtime, previous.id, state.player.track_started_at, state.player.position_secs.round() as i32, false, ); } same_track.then_some(previous_started_at).flatten() } else { None }; state.player.current = Some(track.clone()); state.player.playing = true; state.player.paused = paused; state.player.position_secs = position_secs.max(0.0); state.player.audio_analysis = player::AudioAnalysisSnapshot::default(); state.player.track_started_at = same_track_started_at.or_else(|| Some(now_epoch_seconds())); state.player.prefetched_pos = None; state.status_message = Some(format!("▶ {} — {}", track.title, track.artist_line())); runtime.player_start_pending = true; if track.is_fed_pending() { // A federated track that is not on disk yet: silence the previous // audio, download it and resume through FedTrackResolved. runtime.player.stop(); state.status_message = Some(format!("federation: fetching \"{}\"…", track.title)); spawn_fed_resolve(runtime, &track); return; } let controller = runtime.player.clone(); let volume = player::amplitude(state.player.volume); let tx = runtime.event_tx.clone(); tokio::task::spawn_blocking(move || match open_track_file(&track.file_path) { Ok((reader, byte_len)) => { controller.play(reader, byte_len, volume); if position_secs > 0.0 { controller.seek(std::time::Duration::from_secs_f64(position_secs)); } if paused { controller.pause(); } } Err(err) => { tracing::warn!( track_id = track.id, title = %track.title, file = %track.file_path, %err, "cannot open track file" ); let _ = tx.send(AppEvent::Player(player::PlayerEvent::Failed(format!( "playback failed: {err}" )))); } }); } fn open_track_file(path: &str) -> std::io::Result<(player::TrackReader, Option)> { let file = std::fs::File::open(path)?; let byte_len = file.metadata().ok().map(|meta| meta.len()); Ok((std::io::BufReader::new(file), byte_len)) } /// Open the next queue item ~30s before the current track ends and append /// it in the audio thread, so rodio switches sources without a device gap. fn maybe_prefetch_next(state: &mut AppState, runtime: &Runtime) { const PREFETCH_MARGIN_SECS: f64 = 30.0; let player = &state.player; if !player.playing || player.paused || player.prefetched_pos.is_some() { return; } let Some(track) = &player.current else { return; }; if track.duration_seconds <= 0.0 || track.duration_seconds - player.position_secs > PREFETCH_MARGIN_SECS { return; } let Some(next_pos) = update::peek_next_pos(player) else { return; }; let Some(next) = player.queue.get(next_pos).cloned() else { return; }; if next.is_fed_pending() { // Download the upcoming federated track ahead of time; the gapless // enqueue happens on a later tick once it resolved to a file. spawn_fed_resolve(runtime, &next); return; } state.player.prefetched_pos = Some(next_pos); tracing::debug!(title = %next.title, "prefetching next track"); let controller = runtime.player.clone(); let tx = runtime.event_tx.clone(); tokio::task::spawn_blocking(move || match open_track_file(&next.file_path) { Ok((reader, byte_len)) => controller.enqueue(reader, byte_len), Err(err) => { tracing::warn!(%err, "prefetch failed; falling back to a normal switch"); let _ = tx.send(AppEvent::PrefetchFailed { pos: next_pos }); } }); } /// Record a finished/aborted listen in the local history; listens shorter /// than 5s are noise. fn report_history( runtime: &Runtime, track_id: i64, started_at: Option, listened: i32, completed: bool, ) { // Ephemeral federated tracks are not library rows; no history for them. if listened < 5 || track_id < 0 { return; } let library = Arc::clone(&runtime.library); tokio::task::spawn_blocking(move || { if let Err(err) = library.add_history(track_id, started_at, listened, completed) { tracing::warn!(%err, "history write failed"); } }); } /// Persist the Federation-tab settings and (re)start or stop the node. pub(crate) fn fed_apply_settings(state: &mut AppState, runtime: &Runtime) { let settings = state.federation.settings.clone(); let fed = Arc::clone(&runtime.federation); let tx = runtime.event_tx.clone(); tokio::spawn(async move { if let Err(err) = fed.apply_settings(settings).await { let _ = tx.send(AppEvent::StatusMessage(format!("federation: {err:#}"))); } let _ = tx.send(AppEvent::FederationStatus(fed.status().await)); }); } /// Connect to a peer by its pasted ticket (manual peering). pub(crate) fn fed_connect(runtime: &Runtime, ticket: String) { let fed = Arc::clone(&runtime.federation); let tx = runtime.event_tx.clone(); tokio::spawn(async move { let message = match fed.connect(&ticket).await { Ok(peer) => format!("federation: connected to {}…", &peer[..peer.len().min(10)]), Err(err) => format!("federation: {err:#}"), }; let _ = tx.send(AppEvent::FederationStatus(fed.status().await)); let _ = tx.send(AppEvent::StatusMessage(message)); }); } /// Pair with another trusted client by an opaque `frid://i/...` invite. pub(crate) fn device_connect(runtime: &Runtime, invite: String) { let fed = Arc::clone(&runtime.federation); let devices = Arc::clone(&runtime.devices); let tx = runtime.event_tx.clone(); tokio::spawn(async move { let _ = tx.send(AppEvent::StatusMessage( "waiting for device confirmation...".to_string(), )); let result = fed .device_connect(&invite) .await .map_err(|err| format!("{err:#}")); let _ = tx.send(AppEvent::FederationStatus(fed.status().await)); let _ = tx.send(AppEvent::DeviceSyncStatus(devices.status())); let _ = tx.send(AppEvent::DeviceConnectResult(result)); }); } /// Downloads one pending federated track (into the cache, or the library /// when save-on-listen is enabled) and reports back with the placeholder id /// so the queue can swap the resolved track in. fn spawn_fed_resolve(runtime: &Runtime, track: &crate::library::models::TrackItem) { let Some(fed_track) = track.fed.clone() else { return; }; { let mut resolving = runtime .fed_resolving .lock() .unwrap_or_else(std::sync::PoisonError::into_inner); if !resolving.insert(track.id) { return; } } let placeholder_id = track.id; let fed = Arc::clone(&runtime.federation); let tx = runtime.event_tx.clone(); tokio::spawn(async move { let result = fed .prepare_playback(&fed_track) .await .map(Box::new) .map_err(|err| format!("{err:#}")); let _ = tx.send(AppEvent::FedTrackResolved { placeholder_id, result, }); }); } /// Downloads federated tracks into the library one by one (with progress in /// the status bar) and optionally links them to a playlist afterwards. pub(crate) fn fed_download_spawn( runtime: &Runtime, tracks: Vec, playlist: Option<(i64, String)>, ) { if tracks.is_empty() { return; } let fed = Arc::clone(&runtime.federation); let library = Arc::clone(&runtime.library); let devices = Arc::clone(&runtime.devices); let tx = runtime.event_tx.clone(); tokio::spawn(async move { let total = tracks.len(); let mut imported_ids = Vec::new(); let mut imported_fed_tracks = Vec::new(); let mut failed = 0usize; for (index, track) in tracks.iter().enumerate() { let _ = tx.send(AppEvent::StatusMessage(format!( "federation: downloading {}/{total}: {}", index + 1, track.title ))); match fed.download_to_library(track).await { Ok(imported) => { imported_ids.push(imported.id); imported_fed_tracks.push(track.clone()); } Err(err) => { failed += 1; tracing::warn!(title = %track.title, "federated download failed: {err:#}"); } } } let mut message = format!("federation: downloaded {} of {total}", imported_ids.len()); if failed > 0 { message.push_str(&format!(" ({failed} failed)")); } if let Some((playlist_id, playlist_title)) = playlist && !imported_ids.is_empty() { let library = Arc::clone(&library); let devices = Arc::clone(&devices); let tx_add = tx.clone(); let title = playlist_title.clone(); tokio::task::spawn_blocking(move || { let result = library .add_tracks_to_playlist(playlist_id, &imported_ids) .map_err(|err| format!("{err:#}")); if result.is_ok() && let Err(err) = devices.record_playlist_fed_tracks_added(playlist_id, &imported_fed_tracks) { tracing::warn!(%err, playlist_id, "recording synced playlist add failed"); } let _ = tx_add.send(AppEvent::PlaylistTracksAdded { playlist_id, playlist_title: title, result, }); }); } let _ = tx.send(AppEvent::LibraryChanged { message: Some(message), }); }); } /// Request a fresh status snapshot for the Federation tab. fn fed_spawn_status(runtime: &Runtime) { let fed = Arc::clone(&runtime.federation); let devices = Arc::clone(&runtime.devices); let tx = runtime.event_tx.clone(); tokio::spawn(async move { let _ = tx.send(AppEvent::FederationStatus(fed.status().await)); let _ = tx.send(AppEvent::DeviceSyncStatus(devices.status())); }); } /// Clear a timed-out quit confirmation and its status-bar hint. fn expire_quit_confirmation(state: &mut AppState) { if state .quit_armed_until .is_some_and(|deadline| std::time::Instant::now() > deadline) { state.quit_armed_until = None; if state.status_message.as_deref() == Some(update::QUIT_CONFIRM_HINT) { state.status_message = None; } } } fn handle_terminal_event( state: &mut AppState, keymap: &mut Keymap, runtime: &mut Runtime, event: TermEvent, ) { match event { TermEvent::Key(key) => { // Kitty-enhanced terminals and Windows also deliver Release // events; acting on them would double-fire every binding. if !matches!(key.kind, KeyEventKind::Press | KeyEventKind::Repeat) { return; } if state.popup.is_some() { popup::handle_key(state, runtime, key); } else if state.cmdline.active { cmdline::handle_key(state, runtime, key); } else { handle_main_key(state, keymap, runtime, key); } } TermEvent::Paste(pasted) => { if state.popup.is_some() { popup::handle_paste(state, &pasted); } else if state.cmdline.active { cmdline::handle_paste(state, runtime, &pasted); } } _ => {} } } fn handle_main_key( state: &mut AppState, keymap: &mut Keymap, runtime: &mut Runtime, key: KeyEvent, ) { let combo = KeyCombination::from(key); match keymap.resolve(combo, state.active_tab.key_context()) { KeyResolution::Action(action) => { state.pending_keys = None; // trace, not debug: on the Logs tab every keypress would // otherwise append a line and pollute what's being read. tracing::trace!(?action, "key resolved"); if let Some(effect) = update(state, action) { perform_effect(state, runtime, effect); } } KeyResolution::Pending(keys) => state.pending_keys = Some(keys), KeyResolution::Unmatched => state.pending_keys = None, } } /// `:import ` — run a library import in the background, reporting /// progress into the status bar. pub(super) fn spawn_import(state: &mut AppState, runtime: &Runtime, path: &str) { let expanded = expand_tilde(path); state.status_message = Some(format!("importing {}…", expanded.display())); let library = Arc::clone(&runtime.library); let tx = runtime.event_tx.clone(); tokio::task::spawn_blocking(move || { let progress_tx = tx.clone(); let mut last_refresh = std::time::Instant::now(); let result = crate::library::import::import_path(&library, &expanded, |done, total, name| { let _ = progress_tx.send(AppEvent::ImportProgress { done, total, current: name.to_string(), }); // Long imports show up in the views as they go, not only at // the end: refresh about once a second. if done < total && last_refresh.elapsed() >= Duration::from_secs(1) { last_refresh = std::time::Instant::now(); let _ = progress_tx.send(AppEvent::LibraryChanged { message: None }); } }); let event = match result { Ok(outcome) => { for (file, error) in &outcome.failed { tracing::warn!(file = %file.display(), error, "file was not imported"); } AppEvent::LibraryChanged { message: Some(outcome.summary()), } } Err(err) => AppEvent::StatusMessage(format!("import failed: {err:#}")), }; let _ = tx.send(event); }); } fn expand_tilde(path: &str) -> PathBuf { if let Some(rest) = path.strip_prefix("~/") && let Some(home) = std::env::home_dir() { return home.join(rest); } PathBuf::from(path) } /// The library changed (import, edit, delete): reload everything that is /// currently on screen or cached, replacing data in place so nothing /// flashes "loading" while the user keeps browsing. Runs repeatedly during /// long imports, so every step must be cheap and non-disruptive. fn on_library_changed(state: &mut AppState, runtime: &mut Runtime) { refresh_artists(state, runtime); // Refresh every cached drill-down view in place; the handlers replace // the entries when the fresh data arrives. for id in state.artist_views.keys().copied().collect::>() { let library = Arc::clone(&runtime.library); let tx = runtime.event_tx.clone(); tokio::task::spawn_blocking(move || { let result = library.artist(id).map_err(err_string); let _ = tx.send(AppEvent::ArtistViewLoaded { id, result }); }); } for id in state.release_views.keys().copied().collect::>() { let library = Arc::clone(&runtime.library); let tx = runtime.event_tx.clone(); tokio::task::spawn_blocking(move || { let result = library.release(id).map_err(err_string); let _ = tx.send(AppEvent::ReleaseViewLoaded { id, result }); }); } for id in state.playlist_views.keys().copied().collect::>() { let library = Arc::clone(&runtime.library); let tx = runtime.event_tx.clone(); tokio::task::spawn_blocking(move || { let result = library.playlist(id).map_err(err_string); let _ = tx.send(AppEvent::PlaylistViewLoaded { id, result }); }); } if state.playlists.list.is_some() { let library = Arc::clone(&runtime.library); let tx = runtime.event_tx.clone(); tokio::task::spawn_blocking(move || { let result = library.playlists().map_err(err_string); let _ = tx.send(AppEvent::PlaylistsLoaded(result)); }); } // Likes reload on the next maintenance pass; the old set stays visible // until then. state.likes_loaded = false; // Fresh copies of whatever sits in the queue. Federated placeholders // and ephemeral tracks (negative ids) are not library rows and keep // their in-memory copies. let ids: Vec = state .player .queue .iter() .map(|track| track.id) .filter(|id| *id >= 0) .collect(); if !ids.is_empty() { let library = Arc::clone(&runtime.library); let tx = runtime.event_tx.clone(); tokio::task::spawn_blocking(move || match library.tracks_by_ids(&ids) { Ok(tracks) => { let _ = tx.send(AppEvent::QueueTracksRefreshed { tracks }); } Err(err) => tracing::warn!(%err, "queue refresh failed"), }); } // A live search view shows stale rows now; run the query again. if state .global .stack .iter() .any(|view| matches!(view, state::GlobalView::Search { .. })) && !state.search.query.is_empty() { cmdline::schedule_search(state, runtime); } } /// Reload the artist grid atomically: fetch everything that is loaded now /// as one page and swap it in when it arrives, so the grid never shows an /// empty "loading" state in between. Pagination then continues from page 2. fn refresh_artists(state: &mut AppState, runtime: &Runtime) { let global = &mut state.global; let needed = artist_grid_capacity() + ARTISTS_PREFETCH_MARGIN; let limit = (global.artists.len().max(needed) as i64).clamp(48, 1000); let hide_featured_only = global.filters.hide_featured_only; global.reloading = true; let library = Arc::clone(&runtime.library); let tx = runtime.event_tx.clone(); tokio::task::spawn_blocking(move || { let event = match library.artists(1, limit, hide_featured_only) { Ok(page) => AppEvent::ArtistsReloaded { page, limit }, Err(err) => AppEvent::ArtistsLoaded(Err(err_string(err))), }; let _ = tx.send(event); }); } fn save_app_settings(state: &AppState) { let settings = crate::config::settings::AppSettings { volume: state.player.volume, library: state.global.filters, }; if let Err(err) = crate::config::settings::save(&settings) { tracing::warn!(%err, "saving app settings failed"); } } fn reset_artist_pagination(state: &mut AppState) { let global = &mut state.global; global.artists.clear(); global.total = 0; global.has_more = true; global.next_page = 1; global.loading = false; global.error = None; global.selected = 0; global.page_limit = None; global.reloading = false; } /// Swap queue entries for their fresh library copies; tracks that were /// deleted leave the queue. fn apply_queue_refresh( state: &mut AppState, runtime: &mut Runtime, tracks: Vec, ) { let by_id: std::collections::HashMap = tracks.into_iter().map(|track| (track.id, track)).collect(); let current_id = state.player.current.as_ref().map(|track| track.id); for track in &mut state.player.queue { if let Some(fresh) = by_id.get(&track.id) { *track = fresh.clone(); } } let had_missing = state .player .queue .iter() .any(|track| track.id >= 0 && !by_id.contains_key(&track.id)); if had_missing { state .player .queue .retain(|track| track.id < 0 || by_id.contains_key(&track.id)); state.player.prefetched_pos = None; } if state.player.queue.is_empty() { if state.player.playing { runtime.player_start_pending = false; runtime.player.stop(); } state.player = state::PlayerBar::default(); state.queue_tab.cursor = 0; return; } state.queue_tab.cursor = state.queue_tab.cursor.min(state.player.queue.len() - 1); match current_id { Some(id) if by_id.contains_key(&id) => { if let Some(position) = state.player.queue.iter().position(|track| track.id == id) { state.player.queue_pos = position; } state.player.current = by_id.get(&id).cloned(); push_media_metadata(state, runtime); } Some(_) => { // The playing track was deleted from the library. runtime.player_start_pending = false; runtime.player.stop(); state.player.queue_pos = state.player.queue_pos.min(state.player.queue.len() - 1); state.player.current = None; state.player.playing = false; state.player.paused = false; push_media_update(state, runtime, true); } None => { state.player.queue_pos = state.player.queue_pos.min(state.player.queue.len() - 1); } } } fn handle_device_playback_snapshot( state: &mut AppState, runtime: &mut Runtime, snapshot: crate::devices::PlaybackSnapshot, ) { if snapshot.device_id == state.device_playback.self_device_id { return; } state .device_playback .remote .insert(snapshot.device_id.clone(), snapshot.clone()); let now = unix_time_ms(); state.device_playback.online_devices = state .device_playback .remote .values() .filter(|snapshot| { now.saturating_sub(snapshot.updated_at_ms) <= state::DEVICE_ONLINE_TTL_MS }) .count() + 1; if !snapshot.active { return; } let lease_expired = active_idle_lease_expired(&snapshot, now); let already_controls_this_device = state.device_playback.is_control() && state.device_playback.active_device_id.as_deref() == Some(snapshot.device_id.as_str()); if !lease_expired || already_controls_this_device { let was_active = state.device_playback.role == state::DevicePlaybackRole::Active; let was_paused = state.player.playing && state.player.paused; become_control_device(state, runtime, snapshot.clone()); if was_active && was_paused { state.popup = Some(state::Popup::FedText { title: "Active device moved".to_string(), text: format!("Playback is now controlled by {}.", snapshot.device_name), }); } return; } if state.device_playback.is_control() { return; } become_active_device(state, runtime, false); state.status_message = Some(format!( "active playback moved here; {} was idle for 5m", snapshot.device_name )); } fn handle_playback_command( state: &mut AppState, runtime: &mut Runtime, command: crate::devices::PlaybackCommand, ) { match command { crate::devices::PlaybackCommand::SetState { state: wire } => { let old_current_id = state.player.current.as_ref().map(|track| track.id); let old_playing = state.player.playing; let old_paused = state.player.paused; become_active_device(state, runtime, false); apply_playback_state_to_ui(state, &wire, Some(runtime.library.as_ref())); runtime .player .set_volume(player::amplitude(state.player.volume)); if !state.player.playing { runtime.player_start_pending = false; runtime.player.stop(); push_media_update(state, runtime, true); publish_playback_snapshot(state, runtime); return; } let current_id = state.player.current.as_ref().map(|track| track.id); if !old_playing || old_current_id != current_id { start_current_audio( state, runtime, state.player.position_secs, state.player.paused, ); push_media_metadata(state, runtime); } else { runtime.player.seek(std::time::Duration::from_secs_f64( state.player.position_secs, )); if state.player.paused && !old_paused { runtime.player.pause(); } else if !state.player.paused && old_paused { runtime.player.resume(); } } push_media_update(state, runtime, true); publish_playback_snapshot(state, runtime); } } } fn handle_app_event(state: &mut AppState, runtime: &mut Runtime, event: AppEvent) { match event { AppEvent::StatusMessage(message) => state.status_message = Some(message), AppEvent::FederationStatus(status) => { state.federation.status = Some(status); } AppEvent::DeviceSyncStatus(status) => { state.federation.devices = Some(status); let now = unix_time_ms(); state.device_playback.online_devices = state .federation .devices .as_ref() .map(|status| { status .devices .iter() .filter(|device| { state::device_presence_section(state, device, now) == state::DevicePresenceSection::Online }) .count() }) .unwrap_or(1) .max(1); let active_revoked = state.device_playback.is_control() && state .device_playback .active_device_id .as_ref() .is_some_and(|active| { state .federation .devices .as_ref() .and_then(|status| { status .devices .iter() .find(|device| device.device_id == *active) }) .is_some_and(|device| device.revoked) }); if active_revoked { become_active_device(state, runtime, false); state.status_message = Some("active playback moved to this device".into()); } clamp_settings_cursor(state); } AppEvent::DeviceInvite(result) => match result { Ok(invite) => { let copied = copy_text_to_clipboard(&invite); state.popup = Some(state::Popup::FedText { title: "Device invite".to_string(), text: invite, }); state.status_message = Some(if copied { "device invite copied to clipboard".to_string() } else { "device invite generated".to_string() }); } Err(message) => state.status_message = Some(format!("device invite: {message}")), }, AppEvent::DeviceConnectResult(result) => match result { Ok(message) => state.status_message = Some(message), Err(message) => state.status_message = Some(format!("connect failed: {message}")), }, AppEvent::DevicePairingRequest(request) => { state.popup = Some(state::Popup::DevicePairing { request_id: request.request_id, device_id: request.device_id, name: request.name, client_version: request.client_version, requester_group_id: request.requester_group_id, requester_group_active_devices: request.requester_group_active_devices, }); state.federation.devices = Some(runtime.devices.status()); } AppEvent::DevicePlayback(snapshot) => { handle_device_playback_snapshot(state, runtime, snapshot); } AppEvent::PlaybackCommand(command) => { handle_playback_command(state, runtime, command); } AppEvent::FedSearchLoaded { seq, result } => { if runtime.search_seq.load(std::sync::atomic::Ordering::SeqCst) != seq { return; } state.search.fed_loading = false; match result { Ok(results) => { state.search.fed_artists = results.artists; state.search.fed_tracks = results.tracks; } Err(message) => tracing::warn!(%message, "federated search failed"), } } AppEvent::FedTrackResolved { placeholder_id, result, } => { runtime .fed_resolving .lock() .unwrap_or_else(std::sync::PoisonError::into_inner) .remove(&placeholder_id); if state.device_playback.is_control() { tracing::debug!( placeholder_id, "ignored local federated track resolution while controlling remote playback" ); return; } match result { Ok(playable) => { if playable.imported { // Save-on-listen imported the file; refresh the // library views through the standard change path. let _ = runtime.event_tx.send(AppEvent::LibraryChanged { message: Some(format!( "saved \"{}\" to the library", playable.track.title )), }); } let resolved = playable.track.clone(); // Swap the placeholder for the real track everywhere it // sits in the queue. for slot in &mut state.player.queue { if slot.id == placeholder_id { *slot = resolved.clone(); } } let waiting = state .player .current .as_ref() .is_some_and(|current| current.id == placeholder_id); if waiting { // Playback was parked on this track; start it now. let paused = state.player.paused; start_current_audio(state, runtime, 0.0, paused); push_media_metadata(state, runtime); push_media_update(state, runtime, true); } } Err(message) => { state.status_message = Some(format!("federation: {message}")); let waiting = state .player .current .as_ref() .is_some_and(|current| current.id == placeholder_id); if waiting { // Skip the failed track instead of stalling the queue. if state.player.queue_pos + 1 < state.player.queue.len() { state.player.queue_pos += 1; start_current_audio(state, runtime, 0.0, state.player.paused); push_media_metadata(state, runtime); push_media_update(state, runtime, true); } else { runtime.player_start_pending = false; state.player.playing = false; state.player.current = None; runtime.player.stop(); } } } } } AppEvent::FedTrackInfoLoaded { placeholder_id, item_id, result, } => match result { Ok(enriched) => { if let Some(state::Popup::TrackInfo { tracks, .. }) = &mut state.popup && let Some(slot) = tracks.iter_mut().find(|track| { track.id == placeholder_id && track.fed.as_ref().is_some_and(|fed| fed.item_id == item_id) }) { *slot = enriched; } } Err(message) => { state.status_message = Some(format!("federation metadata: {message}")); } }, AppEvent::FedArtistLoaded { name, result } => { if let Some((current, data)) = &mut state.fed_artist_view && *current == name { *data = match result { Ok(card) => state::Loadable::Ready(card), Err(message) => state::Loadable::Failed(message), }; } } AppEvent::FedCardArt { name, release, path, } => { if let Some((current, state::Loadable::Ready(card))) = &mut state.fed_artist_view && *current == name { match release { None => card.image_path = Some(path), Some(title) => { if let Some(slot) = card.releases.iter_mut().find(|r| r.title == title) { slot.cover_path = Some(path); } } } } } AppEvent::FedTicket(result) => match result { Ok(ticket) => { state.popup = Some(state::Popup::FedText { title: "Federation ticket (share with a peer)".to_string(), text: ticket, }); } Err(message) => state.status_message = Some(message), }, AppEvent::ArtistsLoaded(Ok(page)) => { let global = &mut state.global; if global.reloading { // A stale page of the pre-change pagination; the pending // ArtistsReloaded swap supersedes it. return; } global.loading = false; global.total = page.total; global.has_more = page.has_more; global.next_page = page.page + 1; global.artists.extend(page.items); if !global.has_more && !global.artists.is_empty() { global.selected = global.selected.min(global.artists.len() - 1); } } AppEvent::ArtistsLoaded(Err(message)) => { tracing::warn!(%message, "artists page load failed"); state.global.reloading = false; state.global.loading = false; state.global.error = Some(message.clone()); state.status_message = Some(message); } AppEvent::ArtistsReloaded { page, limit } => { let global = &mut state.global; global.reloading = false; global.loading = false; global.error = None; global.total = page.total; global.has_more = page.has_more; global.next_page = 2; global.page_limit = Some(limit); global.artists = page.items; if !global.artists.is_empty() { global.selected = global.selected.min(global.artists.len() - 1); } else { global.selected = 0; } } AppEvent::ArtistViewLoaded { id, result } => { let entry = match result { Ok(detail) => state::Loadable::Ready(detail), Err(message) => { tracing::warn!(artist = id, %message, "artist view load failed"); state::Loadable::Failed(message) } }; state.artist_views.insert(id, entry); } AppEvent::ReleaseViewLoaded { id, result } => { let entry = match result { Ok(detail) => state::Loadable::Ready(detail), Err(message) => { tracing::warn!(release = id, %message, "release view load failed"); state::Loadable::Failed(message) } }; state.release_views.insert(id, entry); // A Shift-J jump was waiting for this release: focus its track. if let Some((release_id, track_id)) = state.pending_release_focus && release_id == id { state.pending_release_focus = None; if let Some(state::Loadable::Ready(detail)) = state.release_views.get(&id) { let position = detail .tracks .iter() .position(|t| t.id == track_id) .unwrap_or(0); if let Some(state::GlobalView::Release { id: top, cursor }) = state.global.stack.last_mut() && *top == release_id { *cursor = position; } } } } AppEvent::SearchLoaded { seq, result } => { if seq != runtime.search_seq.load(std::sync::atomic::Ordering::SeqCst) { return; } state.search.loading = false; match result { Ok(results) => state.search.results = Some(results), Err(message) => state.status_message = Some(message), } } AppEvent::ArtLoaded { key, art } => { let entry = match art { Some(image) => state::ArtState::Ready(image), None => state::ArtState::Failed, }; state.art.insert(key, entry); } AppEvent::Player(event) if state.device_playback.is_control() => { runtime.player_start_pending = false; tracing::debug!( ?event, "ignored local player event while controlling remote playback" ); } AppEvent::Player(player::PlayerEvent::Started) => { runtime.player_start_pending = false; } AppEvent::Player(player::PlayerEvent::TrackFinished { has_next }) => { // The finished track gets a full-duration, completed entry. if let Some(finished) = state.player.current.clone() { report_history( runtime, finished.id, state.player.track_started_at, finished.duration_seconds.round() as i32, true, ); } if has_next { // A prefetched source is already playing; just realign state. let next_pos = state .player .prefetched_pos .take() .unwrap_or(state.player.queue_pos + 1); state.player.queue_pos = next_pos.min(state.player.queue.len().saturating_sub(1)); state.player.current = state.player.queue.get(state.player.queue_pos).cloned(); state.player.position_secs = 0.0; state.player.track_started_at = Some(now_epoch_seconds()); push_media_metadata(state, runtime); push_media_update(state, runtime, true); } else { state.player.current = None; state.player.prefetched_pos = None; if let Some(effect) = update::advance_after_finish(state) { perform_effect(state, runtime, effect); } } } AppEvent::Player(player::PlayerEvent::Failed(message)) => { runtime.player_start_pending = false; tracing::error!(%message, "playback failed"); state.player.playing = false; state.player.paused = false; state.status_message = Some(message); } AppEvent::PrefetchFailed { pos } => { if state.device_playback.is_control() { return; } if state.player.prefetched_pos == Some(pos) { state.player.prefetched_pos = None; } } AppEvent::PlaylistsLoaded(result) => { state.playlists.list = Some(match result { Ok(list) => { state.playlists.selected = state.playlists.selected.min(list.len().saturating_sub(1)); state::Loadable::Ready(list) } Err(message) => { tracing::warn!(%message, "playlists load failed"); state::Loadable::Failed(message) } }); } AppEvent::PlaylistViewLoaded { id, result } => { let entry = match result { Ok(detail) => state::Loadable::Ready(detail), Err(message) => state::Loadable::Failed(message), }; state.playlist_views.insert(id, entry); } AppEvent::LikesLoaded(result) => match result { Ok(ids) => { state.likes = ids.into_iter().collect(); } Err(message) => tracing::warn!(%message, "likes load failed"), }, AppEvent::FedLikesLoaded(result) => match result { Ok(ids) => state.fed_likes = ids.into_iter().collect(), Err(message) => tracing::warn!(%message, "federated likes load failed"), }, AppEvent::FedLikeToggled { item_id, content_id, liked, } => { if liked { state.fed_likes.insert(item_id); if let Some(content_id) = content_id.and_then(|id| music_dht::normalize_content_id(&id)) { state.fed_likes.insert(content_id); } } else { state.fed_likes.remove(&item_id); if let Some(content_id) = content_id.and_then(|id| music_dht::normalize_content_id(&id)) { state.fed_likes.remove(&content_id); } } // The virtual Likes playlist is stale now; refetch on next open. state.playlist_views.remove(&state::LIKES_PLAYLIST_ID); state.playlists.list = None; state.status_message = Some(if liked { "♥ liked (federation)".to_string() } else { "like removed".to_string() }); } AppEvent::LikeToggled { track_id, liked } => { if liked { state.likes.insert(track_id); } else { state.likes.remove(&track_id); } // The virtual Likes playlist is stale now; refetch on next open. state.playlist_views.remove(&state::LIKES_PLAYLIST_ID); state.playlists.list = None; state.status_message = Some(if liked { "♥ liked".to_string() } else { "like removed".to_string() }); } AppEvent::EnqueueTracks { tracks, next } => { let count = tracks.len(); update::enqueue_tracks(state, tracks, next); record_control_playback_state(state, runtime); state.status_message = Some(if next { format!("{count} tracks queued next") } else { format!("{count} tracks queued") }); } AppEvent::PlaylistCreated { result, add_target } => match result { Ok(playlist) => { tracing::info!(title = %playlist.title, "playlist created"); state.status_message = Some(format!("playlist \"{}\" created", playlist.title)); state.popup = None; // The list is stale; refetch when next needed. state.playlists.list = None; if let Some(target) = add_target { popup::spawn_add_target(runtime, playlist.id, playlist.title.clone(), target); } } Err(message) => { tracing::warn!(%message, "playlist creation failed"); state.status_message = Some(format!("create failed: {message}")); if let Some(state::Popup::NewPlaylist { busy, .. }) = &mut state.popup { *busy = false; } } }, AppEvent::PlaylistTracksAdded { playlist_id, playlist_title, result, } => { state.popup = None; match result { Ok(()) => { state.status_message = Some(format!("added to \"{playlist_title}\"")); // Counts and contents changed; refetch lazily. state.playlist_views.remove(&playlist_id); state.playlists.list = None; } Err(message) => { tracing::warn!(%message, playlist_id, "adding to playlist failed"); state.status_message = Some(format!("add failed: {message}")); } } } AppEvent::LibraryChanged { message } => { on_library_changed(state, runtime); if let Some(message) = message { state.status_message = Some(message); } } AppEvent::ImportProgress { done, total, current, } => { state.status_message = Some(format!("importing {done}/{total}: {current}")); } AppEvent::QueueTracksRefreshed { tracks } => { if state.device_playback.is_control() { return; } apply_queue_refresh(state, runtime, tracks); } AppEvent::Media(command) => { use crate::media::MediaCommand; tracing::debug!(?command, "media key"); let action = match command { MediaCommand::TogglePause => action::Action::PlayPause, MediaCommand::Play if state.player.paused || state.player.current.is_none() => { action::Action::PlayPause } MediaCommand::Play => return, MediaCommand::Pause if state.player.current.is_some() && !state.player.paused => { action::Action::PlayPause } MediaCommand::Pause => return, MediaCommand::Next => action::Action::NextTrack, MediaCommand::Previous => action::Action::PrevTrack, MediaCommand::Stop => { runtime.player_start_pending = false; state.player.playing = false; state.player.paused = false; state.player.current = None; if state.device_playback.is_control() { record_control_playback_state(state, runtime); } else { runtime.player.stop(); push_media_update(state, runtime, true); } return; } }; if let Some(effect) = update(state, action) { perform_effect(state, runtime, effect); } } } } /// Mirror the playback state to the OS now-playing surface. `force` skips /// the position throttle (track switches, pauses). fn push_media_update(state: &AppState, runtime: &mut Runtime, force: bool) { use crate::media::MediaUpdate; const POSITION_INTERVAL: Duration = Duration::from_secs(2); if !force && runtime .last_media_push .is_some_and(|at| at.elapsed() < POSITION_INTERVAL) { return; } runtime.last_media_push = Some(std::time::Instant::now()); let player = &state.player; if !player.playing { let _ = runtime.media_tx.send(MediaUpdate::Stopped); return; } let _ = runtime.media_tx.send(MediaUpdate::Playback { playing: player.playing, paused: player.paused, position_secs: player.position_secs, }); } fn push_media_metadata(state: &AppState, runtime: &Runtime) { use crate::media::MediaUpdate; if let Some(track) = &state.player.current { let _ = runtime.media_tx.send(MediaUpdate::Metadata { title: track.title.clone(), artist: track.artist_line(), album: track.release_title.clone(), duration_secs: track.duration_seconds, }); } }