diff --git a/Cargo.lock b/Cargo.lock index 1e30f69..d9f705b 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -524,9 +524,9 @@ dependencies = [ [[package]] name = "cc" -version = "1.3.0" +version = "1.4.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "c89588d05638b5b4594a3348a2d6c20277e43a7f5c5202b05cc56888475a47b8" +checksum = "5add81bb678e6cb321aff7fa0dc7689ad82b112dbc032cea19f91d6b8e3582b9" dependencies = [ "find-msvc-tools", "shlex", @@ -862,9 +862,9 @@ checksum = "790eea4361631c5e7d22598ecd5723ff611904e3344ce8720784c93e3d83d40b" [[package]] name = "crokey" -version = "1.4.0" +version = "1.5.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "04a63daf06a168535c74ab97cdba3ed4fa5d4f32cb36e437dcceb83d66854b7c" +checksum = "074bc31bff1084f74e5a145ad5f6ac74469fa312159faa5144f0b1c4506b8d80" dependencies = [ "crokey-proc_macros", "crossterm", @@ -875,9 +875,9 @@ dependencies = [ [[package]] name = "crokey-proc_macros" -version = "1.4.0" +version = "1.5.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "847f11a14855fc490bd5d059821895c53e77eeb3c2b73ee3dded7ce77c93b231" +checksum = "b88c4ff2d5717616c3bb4a31d06b0b241ed811c7c0c41d077d77631bee7caad0" dependencies = [ "crossterm", "proc-macro2", @@ -1275,9 +1275,9 @@ dependencies = [ [[package]] name = "either" -version = "1.16.0" +version = "1.17.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "91622ff5e7162018101f2fea40d6ebf4a78bbe5a49736a2020649edf9693679e" +checksum = "9e5e8f6c15a24b9a3ee5efec809ccd006d3b30e8b3bb63c39af737c7f87daa1d" [[package]] name = "embedded-io" @@ -1456,7 +1456,7 @@ dependencies = [ [[package]] name = "federation-net" version = "0.1.0" -source = "git+https://gt.hexor.cy/ab/frid.git#512a818a6a52ec713678e9a4e1cf0f50bb1e34ab" +source = "git+https://gt.hexor.cy/ab/frid.git#085a4752da25a8d3fe7eac673af081a4a73c08bd" dependencies = [ "blake3", "data-encoding", @@ -2990,7 +2990,7 @@ dependencies = [ [[package]] name = "music-dht" version = "0.1.0" -source = "git+https://gt.hexor.cy/ab/frid.git#512a818a6a52ec713678e9a4e1cf0f50bb1e34ab" +source = "git+https://gt.hexor.cy/ab/frid.git#085a4752da25a8d3fe7eac673af081a4a73c08bd" dependencies = [ "async-trait", "blake3", @@ -3001,6 +3001,7 @@ dependencies = [ "rand 0.9.5", "rusqlite", "serde", + "serde_json", "thiserror 2.0.19", "tokio", "tracing", diff --git a/src/app/event.rs b/src/app/event.rs index 3475dea..9811d43 100644 --- a/src/app/event.rs +++ b/src/app/event.rs @@ -50,10 +50,10 @@ pub enum AppEvent { id: i64, result: Result, }, - /// Liked track ids for the ♥ markers. - LikesLoaded(Result, String>), + /// Liked local content ids for the ♥ markers. + LikesLoaded(Result, String>), LikeToggled { - track_id: i64, + content_id: String, liked: bool, }, /// Liked federated item ids and content ids for the ♥ markers. @@ -118,6 +118,7 @@ pub enum AppEvent { /// queue swaps the placeholder for the resolved track. FedTrackResolved { placeholder_id: i64, + resolve_key: String, result: Result, String>, }, /// Rich metadata for a federated track-info preview arrived without diff --git a/src/app/mod.rs b/src/app/mod.rs index c779058..685604e 100644 --- a/src/app/mod.rs +++ b/src/app/mod.rs @@ -44,8 +44,10 @@ pub struct Runtime { /// 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>, + /// Stable playback keys of federated tracks being resolved right now. + pub fed_resolving: std::sync::Mutex>, + /// Stable playback keys that already started through a streaming reader. + pub fed_streaming: Arc>>, /// 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. @@ -58,6 +60,13 @@ pub struct Runtime { pub last_media_push: Option, } +#[derive(Debug, Clone, Copy)] +struct StreamingPlaybackRequest { + volume: u8, + paused: bool, + position_secs: f64, +} + fn now_epoch_seconds() -> i64 { std::time::SystemTime::now() .duration_since(std::time::UNIX_EPOCH) @@ -176,6 +185,7 @@ pub async fn run( 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()), + fed_streaming: Arc::new(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)), @@ -221,6 +231,7 @@ pub async fn run( }, Some(app_event) = event_rx.recv() => handle_app_event(&mut state, &mut runtime, app_event), _ = tick.tick() => { + state.advance_spinner(); expire_quit_confirmation(&mut state); sync_player_shared(&mut state, &runtime); maybe_prefetch_next(&mut state, &runtime); @@ -313,24 +324,7 @@ fn playback_track_to_ui( } fn track_playback_key(track: &crate::library::models::TrackItem) -> String { - if let Some(content_id) = track - .content_id - .as_deref() - .and_then(music_dht::normalize_content_id) - .or_else(|| { - track - .fed - .as_ref() - .and_then(|fed| fed.content_id.as_deref()) - .and_then(music_dht::normalize_content_id) - }) - { - return format!("content:{content_id}"); - } - if let Some(fed) = &track.fed { - return format!("fed:{}:{}", fed.owner, fed.item_id); - } - format!("local:{}", track.id) + state::track_key(track) } fn apply_playback_state_to_ui( @@ -545,6 +539,53 @@ pub(crate) fn transfer_active_to_this_device(state: &mut AppState, runtime: &mut request_urgent_device_sync(runtime); } +pub(crate) fn transfer_active_to_remote_device( + state: &mut AppState, + runtime: &mut Runtime, + target_device_id: String, + target_device_name: String, +) { + if target_device_id.trim().is_empty() + || target_device_id == state.device_playback.self_device_id + { + transfer_active_to_this_device(state, runtime); + return; + } + extrapolate_control_position(state); + let previous_active_id = state.device_playback.active_device_id.clone(); + if state.player.current.is_none() && !state.player.queue.is_empty() { + state.player.current = state.player.queue.get(state.player.queue_pos).cloned(); + } + let wire = playback_state_from_ui(state); + let command = crate::devices::PlaybackCommand::ActiveChanged { + active_device_id: target_device_id.clone(), + active_device_name: target_device_name.clone(), + state: wire.clone(), + }; + record_playback_command_async( + runtime, + target_device_id.clone(), + command.clone(), + "device handoff", + ); + if let Some(previous) = previous_active_id + && previous != target_device_id + && previous != state.device_playback.self_device_id + { + record_playback_command_async(runtime, previous, command.clone(), "device handoff"); + } + let snapshot = crate::devices::PlaybackSnapshot { + device_id: target_device_id.clone(), + device_name: target_device_name.clone(), + active: true, + updated_at_ms: unix_time_ms(), + state: wire, + }; + become_control_device(state, runtime, snapshot); + request_urgent_device_sync(runtime); + state.status_message = Some(format!("active playback moved to {target_device_name}")); +} + fn record_active_handoff( state: &mut AppState, runtime: &Runtime, @@ -733,7 +774,7 @@ fn maintenance(state: &mut AppState, runtime: &mut Runtime) { 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 result = library.liked_content_ids().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)); @@ -952,12 +993,28 @@ fn perform_effect(state: &mut AppState, runtime: &mut Runtime, effect: Effect) { let tx = runtime.event_tx.clone(); tokio::task::spawn_blocking(move || { for track_id in track_ids { - match library.toggle_like(track_id) { + let content_id = match library.track_content_id_by_id(track_id) { + Ok(Some(content_id)) => content_id, + Ok(None) => { + tracing::warn!(track_id, "cannot toggle like without content id"); + let _ = tx.send(AppEvent::StatusMessage( + "like failed: track has no content id".to_string(), + )); + continue; + } + Err(err) => { + tracing::warn!(%err, track_id, "loading track content id failed"); + let _ = + tx.send(AppEvent::StatusMessage(format!("like failed: {err:#}"))); + break; + } + }; + match library.toggle_like_by_content_id(&content_id) { Ok(liked) => { - if let Err(err) = devices.record_track_like(track_id, liked) { + if let Err(err) = devices.record_content_like(&content_id, liked) { tracing::warn!(%err, track_id, "recording synced like failed"); } - let _ = tx.send(AppEvent::LikeToggled { track_id, liked }); + let _ = tx.send(AppEvent::LikeToggled { content_id, liked }); } Err(err) => { tracing::warn!(%err, track_id, "like toggle failed"); @@ -987,6 +1044,12 @@ fn perform_effect(state: &mut AppState, runtime: &mut Runtime, effect: Effect) { } } } + let _ = tx.send(AppEvent::LikesLoaded( + library.liked_content_ids().map_err(err_string), + )); + let _ = tx.send(AppEvent::FedLikesLoaded( + library.fed_like_ids().map_err(err_string), + )); }); } Effect::RemoveFromPlaylist { @@ -1348,13 +1411,13 @@ fn start_current_audio( if let Some(local) = local_track_for_playback(runtime, &track) { state.player.queue[state.player.queue_pos] = local.clone(); track = local; - } else if let Some(fed) = track.fed.clone() { + } else if let Some(fed) = direct_fed_source(&track) { let mut pending = crate::federation::pending_track(&fed); pending.id = track.id; pending.play_count = track.play_count; state.player.queue[state.player.queue_pos] = pending.clone(); track = pending; - } else if track.content_id.is_some() { + } else if track_content_id(&track).is_some() { state.player.current = Some(track.clone()); state.player.playing = true; state.player.paused = paused; @@ -1368,15 +1431,30 @@ fn start_current_audio( "federation: locating \"{}\" for this device…", track.title )); - spawn_content_id_resolve(runtime, &track); + spawn_content_id_resolve( + runtime, + &track, + Some(StreamingPlaybackRequest { + volume: state.player.volume, + paused, + position_secs, + }), + ); return; + } else if let Some(fed) = track.fed.clone() { + let mut pending = crate::federation::pending_track(&fed); + pending.id = track.id; + pending.play_count = track.play_count; + state.player.queue[state.player.queue_pos] = pending.clone(); + track = pending; } } // The track that was playing until now was cut short by this switch. let previous_started_at = state.player.track_started_at; + let next_key = track_playback_key(&track); 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 { + let same_track = track_playback_key(&previous) == next_key; + if state.player.playing && !same_track { report_history( runtime, previous.id, @@ -1404,7 +1482,15 @@ fn start_current_audio( // audio, download it and resume through FedTrackResolved. runtime.player.stop(); state.status_message = Some(format!("federation: fetching \"{}\"…", track.title)); - spawn_fed_resolve(runtime, &track); + spawn_fed_resolve( + runtime, + &track, + Some(StreamingPlaybackRequest { + volume: state.player.volume, + paused, + position_secs, + }), + ); return; } if track_file_missing(&track) { @@ -1454,19 +1540,32 @@ fn local_track_for_playback( runtime: &Runtime, track: &crate::library::models::TrackItem, ) -> Option { - let content_id = track.content_id.as_deref()?; + let content_id = track_content_id(track)?; let local = runtime .library - .track_by_content_id(content_id) + .track_by_content_id(&content_id) .ok() .flatten()?; Path::new(&local.file_path).is_file().then_some(local) } +fn direct_fed_source( + track: &crate::library::models::TrackItem, +) -> Option { + let fed = track.fed.clone()?; + let owner_ok = fed.owner.parse::().is_ok(); + let item_ok = fed.item_id.len() == 64 && fed.item_id.chars().all(|c| c.is_ascii_hexdigit()); + (owner_ok && item_ok).then_some(fed) +} + +fn track_content_id(track: &crate::library::models::TrackItem) -> Option { + state::track_content_id(track) +} + 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)) + Ok((Box::new(std::io::BufReader::new(file)), byte_len)) } /// Open the next queue item ~30s before the current track ends and append @@ -1488,14 +1587,28 @@ fn maybe_prefetch_next(state: &mut AppState, runtime: &Runtime) { let Some(next_pos) = update::peek_next_pos(player) else { return; }; - let Some(next) = player.queue.get(next_pos).cloned() else { + let Some(mut 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; + if track_file_missing(&next) { + if let Some(local) = local_track_for_playback(runtime, &next) { + state.player.queue[next_pos] = local.clone(); + next = local; + } else if let Some(fed) = direct_fed_source(&next) { + next = crate::federation::pending_track(&fed); + spawn_fed_resolve(runtime, &next, None); + return; + } else if track_content_id(&next).is_some() { + // Resolve by content id first: the actual owner/item id may be a + // synced placeholder from another client. + spawn_content_id_resolve(runtime, &next, None); + return; + } else if next.is_fed_pending() { + spawn_fed_resolve(runtime, &next, None); + return; + } else { + return; + } } state.player.prefetched_pos = Some(next_pos); tracing::debug!(title = %next.title, "prefetching next track"); @@ -1577,64 +1690,235 @@ pub(crate) fn device_connect(runtime: &Runtime, invite: String) { }); } +fn download_progress_sender( + tx: mpsc::UnboundedSender, + title: String, +) -> impl FnMut(crate::federation::DownloadProgress) + Send + 'static { + let mut last_sent: Option = None; + let started = std::time::Instant::now(); + move |progress| { + let now = std::time::Instant::now(); + let complete = progress.total > 0 && progress.received >= progress.total; + let first = progress.received == 0; + let due = last_sent.is_none_or(|last| now.duration_since(last) >= TICK_INTERVAL); + if first || complete || due { + last_sent = Some(now); + let elapsed_secs = now.duration_since(started).as_secs_f64(); + let bytes_per_sec = if elapsed_secs >= 0.25 && progress.received > 0 { + Some(progress.received as f64 / elapsed_secs) + } else { + None + }; + let _ = tx.send(AppEvent::StatusMessage(format_download_progress( + &title, + progress, + bytes_per_sec, + ))); + } + } +} + +fn format_download_progress( + title: &str, + progress: crate::federation::DownloadProgress, + bytes_per_sec: Option, +) -> String { + let speed = bytes_per_sec + .map(|bytes| format!(" · {}", format_transfer_rate(bytes))) + .unwrap_or_default(); + if progress.total > 0 { + let percent = (progress.received as f64 / progress.total as f64 * 100.0).clamp(0.0, 100.0); + format!( + "federation: downloading \"{title}\" {:.0}% · {}/{}{speed}", + percent, + format_bytes(progress.received), + format_bytes(progress.total) + ) + } else { + format!( + "federation: downloading \"{title}\" · {}{speed}", + format_bytes(progress.received), + ) + } +} + +fn format_bytes(bytes: u64) -> String { + const KIB: f64 = 1024.0; + const MIB: f64 = KIB * 1024.0; + if bytes >= 1024 * 1024 { + format!("{:.1} MB", bytes as f64 / MIB) + } else if bytes >= 1024 { + format!("{:.0} KB", bytes as f64 / KIB) + } else { + format!("{bytes} B") + } +} + +fn format_transfer_rate(bytes_per_sec: f64) -> String { + const KIB: f64 = 1024.0; + const MIB: f64 = KIB * 1024.0; + if bytes_per_sec >= MIB { + format!("{:.1} MB/s", bytes_per_sec / MIB) + } else if bytes_per_sec >= KIB { + format!("{:.0} KB/s", bytes_per_sec / KIB) + } else { + format!("{:.0} B/s", bytes_per_sec) + } +} + /// 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) { +fn spawn_fed_resolve( + runtime: &Runtime, + track: &crate::library::models::TrackItem, + stream_playback: Option, +) { let Some(fed_track) = track.fed.clone() else { return; }; + let resolve_key = track_playback_key(track); { let mut resolving = runtime .fed_resolving .lock() .unwrap_or_else(std::sync::PoisonError::into_inner); - if !resolving.insert(track.id) { + if !resolving.insert(resolve_key.clone()) { return; } } let placeholder_id = track.id; + let title = track.title.clone(); let fed = Arc::clone(&runtime.federation); let tx = runtime.event_tx.clone(); + let controller = runtime.player.clone(); + let streaming = Arc::clone(&runtime.fed_streaming); tokio::spawn(async move { - let result = fed - .prepare_playback(&fed_track) + let progress = download_progress_sender(tx.clone(), title.clone()); + let result = if let Some(playback) = + stream_playback.filter(|playback| playback.position_secs <= 0.5) + { + let stream_tx = tx.clone(); + let stream_title = title.clone(); + let stream_controller = controller.clone(); + let stream_resolve_key = resolve_key.clone(); + let stream_markers = Arc::clone(&streaming); + let mut started = false; + fed.prepare_playback_streaming_with_progress(&fed_track, progress, move |stream| { + if started { + return; + } + started = true; + stream_markers + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .insert(stream_resolve_key.clone()); + stream_controller.play_stream( + Box::new(stream.reader), + Some(stream.mime_type), + player::amplitude(playback.volume), + ); + if playback.paused { + stream_controller.pause(); + } + let _ = stream_tx.send(AppEvent::StatusMessage(format!( + "federation: streaming \"{stream_title}\" while downloading…" + ))); + }) .await - .map(Box::new) - .map_err(|err| format!("{err:#}")); + } else { + fed.prepare_playback_with_progress(&fed_track, progress) + .await + } + .map(Box::new) + .map_err(|err| format!("{err:#}")); let _ = tx.send(AppEvent::FedTrackResolved { placeholder_id, + resolve_key, result, }); }); } -fn spawn_content_id_resolve(runtime: &Runtime, track: &crate::library::models::TrackItem) { - let Some(content_id) = track.content_id.clone() else { +fn spawn_content_id_resolve( + runtime: &Runtime, + track: &crate::library::models::TrackItem, + stream_playback: Option, +) { + let Some(content_id) = track_content_id(track) else { return; }; + let resolve_key = track_playback_key(track); { let mut resolving = runtime .fed_resolving .lock() .unwrap_or_else(std::sync::PoisonError::into_inner); - if !resolving.insert(track.id) { + if !resolving.insert(resolve_key.clone()) { return; } } let placeholder_id = track.id; let label = format!("{} {}", track.artist_line(), track.title); + let title = track.title.clone(); let fed = Arc::clone(&runtime.federation); let tx = runtime.event_tx.clone(); + let controller = runtime.player.clone(); + let streaming = Arc::clone(&runtime.fed_streaming); tokio::spawn(async move { let result = match fed.track_by_content_id(&content_id, Some(&label)).await { - Ok(fed_track) => fed.prepare_playback(&fed_track).await, + Ok(fed_track) => { + let _ = tx.send(AppEvent::StatusMessage(format!( + "federation: found source for \"{title}\", downloading…" + ))); + let progress = download_progress_sender(tx.clone(), title.clone()); + if let Some(playback) = + stream_playback.filter(|playback| playback.position_secs <= 0.5) + { + let stream_tx = tx.clone(); + let stream_title = title.clone(); + let stream_controller = controller.clone(); + let stream_resolve_key = resolve_key.clone(); + let stream_markers = Arc::clone(&streaming); + let mut started = false; + fed.prepare_playback_streaming_with_progress( + &fed_track, + progress, + move |stream| { + if started { + return; + } + started = true; + stream_markers + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .insert(stream_resolve_key.clone()); + stream_controller.play_stream( + Box::new(stream.reader), + Some(stream.mime_type), + player::amplitude(playback.volume), + ); + if playback.paused { + stream_controller.pause(); + } + let _ = stream_tx.send(AppEvent::StatusMessage(format!( + "federation: streaming \"{stream_title}\" while downloading…" + ))); + }, + ) + .await + } else { + fed.prepare_playback_with_progress(&fed_track, progress) + .await + } + } Err(err) => Err(err), } .map(Box::new) .map_err(|err| format!("{err:#}")); let _ = tx.send(AppEvent::FedTrackResolved { placeholder_id, + resolve_key, result, }); }); @@ -1660,12 +1944,14 @@ pub(crate) fn fed_download_spawn( let mut imported_fed_tracks = Vec::new(); let mut failed = 0usize; for (index, track) in tracks.iter().enumerate() { + let title = track.title.clone(); let _ = tx.send(AppEvent::StatusMessage(format!( "federation: downloading {}/{total}: {}", index + 1, track.title ))); - match fed.download_to_library(track).await { + let progress = download_progress_sender(tx.clone(), title); + match fed.download_to_library_with_progress(track, progress).await { Ok(imported) => { imported_ids.push(imported.id); imported_fed_tracks.push(track.clone()); @@ -1967,10 +2253,18 @@ fn apply_queue_refresh( ) { let by_id: std::collections::HashMap = tracks.into_iter().map(|track| (track.id, track)).collect(); + let by_key: std::collections::HashMap = by_id + .values() + .cloned() + .map(|track| (track_playback_key(&track), track)) + .collect(); let current_id = state.player.current.as_ref().map(|track| track.id); + let current_key = state.player.current.as_ref().map(track_playback_key); for track in &mut state.player.queue { if let Some(fresh) = by_id.get(&track.id) { *track = fresh.clone(); + } else if let Some(fresh) = by_key.get(&track_playback_key(track)) { + *track = fresh.clone(); } } let had_missing = state @@ -1995,15 +2289,29 @@ fn apply_queue_refresh( 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) { + match (current_id, current_key) { + (Some(id), _) if id < 0 => { + // The current source can be an ephemeral federated stream while + // the queue is refreshed after save-on-listen. It is not a + // deleted library row, so keep playback untouched. + state.player.queue_pos = state.player.queue_pos.min(state.player.queue.len() - 1); + } + (Some(id), Some(key)) if by_id.contains_key(&id) || by_key.contains_key(&key) => { + if let Some(position) = state + .player + .queue + .iter() + .position(|track| track_playback_key(track) == key) + { state.player.queue_pos = position; } - state.player.current = by_id.get(&id).cloned(); + state.player.current = by_id + .get(&id) + .cloned() + .or_else(|| by_key.get(&key).cloned()); push_media_metadata(state, runtime); } - Some(_) => { + (Some(_), _) => { // The playing track was deleted from the library. runtime.player_start_pending = false; runtime.player.stop(); @@ -2013,7 +2321,7 @@ fn apply_queue_refresh( state.player.paused = false; push_media_update(state, runtime, true); } - None => { + (None, _) => { state.player.queue_pos = state.player.queue_pos.min(state.player.queue.len() - 1); } } @@ -2251,18 +2559,29 @@ fn handle_app_event(state: &mut AppState, runtime: &mut Runtime, event: AppEvent } AppEvent::FedTrackResolved { placeholder_id, + resolve_key, result, } => { runtime .fed_resolving .lock() .unwrap_or_else(std::sync::PoisonError::into_inner) - .remove(&placeholder_id); + .remove(&resolve_key); + let stream_started = runtime + .fed_streaming + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .remove(&resolve_key); if state.device_playback.is_control() { + runtime.player_start_pending = false; tracing::debug!( placeholder_id, + resolve_key, "ignored local federated track resolution while controlling remote playback" ); + if let Err(message) = result { + state.status_message = Some(format!("federation: {message}")); + } return; } let placeholder_key = state @@ -2277,7 +2596,8 @@ fn handle_app_event(state: &mut AppState, runtime: &mut Runtime, event: AppEvent .as_ref() .filter(|track| track.id == placeholder_id) }) - .map(track_playback_key); + .map(track_playback_key) + .or(Some(resolve_key.clone())); let queue_pos_waiting = state .player @@ -2301,6 +2621,7 @@ fn handle_app_event(state: &mut AppState, runtime: &mut Runtime, event: AppEvent } else { 0.0 }; + let already_streaming = waiting && stream_started && state.player.playing; match result { Ok(playable) => { if playable.imported { @@ -2327,7 +2648,8 @@ fn handle_app_event(state: &mut AppState, runtime: &mut Runtime, event: AppEvent } } if waiting { - // Playback was parked on this track; start it now. + // Playback was either parked on this track or already + // running from the streaming reader. let paused = state.player.paused; if state .player @@ -2342,11 +2664,19 @@ fn handle_app_event(state: &mut AppState, runtime: &mut Runtime, event: AppEvent { state.player.queue_pos = pos; } - state.player.current = - state.player.queue.get(state.player.queue_pos).cloned(); - start_current_audio(state, runtime, resume_position_secs, paused); - push_media_metadata(state, runtime); - push_media_update(state, runtime, true); + if already_streaming { + // Do not swap the currently playing item under + // the audio pipeline. The streaming reader will + // play through the completed cache file naturally; + // the resolved queue entry is for future starts. + push_media_update(state, runtime, true); + } else { + state.player.current = + state.player.queue.get(state.player.queue_pos).cloned(); + start_current_audio(state, runtime, resume_position_secs, paused); + push_media_metadata(state, runtime); + push_media_update(state, runtime, true); + } } } Err(message) => { @@ -2624,6 +2954,7 @@ fn handle_app_event(state: &mut AppState, runtime: &mut Runtime, event: AppEvent content_id.and_then(|id| music_dht::normalize_content_id(&id)) { state.fed_likes.remove(&content_id); + state.likes.remove(&content_id); } } // The virtual Likes playlist is stale now; refetch on next open. @@ -2635,11 +2966,13 @@ fn handle_app_event(state: &mut AppState, runtime: &mut Runtime, event: AppEvent "like removed".to_string() }); } - AppEvent::LikeToggled { track_id, liked } => { + AppEvent::LikeToggled { content_id, liked } => { if liked { - state.likes.insert(track_id); + state.fed_likes.remove(&content_id); + state.likes.insert(content_id); } else { - state.likes.remove(&track_id); + state.likes.remove(&content_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); diff --git a/src/app/popup.rs b/src/app/popup.rs index a034bfa..12320e9 100644 --- a/src/app/popup.rs +++ b/src/app/popup.rs @@ -190,10 +190,18 @@ fn handle_connected_devices( super::transfer_active_to_this_device(state, runtime); state.status_message = Some("active playback moved to this device".into()); } else if let Some(row) = other_rows.get(cursor.saturating_sub(1).min(last)) { - if let Some(snapshot) = state.device_playback.remote.get(&row.device_id).cloned() { + if row.active + && let Some(snapshot) = + state.device_playback.remote.get(&row.device_id).cloned() + { super::become_control_device(state, runtime, snapshot); } else { - state.status_message = Some("device has no playback snapshot yet".into()); + super::transfer_active_to_remote_device( + state, + runtime, + row.device_id.clone(), + row.name.clone(), + ); } } } diff --git a/src/app/state.rs b/src/app/state.rs index f93d26a..d9a44a4 100644 --- a/src/app/state.rs +++ b/src/app/state.rs @@ -900,9 +900,9 @@ pub struct PlayerBar { pub prefetched_pos: Option, pub volume: u8, pub shuffle: bool, - /// Track ids in pre-shuffle order; restores the queue when shuffle is - /// turned off. - pub original_order: Option>, + /// Stable track keys in pre-shuffle order; restores the queue when + /// shuffle is turned off. + pub original_order: Option>, pub repeat: RepeatMode, } @@ -980,6 +980,7 @@ pub struct AppState { pub help_visible: bool, pub pending_keys: Option, pub status_message: Option, + pub spinner_frame: usize, pub settings_cursor: usize, pub player: PlayerBar, pub device_playback: DevicePlaybackState, @@ -989,8 +990,8 @@ pub struct AppState { pub release_views: HashMap>, pub playlists: PlaylistsTab, pub playlist_views: HashMap>, - /// Liked track ids, for the ♥ markers everywhere tracks are shown. - pub likes: std::collections::HashSet, + /// Liked local-library content ids, for the ♥ markers everywhere tracks are shown. + pub likes: std::collections::HashSet, /// Liked federated tracks (DHT item ids and content ids) — likes that /// reference peers' tracks without downloading them. pub fed_likes: std::collections::HashSet, @@ -1022,6 +1023,15 @@ pub struct AppState { } impl AppState { + pub fn advance_spinner(&mut self) { + self.spinner_frame = self.spinner_frame.wrapping_add(1); + } + + pub fn spinner(&self) -> &'static str { + const FRAMES: &[&str] = &["⠋", "⠙", "⠹", "⠸", "⠼", "⠴", "⠦", "⠧", "⠇", "⠏"]; + FRAMES[self.spinner_frame % FRAMES.len()] + } + pub fn connected_devices_enabled(&self) -> bool { self.federation.settings.enabled && !self.federation.settings.network_id.trim().is_empty() } @@ -1032,7 +1042,9 @@ impl AppState { .content_id .as_deref() .and_then(music_dht::normalize_content_id) - .is_some_and(|content_id| self.fed_likes.contains(&content_id)) + .is_some_and(|content_id| { + self.likes.contains(&content_id) || self.fed_likes.contains(&content_id) + }) } pub fn fed_card_track_liked(&self, track: &crate::federation::FedCardTrack) -> bool { @@ -1040,10 +1052,46 @@ impl AppState { .content_id .as_deref() .and_then(music_dht::normalize_content_id) - .is_some_and(|content_id| self.fed_likes.contains(&content_id)) + .is_some_and(|content_id| { + self.likes.contains(&content_id) || self.fed_likes.contains(&content_id) + }) || track .sources .iter() .any(|(_, item_id)| self.fed_likes.contains(item_id)) } + + pub fn track_liked(&self, track: &TrackItem) -> bool { + if let Some(content_id) = track_content_id(track) { + return self.likes.contains(&content_id) || self.fed_likes.contains(&content_id); + } + track + .fed + .as_ref() + .is_some_and(|fed| self.fed_track_liked(fed)) + } +} + +pub fn track_content_id(track: &TrackItem) -> Option { + track + .content_id + .as_deref() + .and_then(music_dht::normalize_content_id) + .or_else(|| { + track + .fed + .as_ref() + .and_then(|fed| fed.content_id.as_deref()) + .and_then(music_dht::normalize_content_id) + }) +} + +pub fn track_key(track: &TrackItem) -> String { + if let Some(content_id) = track_content_id(track) { + return format!("content:{content_id}"); + } + if let Some(fed) = &track.fed { + return format!("fed:{}:{}", fed.owner, fed.item_id); + } + format!("local:{}", track.id) } diff --git a/src/app/update.rs b/src/app/update.rs index a031d21..cffb9f2 100644 --- a/src/app/update.rs +++ b/src/app/update.rs @@ -6,7 +6,7 @@ use crate::library::models::TrackItem; use super::state::{ AppState, GlobalView, Loadable, OpenedPlaylist, SearchState, TILE_HEIGHT, TILE_WIDTH, Tab, TrackSelectionScope, ViewMode, fed_release_display_order, fed_release_rows, - release_display_order, release_rows, settings_rows, + release_display_order, release_rows, settings_rows, track_content_id, track_key, }; pub const QUIT_CONFIRM_WINDOW: Duration = Duration::from_millis(1500); @@ -283,26 +283,31 @@ pub fn update(state: &mut AppState, action: Action) -> Option { .map(|track| vec![track.clone()]) .unwrap_or_default(); } - // Local library rows are liked by id; federated tracks (pending - // or cached) are liked as references into the federation. - let mut track_ids: Vec = Vec::new(); + let mut local_tracks: Vec = Vec::new(); let mut fed_tracks: Vec = Vec::new(); + let mut seen = std::collections::HashSet::new(); for track in tracks { + if !seen.insert(track_key(&track)) { + continue; + } match &track.fed { Some(fed) => fed_tracks.push(fed.clone()), - None if track.id >= 0 => track_ids.push(track.id), + None if track.id >= 0 && track_content_id(&track).is_some() => { + local_tracks.push(track) + } None => {} } } - if track_ids.is_empty() && fed_tracks.is_empty() { + if local_tracks.is_empty() && fed_tracks.is_empty() { state.status_message = Some("no track selected".into()); None } else { - let should_like = track_ids.iter().any(|id| !state.likes.contains(id)) + let should_like = local_tracks.iter().any(|track| !state.track_liked(track)) || fed_tracks.iter().any(|fed| !state.fed_track_liked(fed)); - let toggles: Vec = track_ids + let toggles: Vec = local_tracks .into_iter() - .filter(|id| state.likes.contains(id) != should_like) + .filter(|track| state.track_liked(track) != should_like) + .map(|track| track.id) .collect(); let fed_toggles: Vec = fed_tracks .into_iter() @@ -355,10 +360,9 @@ pub fn update(state: &mut AppState, action: Action) -> Option { None } Action::DownloadSelected => { - let tracks = selected_fed_tracks(state); + let tracks = selected_downloadable_fed_tracks(state); if tracks.is_empty() { - state.status_message = - Some("select federated tracks first (works in the federation results)".into()); + state.status_message = Some("select federated tracks first".into()); None } else { state.track_selection.clear(); @@ -614,30 +618,38 @@ fn delete_selected(state: &mut AppState) -> Option { return None; } let track_ids: Vec = tracks.iter().map(|track| track.id).collect(); - let content_ids: Vec = tracks - .iter() - .filter_map(|track| { - track - .content_id - .as_deref() - .and_then(music_dht::normalize_content_id) - }) - .collect(); + let content_ids: Vec = tracks.iter().filter_map(track_content_id).collect(); state.track_selection.clear(); if opened.id == super::state::LIKES_PLAYLIST_ID { - let liked: Vec = track_ids - .into_iter() - .filter(|id| state.likes.contains(id)) - .collect(); - state.status_message = Some(format!("removing {} like(s)", liked.len())); + let mut seen = std::collections::HashSet::new(); + let mut liked = Vec::new(); + let mut fed_tracks = Vec::new(); + for track in tracks { + if !state.track_liked(&track) || !seen.insert(track_key(&track)) { + continue; + } + if let Some(fed) = track.fed { + fed_tracks.push(fed); + } else if track.id >= 0 && track_content_id(&track).is_some() { + liked.push(track.id); + } + } + let total = liked.len() + fed_tracks.len(); + state.status_message = Some(format!("removing {} like(s)", total)); return Some(Effect::ToggleLikes { track_ids: liked, - fed_tracks: vec![], + fed_tracks, }); } + let content_ids: Vec = content_ids + .into_iter() + .collect::>() + .into_iter() + .collect(); + let track_ids: Vec = track_ids.into_iter().filter(|id| *id >= 0).collect(); state.status_message = Some(format!( "removing {} track(s) from playlist", - track_ids.len() + content_ids.len().max(track_ids.len()) )); return Some(Effect::RemoveFromPlaylist { playlist_id: opened.id, @@ -1015,14 +1027,14 @@ fn remove_queue_indices(state: &mut AppState, indices: &[usize]) -> QueueRemoval } let old_queue_pos = state.player.queue_pos; - let current_id = state.player.current.as_ref().map(|track| track.id); - let current_removed = current_id.is_some_and(|id| { + let current_key = state.player.current.as_ref().map(track_key); + let current_removed = current_key.as_ref().is_some_and(|key| { unique.iter().any(|index| { state .player .queue .get(*index) - .is_some_and(|track| track.id == id) + .is_some_and(|track| track_key(track) == *key) }) }); let removed_before_current = unique @@ -1060,19 +1072,24 @@ fn remove_queue_indices(state: &mut AppState, indices: &[usize]) -> QueueRemoval }; } - if let Some(id) = current_id { - if let Some(position) = state.player.queue.iter().position(|track| track.id == id) { + if let Some(key) = current_key.as_ref() { + if let Some(position) = state + .player + .queue + .iter() + .position(|track| track_key(track) == *key) + { state.player.queue_pos = position; } } else { state.player.queue_pos = state.player.queue_pos.min(state.player.queue.len() - 1); } - state.player.current = current_id.and_then(|id| { + state.player.current = current_key.and_then(|key| { state .player .queue .iter() - .find(|track| track.id == id) + .find(|track| track_key(track) == key) .cloned() }); state.queue_tab.cursor = state.queue_tab.cursor.min(state.player.queue.len() - 1); @@ -1416,7 +1433,7 @@ pub fn shuffle_upcoming(player: &mut super::state::PlayerBar) { return; } if player.original_order.is_none() { - player.original_order = Some(player.queue.iter().map(|t| t.id).collect()); + player.original_order = Some(player.queue.iter().map(track_key).collect()); } shuffle_range(player, upcoming_start(player)); } @@ -1440,7 +1457,7 @@ pub fn restore_queue_order(player: &mut super::state::PlayerBar) { let key = order .iter() .enumerate() - .position(|(slot, id)| !used[slot] && *id == track.id) + .position(|(slot, key)| !used[slot] && *key == track_key(&track)) .inspect(|&slot| used[slot] = true) .unwrap_or(usize::MAX); (key, position, track) @@ -2212,6 +2229,83 @@ pub(crate) fn selected_fed_tracks(state: &AppState) -> Vec Vec { + let mut tracks = Vec::new(); + for track in selected_tracks(state) { + if let Some(fed) = fed_track_for_download(&track) + && !tracks + .iter() + .any(|existing| same_fed_download(existing, &fed)) + { + tracks.push(fed); + } + } + for fed in selected_fed_tracks(state) { + if !tracks + .iter() + .any(|existing| same_fed_download(existing, &fed)) + { + tracks.push(fed); + } + } + tracks +} + +fn fed_track_for_download(track: &TrackItem) -> Option { + if let Some(fed) = &track.fed { + return Some(fed.clone()); + } + let content_id = track + .content_id + .as_deref() + .and_then(music_dht::normalize_content_id)?; + Some(crate::federation::FedTrack { + item_id: String::new(), + owner: String::new(), + own: false, + title: track.title.clone(), + artist_names: track + .artists + .iter() + .map(|artist| artist.name.clone()) + .collect(), + featured_artist_names: track + .featured_artists + .iter() + .map(|artist| artist.name.clone()) + .collect(), + year: track.release_year, + duration_seconds: (track.duration_seconds > 0.0) + .then(|| track.duration_seconds.round().max(0.0) as i64), + content_id: Some(content_id), + release_title: (!track.release_title.trim().is_empty()) + .then(|| track.release_title.clone()), + track_number: track.track_number, + disc_number: track.disc_number, + }) +} + +fn same_fed_download( + left: &crate::federation::FedTrack, + right: &crate::federation::FedTrack, +) -> bool { + let left_content = left + .content_id + .as_deref() + .and_then(music_dht::normalize_content_id); + let right_content = right + .content_id + .as_deref() + .and_then(music_dht::normalize_content_id); + if left_content.is_some() && left_content == right_content { + return true; + } + !left.owner.is_empty() + && !left.item_id.is_empty() + && left.owner == right.owner + && left.item_id == right.item_id +} + /// What Enter resolved to in the current view. enum Outcome { Push(GlobalView), @@ -2482,7 +2576,7 @@ pub(super) fn on_new_queue(state: &mut AppState) { let player = &mut state.player; player.original_order = None; if player.shuffle && !player.queue.is_empty() { - player.original_order = Some(player.queue.iter().map(|t| t.id).collect()); + player.original_order = Some(player.queue.iter().map(track_key).collect()); shuffle_range(player, (player.queue_pos + 1).min(player.queue.len())); } } @@ -2610,7 +2704,7 @@ mod tests { release_year: None, cover_path: None, file_path: format!("/s/{id}"), - content_id: None, + content_id: Some(format!("b3:{id:064x}")), audio_format: None, audio_bitrate: None, audio_sample_rate: None, @@ -3162,7 +3256,7 @@ mod tests { ..AppState::default() }; state.player.queue = (1..=3).map(test_track).collect(); - state.likes.insert(1); + state.likes.insert(format!("b3:{:064x}", 1)); update(&mut state, Action::ToggleTrackSelection); update(&mut state, Action::SelectLast); @@ -3174,7 +3268,10 @@ mod tests { }) ); - state.likes = [1, 2, 3].into_iter().collect(); + state.likes = [1, 2, 3] + .into_iter() + .map(|id| format!("b3:{id:064x}")) + .collect(); assert_eq!( update(&mut state, Action::ToggleLike), Some(Effect::ToggleLikes { diff --git a/src/devices.rs b/src/devices.rs index 9ec9222..29967a5 100644 --- a/src/devices.rs +++ b/src/devices.rs @@ -728,6 +728,7 @@ impl DeviceSync { &self, service: Arc, invite_link: &str, + transport_stats: Arc, ) -> Result { let invite = parse_invite(invite_link)?; anyhow::ensure!(invite.expires_at_ms >= now_ms(), "invite expired"); @@ -739,7 +740,12 @@ impl DeviceSync { let mut last_error: Option; loop { match self - .try_connect_invite(Arc::clone(&service), &invite, ticket.clone()) + .try_connect_invite( + Arc::clone(&service), + &invite, + ticket.clone(), + Arc::clone(&transport_stats), + ) .await { Ok(PairAttempt::Accepted(message)) => return Ok(message), @@ -765,6 +771,7 @@ impl DeviceSync { service: Arc, invite: &InviteWire, ticket: PeerTicket, + transport_stats: Arc, ) -> Result { let peer = service.connect(ticket).await?; let own_ticket = service.ticket().await?.to_string(); @@ -777,6 +784,13 @@ impl DeviceSync { let snapshot = self.snapshot()?; let playback = self.local_playback_snapshot(); let mut stream = service.open_stream(peer, SYNC_ALPN).await?; + crate::federation::record_stream_transport( + &transport_stats, + "device-sync", + "outbound", + "pair-open", + &stream, + ); write_msg( &mut stream, &WireMessage::PairRequest { @@ -794,10 +808,17 @@ impl DeviceSync { ) .await?; finish_send(&mut stream).await?; - match read_msg(&mut stream) + let response = read_msg(&mut stream) .await - .context("pairing response was not received")? - { + .context("pairing response was not received")?; + crate::federation::record_stream_transport( + &transport_stats, + "device-sync", + "outbound", + "pair-done", + &stream, + ); + match response { WireMessage::PairResponse { accepted: true, group_id: Some(group_id), @@ -936,14 +957,15 @@ impl DeviceSync { Ok(()) } - pub fn record_track_like(&self, track_id: i64, liked: bool) -> Result<()> { - if let Some(content_id) = self.library.track_content_id_by_id(track_id)? { - self.record_local_op(SyncOpPayload::TrackLikeSet { - content_id, - liked, - fed: None, - })?; - } + pub fn record_content_like(&self, content_id: &str, liked: bool) -> Result<()> { + let Some(content_id) = music_dht::normalize_content_id(content_id) else { + return Ok(()); + }; + self.record_local_op(SyncOpPayload::TrackLikeSet { + content_id, + liked, + fed: None, + })?; Ok(()) } @@ -1057,13 +1079,20 @@ impl DeviceSync { Ok(()) } - pub async fn sync_once(&self, service: Arc) -> Result<()> { + pub async fn sync_once( + &self, + service: Arc, + transport_stats: Arc, + ) -> Result<()> { let devices = self.active_remote_devices()?; for device in devices { if device.endpoint_ticket.trim().is_empty() { continue; } - if let Err(err) = self.sync_device(Arc::clone(&service), &device).await { + if let Err(err) = self + .sync_device(Arc::clone(&service), &device, Arc::clone(&transport_stats)) + .await + { tracing::debug!(device = %device.device_id, "device sync failed: {err:#}"); self.set_last_error(Some(format!("{}: {err:#}", short_id(&device.device_id))))?; } @@ -1076,6 +1105,7 @@ impl DeviceSync { &self, service: Arc, device: &StoredDevice, + transport_stats: Arc, ) -> Result<()> { let ticket: PeerTicket = device.endpoint_ticket.parse()?; let peer = service.connect(ticket).await?; @@ -1088,6 +1118,13 @@ impl DeviceSync { let snapshot = self.snapshot()?; let playback = self.local_playback_snapshot(); let mut stream = service.open_stream(peer, SYNC_ALPN).await?; + crate::federation::record_stream_transport( + &transport_stats, + "device-sync", + "outbound", + "sync-open", + &stream, + ); write_msg( &mut stream, &WireMessage::Hello { @@ -1102,10 +1139,17 @@ impl DeviceSync { ) .await?; finish_send(&mut stream).await?; - match read_msg(&mut stream) + let response = read_msg(&mut stream) .await - .context("device sync response was not received")? - { + .context("device sync response was not received")?; + crate::federation::record_stream_transport( + &transport_stats, + "device-sync", + "outbound", + "sync-done", + &stream, + ); + match response { WireMessage::SyncResponse { accepted: true, devices, @@ -2432,24 +2476,33 @@ pub async fn serve_peers( mut acceptor: StreamAcceptor, sync: Arc, service: Arc, + transport_stats: Arc, ) { while let Some(stream) = acceptor.accept().await { let sync = Arc::clone(&sync); let service = Arc::clone(&service); + let transport_stats = Arc::clone(&transport_stats); tokio::spawn(async move { let peer = stream.peer_id; - if let Err(err) = serve_one(stream, sync, service).await { + if let Err(err) = serve_one(stream, sync, service, transport_stats).await { tracing::warn!(peer = %peer, "personal sync stream failed: {err:#}"); } }); } } -pub async fn sync_loop(sync: Arc, service: Arc) { +pub async fn sync_loop( + sync: Arc, + service: Arc, + transport_stats: Arc, +) { let mut interval = tokio::time::interval(SYNC_INTERVAL); loop { interval.tick().await; - if let Err(err) = sync.sync_once(Arc::clone(&service)).await { + if let Err(err) = sync + .sync_once(Arc::clone(&service), Arc::clone(&transport_stats)) + .await + { tracing::debug!("personal sync tick failed: {err:#}"); } if let Some(tx) = lock(&sync.event_tx).as_ref() { @@ -2462,7 +2515,15 @@ async fn serve_one( mut stream: ByteStream, sync: Arc, service: Arc, + transport_stats: Arc, ) -> Result<()> { + crate::federation::record_stream_transport( + &transport_stats, + "device-sync", + "inbound", + "open", + &stream, + ); match read_msg(&mut stream).await? { WireMessage::PairRequest { invite_id, diff --git a/src/federation/audio.rs b/src/federation/audio.rs index b8c81db..c53cab9 100644 --- a/src/federation/audio.rs +++ b/src/federation/audio.rs @@ -7,6 +7,7 @@ use std::path::{Path, PathBuf}; use std::sync::Arc; +use std::time::Instant; use anyhow::{Context, Result}; use music_dht::{ByteStream, EndpointId, ItemId, ItemKind, MusicDhtService, StreamAcceptor}; @@ -20,6 +21,22 @@ pub const AUDIO_ALPN: &[u8] = b"furumi-fd/audio/1"; /// Maximum size of a JSON protocol line (request or response header). const MAX_PROTOCOL_LINE: usize = 4096; +const STREAM_START_BUFFER_BYTES: u64 = 2 * 1024 * 1024; + +#[derive(Debug, Clone, Copy)] +pub struct DownloadProgress { + pub stream_key: u64, + pub transport_phase: &'static str, + pub received: u64, + pub total: u64, + pub transport: Option, +} + +#[derive(Debug)] +pub struct StreamingStart { + pub reader: crate::streaming::GrowingFileReader, + pub mime_type: String, +} #[derive(Debug, Serialize, Deserialize)] struct AudioRequest { @@ -278,26 +295,31 @@ pub async fn fetch_metadata( }) } -/// Downloads a whole track (with metadata and cover art) from `owner` into -/// `dir/.`. An already complete cached audio file is reused; -/// the metadata and cover still come fresh from the header. -pub async fn download_track( +pub async fn download_track_with_streaming( service: &MusicDhtService, owner: EndpointId, item_id_hex: &str, dir: &Path, stem: &str, -) -> Result { + want_images: bool, + mut progress: F, + mut stream_start: Option<&mut (dyn FnMut(StreamingStart) + Send)>, +) -> Result +where + F: FnMut(DownloadProgress) + Send, +{ + let started = Instant::now(); let mut stream = service .open_stream(owner, AUDIO_ALPN) .await .map_err(|err| anyhow::anyhow!("cannot reach the owner peer: {err}"))?; + let stream_key = super::stream_transport_key(&stream); write_line( &mut stream.send, &AudioRequest { item_id: item_id_hex.to_string(), offset: 0, - want_cover: true, + want_cover: want_images, metadata_only: false, }, ) @@ -311,6 +333,22 @@ pub async fn download_track( header.error.unwrap_or_else(|| "unknown error".to_string()) ); } + progress(DownloadProgress { + stream_key, + transport_phase: "download-open", + received: 0, + total: header.total_size, + transport: Some(stream.connection_stats()), + }); + tracing::info!( + owner = %owner, + item_id = %item_id_hex, + audio_bytes = header.total_size, + cover_bytes = header.cover_size, + artist_image_bytes = header.artist_image_size, + want_images, + "federated audio stream opened" + ); // The image segments precede the audio bytes and are read regardless of // the cache state — they sit first in the stream. @@ -346,6 +384,20 @@ pub async fn download_track( && header.total_size > 0 { // Audio already fully downloaded earlier; no need to fetch again. + progress(DownloadProgress { + stream_key, + transport_phase: "download-done", + received: header.total_size, + total: header.total_size, + transport: Some(stream.connection_stats()), + }); + tracing::info!( + owner = %owner, + item_id = %item_id_hex, + bytes = header.total_size, + path = %path.display(), + "federated audio reused from cache" + ); return Ok(Downloaded { path, mime_type: header.mime_type, @@ -357,12 +409,56 @@ pub async fn download_track( let temp_path = dir.join(format!(".{stem}.{extension}.part")); let mut file = tokio::fs::File::create(&temp_path).await?; + let mut streaming = match stream_start.as_ref() { + Some(_) => match crate::streaming::growing_file(&temp_path) { + Ok((reader, writer)) => Some((Some(reader), writer, false)), + Err(err) => { + tracing::debug!(%err, path = %temp_path.display(), "streaming reader disabled"); + None + } + }, + None => None, + }; + let stream_start_at = stream_start_buffer(header.total_size); let mut received: u64 = 0; let mut chunk = vec![0u8; 64 * 1024]; // quinn's inherent read returns None when the peer finished the stream. while let Some(n) = stream.recv.read(&mut chunk).await? { file.write_all(&chunk[..n]).await?; received += n as u64; + if let Some((reader, writer, started)) = &mut streaming { + writer.add_available(n as u64); + if !*started + && received >= stream_start_at + && let (Some(reader), Some(callback)) = (reader.take(), stream_start.as_mut()) + { + callback(StreamingStart { + reader, + mime_type: header.mime_type.clone(), + }); + *started = true; + } + } + progress(DownloadProgress { + stream_key, + transport_phase: "download", + received, + total: header.total_size, + transport: None, + }); + } + if let Some((reader, writer, started)) = &mut streaming { + if !*started + && received > 0 + && let (Some(reader), Some(callback)) = (reader.take(), stream_start.as_mut()) + { + callback(StreamingStart { + reader, + mime_type: header.mime_type.clone(), + }); + *started = true; + } + writer.finish(); } file.flush().await?; drop(file); @@ -374,6 +470,27 @@ pub async fn download_track( ); } tokio::fs::rename(&temp_path, &path).await?; + progress(DownloadProgress { + stream_key, + transport_phase: "download-done", + received, + total: header.total_size, + transport: Some(stream.connection_stats()), + }); + let elapsed = started.elapsed(); + let kib_per_sec = if elapsed.as_secs_f64() > 0.0 { + received as f64 / 1024.0 / elapsed.as_secs_f64() + } else { + 0.0 + }; + tracing::info!( + owner = %owner, + item_id = %item_id_hex, + bytes = received, + elapsed_ms = elapsed.as_millis(), + kib_per_sec, + "federated audio downloaded" + ); Ok(Downloaded { path, mime_type: header.mime_type, @@ -383,6 +500,14 @@ pub async fn download_track( }) } +fn stream_start_buffer(total_size: u64) -> u64 { + if total_size == 0 { + STREAM_START_BUFFER_BYTES + } else { + total_size.min(STREAM_START_BUFFER_BYTES) + } +} + // --------------------------------------------------------------------------- // Serving side: answer audio requests from other peers // --------------------------------------------------------------------------- @@ -466,19 +591,37 @@ fn resolve_for_serving( /// Runs the accept loop of the audio protocol until the acceptor closes. /// Every track of the local library is downloadable by every peer of the /// network — the libraries of all participants are equal. -pub async fn serve_peers(mut acceptor: StreamAcceptor, library: Arc, own: EndpointId) { +pub async fn serve_peers( + mut acceptor: StreamAcceptor, + library: Arc, + own: EndpointId, + transport_stats: Arc, +) { while let Some(stream) = acceptor.accept().await { let library = Arc::clone(&library); + let transport_stats = Arc::clone(&transport_stats); tokio::spawn(async move { let peer = stream.peer_id; - if let Err(err) = serve_one(stream, library, own).await { + if let Err(err) = serve_one(stream, library, own, transport_stats).await { tracing::warn!(peer = %peer, "audio stream failed: {err:#}"); } }); } } -async fn serve_one(mut stream: ByteStream, library: Arc, own: EndpointId) -> Result<()> { +async fn serve_one( + mut stream: ByteStream, + library: Arc, + own: EndpointId, + transport_stats: Arc, +) -> Result<()> { + crate::federation::record_stream_transport( + &transport_stats, + "audio", + "inbound", + "open", + &stream, + ); let request: AudioRequest = serde_json::from_slice(&read_line(&mut stream.recv).await?)?; tracing::info!( peer = %stream.peer_id, @@ -555,17 +698,54 @@ async fn serve_one(mut stream: ByteStream, library: Arc, own: EndpointI let _ = stream.send.stopped().await; return Ok(()); } + let cover_bytes = cover.as_ref().map_or(0, |(bytes, _)| bytes.len() as u64); + let artist_image_bytes = artist_image + .as_ref() + .map_or(0, |(bytes, _)| bytes.len() as u64); + tracing::info!( + peer = %stream.peer_id, + item = %request.item_id, + audio_bytes = total_size.saturating_sub(offset), + cover_bytes, + artist_image_bytes, + want_images = request.want_cover, + "serving federated audio stream" + ); + let started = Instant::now(); if let Some((bytes, _)) = &cover { stream.send.write_all(bytes).await?; } if let Some((bytes, _)) = &artist_image { stream.send.write_all(bytes).await?; } - tokio::io::copy(&mut file, &mut stream.send).await?; + let audio_sent = tokio::io::copy(&mut file, &mut stream.send).await?; stream.send.finish()?; // Wait until the peer read everything (or gave up) before dropping the // stream, otherwise the tail of the file is lost. let _ = stream.send.stopped().await; + crate::federation::record_stream_transport( + &transport_stats, + "audio", + "inbound", + "done", + &stream, + ); + let elapsed = started.elapsed(); + let total_sent = cover_bytes + artist_image_bytes + audio_sent; + let kib_per_sec = if elapsed.as_secs_f64() > 0.0 { + total_sent as f64 / 1024.0 / elapsed.as_secs_f64() + } else { + 0.0 + }; + tracing::info!( + peer = %stream.peer_id, + item = %request.item_id, + audio_bytes = audio_sent, + total_bytes = total_sent, + elapsed_ms = elapsed.as_millis(), + kib_per_sec, + "served federated audio stream" + ); Ok(()) } diff --git a/src/federation/catalog.rs b/src/federation/catalog.rs index 0f163b4..9d5c834 100644 --- a/src/federation/catalog.rs +++ b/src/federation/catalog.rs @@ -124,19 +124,37 @@ pub struct CatalogTrack { // --------------------------------------------------------------------------- /// Runs the catalog accept loop until the acceptor closes. -pub async fn serve_peers(mut acceptor: StreamAcceptor, library: Arc, own: EndpointId) { +pub async fn serve_peers( + mut acceptor: StreamAcceptor, + library: Arc, + own: EndpointId, + transport_stats: Arc, +) { while let Some(stream) = acceptor.accept().await { let library = Arc::clone(&library); + let transport_stats = Arc::clone(&transport_stats); tokio::spawn(async move { let peer = stream.peer_id; - if let Err(err) = serve_one(stream, library, own).await { + if let Err(err) = serve_one(stream, library, own, transport_stats).await { tracing::warn!(peer = %peer, "catalog request failed: {err:#}"); } }); } } -async fn serve_one(mut stream: ByteStream, library: Arc, own: EndpointId) -> Result<()> { +async fn serve_one( + mut stream: ByteStream, + library: Arc, + own: EndpointId, + transport_stats: Arc, +) -> Result<()> { + crate::federation::record_stream_transport( + &transport_stats, + "catalog", + "inbound", + "open", + &stream, + ); let request: CatalogRequest = serde_json::from_slice(&super::audio::read_line(&mut stream.recv).await?)?; tracing::info!( @@ -190,6 +208,13 @@ async fn serve_one(mut stream: ByteStream, library: Arc, own: EndpointI } stream.send.finish()?; let _ = stream.send.stopped().await; + crate::federation::record_stream_transport( + &transport_stats, + "catalog", + "inbound", + "done", + &stream, + ); Ok(()) } @@ -335,11 +360,19 @@ pub async fn fetch_catalog( service: &MusicDhtService, owner: EndpointId, artist: &str, + transport_stats: &Arc, ) -> Result { let mut stream = service .open_stream(owner, CATALOG_ALPN) .await .map_err(|err| anyhow::anyhow!("cannot reach the peer: {err}"))?; + crate::federation::record_stream_transport( + transport_stats, + "catalog", + "outbound", + "open", + &stream, + ); let mut line = serde_json::to_vec(&CatalogRequest { artist: artist.to_string(), want: None, @@ -360,6 +393,13 @@ pub async fn fetch_catalog( ); let response: CatalogResponse = serde_json::from_slice(&payload).context("malformed catalog response")?; + crate::federation::record_stream_transport( + transport_stats, + "catalog", + "outbound", + "done", + &stream, + ); if !response.ok { anyhow::bail!( "peer refused the catalog: {}", @@ -379,11 +419,19 @@ pub async fn fetch_image( owner: EndpointId, artist: &str, release: Option<&str>, + transport_stats: &Arc, ) -> Result, &'static str)>> { let mut stream = service .open_stream(owner, CATALOG_ALPN) .await .map_err(|err| anyhow::anyhow!("cannot reach the peer: {err}"))?; + crate::federation::record_stream_transport( + transport_stats, + "catalog", + "outbound", + "open", + &stream, + ); let mut line = serde_json::to_vec(&CatalogRequest { artist: artist.to_string(), want: Some(if release.is_some() { @@ -414,6 +462,13 @@ pub async fn fetch_image( .read_exact(&mut bytes) .await .context("stream ended inside the image")?; + crate::federation::record_stream_transport( + transport_stats, + "catalog", + "outbound", + "done", + &stream, + ); let extension = match header.mime_type.as_str() { "image/png" => "png", "image/webp" => "webp", diff --git a/src/federation/mod.rs b/src/federation/mod.rs index a4351a4..3a3f993 100644 --- a/src/federation/mod.rs +++ b/src/federation/mod.rs @@ -15,6 +15,7 @@ mod audio; pub mod catalog; +use std::collections::{HashMap, VecDeque}; use std::path::{Path, PathBuf}; use std::str::FromStr; use std::sync::Arc; @@ -23,8 +24,9 @@ use std::time::Duration; use anyhow::{Context, Result}; use music_dht::{ - EndpointId, ItemKind, ItemSpec, LibraryItem, MusicDhtConfig, MusicDhtService, NetworkId, - PeerTicket, PublishStats, RendezvousConfig, SyncStats, + ByteStream, ByteStreamConnectionStats, EndpointId, ItemKind, ItemSpec, LibraryItem, + MusicDhtConfig, MusicDhtService, NetworkId, PeerTicket, PublishStats, RendezvousConfig, + SyncStats, }; use rusqlite::{Connection, OpenFlags, params}; use serde::{Deserialize, Serialize}; @@ -32,7 +34,7 @@ use serde::{Deserialize, Serialize}; use crate::library::Library; use crate::library::models::{ArtistRef, TrackItem}; -pub use audio::{AUDIO_ALPN, TrackMetadata}; +pub use audio::{AUDIO_ALPN, DownloadProgress, StreamingStart, TrackMetadata}; pub use catalog::{CATALOG_ALPN, FedAppearsOn, FedArtistCard, FedCardTrack, FedRelease}; /// How often the published library is re-synchronized with the local index. @@ -43,12 +45,230 @@ const SYNC_INTERVAL: Duration = Duration::from_secs(60); const CONTENT_LOOKUP_ATTEMPTS: usize = 3; /// Pause between share-link content lookup attempts. const CONTENT_LOOKUP_RETRY_DELAY: Duration = Duration::from_secs(2); +const TRANSPORT_SAMPLE_LIMIT: usize = 16; /// Ephemeral (not-in-library) tracks get negative ids so the rest of the /// app can tell them apart from library rows (history, likes and release /// navigation skip them). static NEXT_EPHEMERAL_ID: AtomicI64 = AtomicI64::new(-1); +#[derive(Debug, Clone, Default)] +pub struct TransportSample { + pub at: String, + pub protocol: &'static str, + pub direction: &'static str, + pub phase: &'static str, + pub peer_id: String, + pub selected_path: String, + pub open_paths: usize, + pub direct_paths: usize, + pub relay_paths: usize, + pub custom_paths: usize, + pub selected_rtt_ms: Option, + pub selected_tx_bytes: u64, + pub selected_rx_bytes: u64, + pub total_tx_bytes: u64, + pub total_rx_bytes: u64, + pub lost_packets: u64, + pub lost_bytes: u64, +} + +#[derive(Debug, Clone, Copy, Default)] +struct TransportTrafficCursor { + tx_bytes: u64, + rx_bytes: u64, + lost_packets: u64, + lost_bytes: u64, +} + +impl TransportSample { + fn from_stats( + protocol: &'static str, + direction: &'static str, + phase: &'static str, + stats: ByteStreamConnectionStats, + ) -> Self { + Self { + at: now_label(), + protocol, + direction, + phase, + peer_id: stats.peer_id.to_string(), + selected_path: stats.selected_path.as_str().to_string(), + open_paths: stats.open_paths, + direct_paths: stats.direct_paths, + relay_paths: stats.relay_paths, + custom_paths: stats.custom_paths, + selected_rtt_ms: stats + .selected_rtt + .map(|duration| duration.as_millis() as u64), + selected_tx_bytes: stats.selected_tx_bytes, + selected_rx_bytes: stats.selected_rx_bytes, + total_tx_bytes: stats.total_tx_bytes, + total_rx_bytes: stats.total_rx_bytes, + lost_packets: stats.lost_packets, + lost_bytes: stats.lost_bytes, + } + } +} + +#[derive(Debug, Clone, Default)] +pub struct TransportStatsSnapshot { + pub total_samples: u64, + pub direct_samples: u64, + pub relay_samples: u64, + pub custom_samples: u64, + pub unknown_samples: u64, + pub audio_samples: u64, + pub catalog_samples: u64, + pub sync_samples: u64, + pub runtime_tx_bytes: u64, + pub runtime_rx_bytes: u64, + pub runtime_lost_packets: u64, + pub runtime_lost_bytes: u64, + pub active_streams: usize, + pub last: Vec, +} + +#[derive(Debug, Default)] +struct TransportStatsState { + total_samples: u64, + direct_samples: u64, + relay_samples: u64, + custom_samples: u64, + unknown_samples: u64, + audio_samples: u64, + catalog_samples: u64, + sync_samples: u64, + runtime_tx_bytes: u64, + runtime_rx_bytes: u64, + runtime_lost_packets: u64, + runtime_lost_bytes: u64, + traffic_cursors: HashMap, + last: VecDeque, +} + +#[derive(Debug, Default)] +pub struct TransportStats { + inner: std::sync::Mutex, +} + +impl TransportStats { + fn reset(&self) { + *lock(&self.inner) = TransportStatsState::default(); + } + + pub fn record( + &self, + stream_key: u64, + protocol: &'static str, + direction: &'static str, + phase: &'static str, + stats: ByteStreamConnectionStats, + ) { + let sample = TransportSample::from_stats(protocol, direction, phase, stats); + let mut state = lock(&self.inner); + state.total_samples += 1; + let previous = if transport_phase_is_open(phase) { + TransportTrafficCursor::default() + } else { + state + .traffic_cursors + .get(&stream_key) + .copied() + .unwrap_or_default() + }; + state.runtime_tx_bytes = state + .runtime_tx_bytes + .saturating_add(sample.total_tx_bytes.saturating_sub(previous.tx_bytes)); + state.runtime_rx_bytes = state + .runtime_rx_bytes + .saturating_add(sample.total_rx_bytes.saturating_sub(previous.rx_bytes)); + state.runtime_lost_packets = state + .runtime_lost_packets + .saturating_add(sample.lost_packets.saturating_sub(previous.lost_packets)); + state.runtime_lost_bytes = state + .runtime_lost_bytes + .saturating_add(sample.lost_bytes.saturating_sub(previous.lost_bytes)); + state.traffic_cursors.insert( + stream_key, + TransportTrafficCursor { + tx_bytes: sample.total_tx_bytes, + rx_bytes: sample.total_rx_bytes, + lost_packets: sample.lost_packets, + lost_bytes: sample.lost_bytes, + }, + ); + if transport_phase_is_terminal(phase) { + state.traffic_cursors.remove(&stream_key); + } + match sample.selected_path.as_str() { + "direct" => state.direct_samples += 1, + "relay" => state.relay_samples += 1, + "custom" => state.custom_samples += 1, + _ => state.unknown_samples += 1, + } + match protocol { + "audio" => state.audio_samples += 1, + "catalog" => state.catalog_samples += 1, + "device-sync" => state.sync_samples += 1, + _ => {} + } + state.last.push_front(sample); + while state.last.len() > TRANSPORT_SAMPLE_LIMIT { + state.last.pop_back(); + } + } + + fn snapshot(&self) -> TransportStatsSnapshot { + let state = lock(&self.inner); + TransportStatsSnapshot { + total_samples: state.total_samples, + direct_samples: state.direct_samples, + relay_samples: state.relay_samples, + custom_samples: state.custom_samples, + unknown_samples: state.unknown_samples, + audio_samples: state.audio_samples, + catalog_samples: state.catalog_samples, + sync_samples: state.sync_samples, + runtime_tx_bytes: state.runtime_tx_bytes, + runtime_rx_bytes: state.runtime_rx_bytes, + runtime_lost_packets: state.runtime_lost_packets, + runtime_lost_bytes: state.runtime_lost_bytes, + active_streams: state.traffic_cursors.len(), + last: state.last.iter().cloned().collect(), + } + } +} + +fn transport_phase_is_terminal(phase: &str) -> bool { + phase == "done" || phase.ends_with("-done") +} + +fn transport_phase_is_open(phase: &str) -> bool { + phase == "open" || phase.ends_with("-open") +} + +pub(crate) fn stream_transport_key(stream: &ByteStream) -> u64 { + stream as *const ByteStream as usize as u64 +} + +pub fn record_stream_transport( + stats: &Arc, + protocol: &'static str, + direction: &'static str, + phase: &'static str, + stream: &ByteStream, +) { + stats.record( + stream_transport_key(stream), + protocol, + direction, + phase, + stream.connection_stats(), + ); +} + // --------------------------------------------------------------------------- // Settings (persisted in /federation.toml) // --------------------------------------------------------------------------- @@ -164,6 +384,7 @@ pub struct FedStatus { pub published_items: usize, pub last_sync: Option, pub last_error: Option, + pub transport: TransportStatsSnapshot, } /// Outcome of preparing a federated track for playback. @@ -196,6 +417,7 @@ pub struct Federation { running: tokio::sync::Mutex>, last_sync: std::sync::Mutex>, last_error: std::sync::Mutex>, + transport_stats: Arc, } #[derive(Debug, Clone)] @@ -318,6 +540,7 @@ impl Federation { running: tokio::sync::Mutex::new(None), last_sync: std::sync::Mutex::new(None), last_error: std::sync::Mutex::new(initial_error), + transport_stats: Arc::new(TransportStats::default()), }) } @@ -384,6 +607,7 @@ impl Federation { } std::fs::create_dir_all(&self.data_dir) .with_context(|| format!("creating {}", self.data_dir.display()))?; + self.transport_stats.reset(); let config = MusicDhtConfig::builder() .data_dir(&self.data_dir) @@ -432,6 +656,7 @@ impl Federation { audio_acceptor, Arc::clone(&self.library), service.endpoint_id(), + Arc::clone(&self.transport_stats), )); // Serve per-artist catalog requests (the federated artist card). let catalog_acceptor = service @@ -441,6 +666,7 @@ impl Federation { catalog_acceptor, Arc::clone(&self.library), service.endpoint_id(), + Arc::clone(&self.transport_stats), )); let sync_acceptor = service .stream_acceptor(crate::devices::SYNC_ALPN) @@ -449,11 +675,13 @@ impl Federation { sync_acceptor, Arc::clone(&self.devices), Arc::clone(&service), + Arc::clone(&self.transport_stats), )); let device_sync = Arc::clone(&self.devices); let device_service = Arc::clone(&service); + let device_transport = Arc::clone(&self.transport_stats); let device_tick_task = tokio::spawn(async move { - crate::devices::sync_loop(device_sync, device_service).await; + crate::devices::sync_loop(device_sync, device_service, device_transport).await; }); *guard = Some(Running { @@ -630,6 +858,7 @@ impl Federation { .map(|items| items.len()) .unwrap_or(0); } + status.transport = self.transport_stats.snapshot(); status } @@ -733,33 +962,36 @@ impl Federation { /// Resolves a share-link content id to one playable federated track. /// - /// Resolution order: the in-session metadata cache, the DHT content key - /// (retried — the DHT is eventually consistent, so a single lookup can - /// transiently come up short), then a name search by the link label: - /// records under the name keys carry content ids too, and failing an - /// exact match, a track whose artists and title all match the label is - /// the same song from another owner. + /// Resolution order: the in-session metadata cache, the DHT content key, + /// then an early name search by the link label while content lookup keeps + /// retrying. The DHT is eventually consistent, so a single lookup can + /// transiently come up short. Records under the name keys carry content + /// ids too, and failing an exact match, a track whose artists and title + /// all match the label is the same song from another owner. pub async fn track_by_content_id( &self, content_id: &str, label: Option<&str>, ) -> Result { + let content_id = + music_dht::normalize_content_id(content_id).context("invalid content id")?; let service = self.service().await?; let own = service.endpoint_id(); for cached in self.cached_metadata_snapshot() { - if cached.fed.content_id.as_deref() == Some(content_id) { + if cached.fed.content_id.as_deref() == Some(content_id.as_str()) { return Ok(cached.to_fed_track()); } } let mut queried_nodes = 0usize; + let mut tried_label_fallback = false; for attempt in 0..CONTENT_LOOKUP_ATTEMPTS { if attempt > 0 { tokio::time::sleep(CONTENT_LOOKUP_RETRY_DELAY).await; } let outcome = service - .search_content_id(content_id) + .search_content_id(&content_id) .await .map_err(|err| anyhow::anyhow!("federated content lookup failed: {err}"))?; queried_nodes = queried_nodes.max(outcome.queried_nodes); @@ -771,41 +1003,13 @@ impl Federation { { return Ok(fed_track_from_item(item, own)); } - } - - if let Some(label) = label { - let normalized = music_dht::normalize_name(label); - if !normalized.is_empty() { - let outcome = service - .search_network(label) - .await - .map_err(|err| anyhow::anyhow!("federated search failed: {err}"))?; - queried_nodes = queried_nodes.max(outcome.queried_nodes); - let candidates: Vec = outcome - .local_results - .into_iter() - .chain(outcome.network_results) - .filter(|item| item.kind == ItemKind::Track) - .collect(); - if let Some(item) = candidates - .iter() - .find(|item| item.content_id.as_deref() == Some(content_id)) + if !tried_label_fallback { + tried_label_fallback = true; + if let Some(item) = self + .track_by_content_label(&service, own, &content_id, label, &mut queried_nodes) + .await? { - return Ok(fed_track_from_item(item.clone(), own)); - } - // "feat" is an artifact of the label format ("A feat. B-Title"), - // not a token of any track record. - let tokens: Vec = music_dht::tokenize(&normalized) - .into_iter() - .filter(|token| token != "feat") - .collect(); - if !tokens.is_empty() - && let Some(item) = candidates.into_iter().find(|item| { - let item_tokens = item.search_tokens(); - tokens.iter().all(|token| item_tokens.contains(token)) - }) - { - return Ok(fed_track_from_item(item, own)); + return Ok(item); } } } @@ -816,6 +1020,89 @@ impl Federation { anyhow::bail!("no peers currently publish this shared track") } + async fn track_by_content_label( + &self, + service: &MusicDhtService, + own: EndpointId, + content_id: &str, + label: Option<&str>, + queried_nodes: &mut usize, + ) -> Result> { + let Some(label) = label else { + return Ok(None); + }; + let normalized = music_dht::normalize_name(label); + if normalized.is_empty() { + return Ok(None); + } + let outcome = service + .search_network(label) + .await + .map_err(|err| anyhow::anyhow!("federated search failed: {err}"))?; + *queried_nodes = (*queried_nodes).max(outcome.queried_nodes); + let candidates: Vec = outcome + .local_results + .into_iter() + .chain(outcome.network_results) + .filter(|item| item.kind == ItemKind::Track) + .collect(); + if let Some(item) = candidates + .iter() + .find(|item| item.content_id.as_deref() == Some(content_id)) + { + return Ok(Some(fed_track_from_item(item.clone(), own))); + } + // "feat" is an artifact of the label format ("A feat. B-Title"), + // not a token of any track record. + let tokens: Vec = music_dht::tokenize(&normalized) + .into_iter() + .filter(|token| token != "feat") + .collect(); + if !tokens.is_empty() + && let Some(item) = candidates.into_iter().find(|item| { + let item_tokens = item.search_tokens(); + tokens.iter().all(|token| item_tokens.contains(token)) + }) + { + return Ok(Some(fed_track_from_item(item, own))); + } + Ok(None) + } + + async fn source_by_content_id_for_playback( + &self, + fed: &FedTrack, + reason: String, + ) -> Result { + let content_id = fed + .content_id + .as_deref() + .and_then(music_dht::normalize_content_id) + .with_context(|| format!("{reason}; no content id fallback is available"))?; + let label = format!("{} {}", fed.artist_names.join(" "), fed.title); + let resolved = self.track_by_content_id(&content_id, Some(&label)).await?; + if resolved.owner == fed.owner && resolved.item_id == fed.item_id { + anyhow::bail!("{reason}"); + } + Ok(resolved) + } + + async fn local_track_by_content_id_for_playback( + self: &Arc, + content_id: &str, + ) -> Result> { + let Some(content_id) = music_dht::normalize_content_id(content_id) else { + return Ok(None); + }; + let library = Arc::clone(&self.library); + tokio::task::spawn_blocking(move || -> Result> { + Ok(library + .track_by_content_id(&content_id)? + .filter(|track| Path::new(&track.file_path).is_file())) + }) + .await? + } + /// Assembles the federated artist card: finds the peers holding the /// artist through the DHT, asks each for its catalog slice directly and /// merges the answers (missing/slow peers are skipped). Role-aware DHT @@ -878,11 +1165,12 @@ impl Federation { let mut requests = Vec::new(); for owner in owners { let service = Arc::clone(&service); + let transport_stats = Arc::clone(&self.transport_stats); let name = name.to_string(); requests.push(tokio::spawn(async move { let result = tokio::time::timeout( Duration::from_secs(5), - catalog::fetch_catalog(&service, owner, &name), + catalog::fetch_catalog(&service, owner, &name, &transport_stats), ) .await; match result { @@ -937,13 +1225,17 @@ impl Federation { running_network_id == network_id, "device invite belongs to a different federation network" ); - self.devices.connect_invite(service, invite).await + self.devices + .connect_invite(service, invite, Arc::clone(&self.transport_stats)) + .await } pub async fn device_sync_now(self: &Arc) -> Result<()> { self.ensure_connected_devices_enabled()?; let service = self.service().await?; - self.devices.sync_once(service).await + self.devices + .sync_once(service, Arc::clone(&self.transport_stats)) + .await } pub async fn connect(&self, ticket: &str) -> Result { @@ -998,7 +1290,7 @@ impl Federation { }; let fetched = tokio::time::timeout( Duration::from_secs(5), - catalog::fetch_image(&service, owner, artist, release), + catalog::fetch_image(&service, owner, artist, release, &self.transport_stats), ) .await; match fetched { @@ -1019,17 +1311,49 @@ impl Federation { /// Downloads a federated track straight into the local library /// (regardless of the save-on-listen setting) and returns the imported /// track. Own/already-local tracks resolve without downloading. - pub async fn download_to_library(self: &Arc, fed: &FedTrack) -> Result { - let playable = self.fetch_playable(fed, true).await?; + pub async fn download_to_library_with_progress( + self: &Arc, + fed: &FedTrack, + progress: F, + ) -> Result + where + F: FnMut(DownloadProgress) + Send, + { + let playable = self + .fetch_playable_with_progress(fed, true, true, progress, None) + .await?; Ok(playable.track) } /// Prepares a federated track for playback: local tracks resolve /// straight to the library; remote tracks are downloaded — into the /// library when save-on-listen is enabled, into the cache otherwise. - pub async fn prepare_playback(self: &Arc, fed: &FedTrack) -> Result { + pub async fn prepare_playback_with_progress( + self: &Arc, + fed: &FedTrack, + progress: F, + ) -> Result + where + F: FnMut(DownloadProgress) + Send, + { let save = self.settings().save_on_listen; - self.fetch_playable(fed, save).await + self.fetch_playable_with_progress(fed, save, false, progress, None) + .await + } + + pub async fn prepare_playback_streaming_with_progress( + self: &Arc, + fed: &FedTrack, + progress: F, + mut stream_start: S, + ) -> Result + where + F: FnMut(DownloadProgress) + Send, + S: FnMut(StreamingStart) + Send, + { + let save = self.settings().save_on_listen; + self.fetch_playable_with_progress(fed, save, false, progress, Some(&mut stream_start)) + .await } /// Fetches rich metadata for a federated track without downloading the @@ -1071,132 +1395,259 @@ impl Federation { Ok(enriched) } - async fn fetch_playable(self: &Arc, fed: &FedTrack, save: bool) -> Result { - let service = self.service().await?; - let item_id = - audio::hex_decode_item_id(&fed.item_id).context("malformed item id in the result")?; + async fn fetch_playable_with_progress( + self: &Arc, + fed: &FedTrack, + save: bool, + want_images: bool, + mut progress: F, + mut stream_start: Option<&mut (dyn FnMut(StreamingStart) + Send)>, + ) -> Result + where + F: FnMut(DownloadProgress) + Send, + { + let mut fed = fed.clone(); + let mut tried_content_lookup = false; - if fed.own { - let library = Arc::clone(&self.library); - let own_id = service.endpoint_id(); - let track = tokio::task::spawn_blocking(move || -> Result> { - let Some(track_id) = audio::resolve_local_track_id(&library, own_id, item_id)? - else { - return Ok(None); - }; - Ok(library.tracks_by_ids(&[track_id])?.into_iter().next()) - }) - .await?? - .context("this track is no longer in the local library")?; + loop { + if let Some(content_id) = fed + .content_id + .as_deref() + .and_then(music_dht::normalize_content_id) + && let Some(track) = self + .local_track_by_content_id_for_playback(&content_id) + .await? + { + return Ok(FedPlayable { + track, + imported: false, + }); + } + + let service = self.service().await?; + let item_id = match audio::hex_decode_item_id(&fed.item_id) { + Some(item_id) => item_id, + None if !tried_content_lookup => { + fed = self + .source_by_content_id_for_playback( + &fed, + "malformed item id in the result".to_string(), + ) + .await?; + tried_content_lookup = true; + continue; + } + None => anyhow::bail!("malformed item id in the result"), + }; + + if fed.own { + let library = Arc::clone(&self.library); + let own_id = service.endpoint_id(); + let track = tokio::task::spawn_blocking(move || -> Result> { + let Some(track_id) = audio::resolve_local_track_id(&library, own_id, item_id)? + else { + return Ok(None); + }; + Ok(library.tracks_by_ids(&[track_id])?.into_iter().next()) + }) + .await??; + if let Some(track) = track.filter(|track| Path::new(&track.file_path).is_file()) { + return Ok(FedPlayable { + track, + imported: false, + }); + } + if !tried_content_lookup { + fed = self + .source_by_content_id_for_playback( + &fed, + "this track is no longer in the local library".to_string(), + ) + .await?; + tried_content_lookup = true; + continue; + } + anyhow::bail!("this track is no longer in the local library"); + } + + let owner = match EndpointId::from_str(&fed.owner) { + Ok(owner) => owner, + Err(_) if !tried_content_lookup => { + fed = self + .source_by_content_id_for_playback( + &fed, + format!("malformed owner id '{}'", fed.owner), + ) + .await?; + tried_content_lookup = true; + continue; + } + Err(_) => anyhow::bail!("malformed owner id '{}'", fed.owner), + }; + let dir = if save { + &self.media_dir + } else { + &self.cache_dir + }; + tokio::fs::create_dir_all(dir).await?; + + let downloaded = match self + .download_track_with_fallback( + &service, + owner, + &fed, + dir, + want_images, + &mut progress, + stream_start + .as_mut() + .map(|callback| &mut **callback as &mut (dyn FnMut(StreamingStart) + Send)), + ) + .await + { + Ok(downloaded) => downloaded, + Err(err) if !tried_content_lookup && fed.content_id.is_some() => { + tracing::warn!( + owner = %fed.owner, + item_id = %fed.item_id, + "federated source failed; resolving by content id: {err:#}" + ); + fed = self + .source_by_content_id_for_playback( + &fed, + "federated source failed".to_string(), + ) + .await?; + tried_content_lookup = true; + continue; + } + Err(err) => return Err(err), + }; + tracing::info!( + path = %downloaded.path.display(), + mime = %downloaded.mime_type, + cover = downloaded.cover.is_some(), + "federated track downloaded" + ); + + if save { + let library = Arc::clone(&self.library); + let import_path = downloaded.path.clone(); + let import_metadata = downloaded.metadata.clone(); + let import_cover = downloaded.cover.clone(); + let artist_image = downloaded.artist_image.clone(); + let fed_item_id = fed.item_id.clone(); + let imported = tokio::task::spawn_blocking(move || -> Result> { + let mut import = crate::library::import::read_file(&import_path)?; + // The owner's database is more authoritative than whatever + // tags the file happens to carry (often none at all). + if let Some(meta) = &import_metadata { + apply_remote_metadata(&mut import, meta); + } + // Same for the cover: the peer's library cover wins over an + // embedded picture; embedded art stays as the fallback. + if import_cover.is_some() { + import.cover = import_cover; + } + let (track_id, _) = crate::library::import::upsert_track(&library, &import)?; + // A like that referenced the federated track moves onto the + // freshly imported local row. + if let Err(err) = library.transfer_fed_like(&fed_item_id, track_id) { + tracing::warn!(%err, "federated like transfer failed"); + } + // The owner's artist image fills the gap for a freshly + // created (or still image-less) main artist. + if let (Some((bytes, extension)), Some(artist_name)) = + (&artist_image, import.artists.first()) + && let Err(err) = save_artist_image(&library, artist_name, bytes, extension) + { + tracing::warn!(%err, "saving the artist image failed"); + } + Ok(library.tracks_by_ids(&[track_id])?.into_iter().next()) + }) + .await?; + match imported { + Ok(Some(track)) => { + return Ok(FedPlayable { + track, + imported: true, + }); + } + Ok(None) => {} + Err(err) => { + tracing::warn!( + "importing the downloaded track failed: {err:#}; playing from the file" + ); + } + } + } + + // Ephemeral playback: put the cover next to the cached audio so + // the views can show it. + let cover_path = match &downloaded.cover { + Some((bytes, extension)) => { + let path = downloaded.path.with_extension(format!("cover.{extension}")); + match tokio::fs::write(&path, bytes).await { + Ok(()) => Some(path.to_string_lossy().into_owned()), + Err(err) => { + tracing::warn!(%err, "saving the cover failed"); + None + } + } + } + None => None, + }; + let mut track = ephemeral_track(&fed, downloaded.metadata.as_ref(), &downloaded.path); + track.cover_path = cover_path; return Ok(FedPlayable { track, imported: false, }); } - - let owner = EndpointId::from_str(&fed.owner) - .map_err(|_| anyhow::anyhow!("malformed owner id '{}'", fed.owner))?; - let dir = if save { - &self.media_dir - } else { - &self.cache_dir - }; - tokio::fs::create_dir_all(dir).await?; - - let downloaded = self - .download_track_with_fallback(&service, owner, fed, dir) - .await?; - tracing::info!( - path = %downloaded.path.display(), - mime = %downloaded.mime_type, - cover = downloaded.cover.is_some(), - "federated track downloaded" - ); - - if save { - let library = Arc::clone(&self.library); - let import_path = downloaded.path.clone(); - let import_metadata = downloaded.metadata.clone(); - let import_cover = downloaded.cover.clone(); - let artist_image = downloaded.artist_image.clone(); - let fed_item_id = fed.item_id.clone(); - let imported = tokio::task::spawn_blocking(move || -> Result> { - let mut import = crate::library::import::read_file(&import_path)?; - // The owner's database is more authoritative than whatever - // tags the file happens to carry (often none at all). - if let Some(meta) = &import_metadata { - apply_remote_metadata(&mut import, meta); - } - // Same for the cover: the peer's library cover wins over an - // embedded picture; embedded art stays as the fallback. - if import_cover.is_some() { - import.cover = import_cover; - } - let (track_id, _) = crate::library::import::upsert_track(&library, &import)?; - // A like that referenced the federated track moves onto the - // freshly imported local row. - if let Err(err) = library.transfer_fed_like(&fed_item_id, track_id) { - tracing::warn!(%err, "federated like transfer failed"); - } - // The owner's artist image fills the gap for a freshly - // created (or still image-less) main artist. - if let (Some((bytes, extension)), Some(artist_name)) = - (&artist_image, import.artists.first()) - && let Err(err) = save_artist_image(&library, artist_name, bytes, extension) - { - tracing::warn!(%err, "saving the artist image failed"); - } - Ok(library.tracks_by_ids(&[track_id])?.into_iter().next()) - }) - .await?; - match imported { - Ok(Some(track)) => { - return Ok(FedPlayable { - track, - imported: true, - }); - } - Ok(None) => {} - Err(err) => { - tracing::warn!( - "importing the downloaded track failed: {err:#}; playing from the file" - ); - } - } - } - - // Ephemeral playback: put the cover next to the cached audio so the - // views can show it. - let cover_path = match &downloaded.cover { - Some((bytes, extension)) => { - let path = downloaded.path.with_extension(format!("cover.{extension}")); - match tokio::fs::write(&path, bytes).await { - Ok(()) => Some(path.to_string_lossy().into_owned()), - Err(err) => { - tracing::warn!(%err, "saving the cover failed"); - None - } - } - } - None => None, - }; - let mut track = ephemeral_track(fed, downloaded.metadata.as_ref(), &downloaded.path); - track.cover_path = cover_path; - Ok(FedPlayable { - track, - imported: false, - }) } - async fn download_track_with_fallback( + async fn download_track_with_fallback( &self, service: &MusicDhtService, owner: EndpointId, fed: &FedTrack, dir: &Path, - ) -> Result { + want_images: bool, + progress: &mut F, + mut stream_start: Option<&mut (dyn FnMut(StreamingStart) + Send)>, + ) -> Result + where + F: FnMut(DownloadProgress) + Send, + { let stem = download_stem(fed); - match audio::download_track(service, owner, &fed.item_id, dir, &stem).await { + let primary = { + let primary_stream_start = stream_start + .as_mut() + .map(|callback| &mut **callback as &mut (dyn FnMut(StreamingStart) + Send)); + audio::download_track_with_streaming( + service, + owner, + &fed.item_id, + dir, + &stem, + want_images, + |event| { + if let Some(stats) = event.transport { + self.transport_stats.record( + event.stream_key, + "audio", + "outbound", + event.transport_phase, + stats, + ); + } + progress(event); + }, + primary_stream_start, + ) + .await + }; + match primary { Ok(downloaded) => return Ok(downloaded), Err(primary_err) => { let Some(content_id) = fed.content_id.as_deref() else { @@ -1224,15 +1675,34 @@ impl Federation { continue; } let candidate_owner = item.owner; - match audio::download_track( - service, - candidate_owner, - &candidate_item_id, - dir, - &stem, - ) - .await - { + let fallback = { + let fallback_stream_start = stream_start.as_mut().map(|callback| { + &mut **callback as &mut (dyn FnMut(StreamingStart) + Send) + }); + audio::download_track_with_streaming( + service, + candidate_owner, + &candidate_item_id, + dir, + &stem, + want_images, + |event| { + if let Some(stats) = event.transport { + self.transport_stats.record( + event.stream_key, + "audio", + "outbound", + event.transport_phase, + stats, + ); + } + progress(event); + }, + fallback_stream_start, + ) + .await + }; + match fallback { Ok(downloaded) => { tracing::info!( owner = %candidate_owner, diff --git a/src/library/mod.rs b/src/library/mod.rs index ee3e6fb..231653e 100644 --- a/src/library/mod.rs +++ b/src/library/mod.rs @@ -1263,6 +1263,56 @@ impl Library { .optional()?) } + pub fn liked_content_ids(&self) -> Result> { + let conn = self.lock(); + let mut statement = conn.prepare( + "SELECT DISTINCT t.content_id + FROM likes k + JOIN tracks t ON t.id = k.track_id + WHERE t.content_id IS NOT NULL", + )?; + let rows = statement + .query_map([], |row| row.get::<_, String>(0))? + .collect::>>()?; + Ok(rows + .into_iter() + .filter_map(|content_id| music_dht::normalize_content_id(&content_id)) + .collect()) + } + + /// Returns the new liked state for this content id. + pub fn toggle_like_by_content_id(&self, content_id: &str) -> Result { + let Some(content_id) = music_dht::normalize_content_id(content_id) else { + anyhow::bail!("invalid content id"); + }; + let mut conn = self.lock(); + let tx = conn.transaction()?; + let removed_local = tx.execute( + "DELETE FROM likes + WHERE track_id IN ( + SELECT id FROM tracks WHERE content_id = ?1 + )", + [&content_id], + )?; + let removed_fed = + tx.execute("DELETE FROM fed_likes WHERE content_id = ?1", [&content_id])?; + if removed_local + removed_fed > 0 { + tx.commit()?; + return Ok(false); + } + let track_id: i64 = tx + .query_row( + "SELECT id FROM tracks WHERE content_id = ?1 ORDER BY id LIMIT 1", + [&content_id], + |row| row.get(0), + ) + .optional()? + .with_context(|| format!("no local track with content id {content_id}"))?; + tx.execute("INSERT INTO likes (track_id) VALUES (?1)", [track_id])?; + tx.commit()?; + Ok(true) + } + pub fn set_synced_like(&self, track_id: i64, liked: bool, liked_hlc_ms: i64) -> Result { let conn = self.lock(); let changed = if liked { @@ -1460,28 +1510,40 @@ impl Library { /// Toggles a like on a federated track; returns the resulting state. pub fn toggle_fed_like(&self, fed: &crate::federation::FedTrack) -> Result { - let conn = self.lock(); let content_id = fed .content_id .as_deref() .and_then(music_dht::normalize_content_id); + let mut conn = self.lock(); + let tx = conn.transaction()?; let removed = match content_id.as_deref() { - Some(content_id) => conn.execute( - "DELETE FROM fed_likes WHERE item_id = ?1 OR content_id = ?2", - params![fed.item_id, content_id], - )?, - None => conn.execute("DELETE FROM fed_likes WHERE item_id = ?1", [&fed.item_id])?, + Some(content_id) => { + let removed_fed = tx.execute( + "DELETE FROM fed_likes WHERE item_id = ?1 OR content_id = ?2", + params![fed.item_id, content_id], + )?; + let removed_local = tx.execute( + "DELETE FROM likes + WHERE track_id IN ( + SELECT id FROM tracks WHERE content_id = ?1 + )", + [content_id], + )?; + removed_fed + removed_local + } + None => tx.execute("DELETE FROM fed_likes WHERE item_id = ?1", [&fed.item_id])?, }; if removed > 0 { + tx.commit()?; return Ok(false); } if let Some(content_id) = content_id.as_deref() { - conn.execute( + tx.execute( "DELETE FROM fed_likes WHERE content_id = ?1 AND item_id != ?2", params![content_id, fed.item_id], )?; } - conn.execute( + tx.execute( "INSERT INTO fed_likes (item_id, owner, title, artist_names, featured_artist_names, year, duration_seconds, content_id, release_title, track_number, disc_number) @@ -1500,6 +1562,7 @@ impl Library { fed.disc_number, ], )?; + tx.commit()?; Ok(true) } @@ -1697,26 +1760,6 @@ impl Library { Ok(true) } - pub fn likes(&self) -> Result> { - let conn = self.lock(); - let mut statement = conn.prepare("SELECT track_id FROM likes")?; - let ids = statement - .query_map([], |row| row.get(0))? - .collect::>>()?; - Ok(ids) - } - - /// Returns the new liked state. - pub fn toggle_like(&self, track_id: i64) -> Result { - let conn = self.lock(); - let removed = conn.execute("DELETE FROM likes WHERE track_id = ?1", [track_id])?; - if removed > 0 { - return Ok(false); - } - conn.execute("INSERT INTO likes (track_id) VALUES (?1)", [track_id])?; - Ok(true) - } - pub fn add_history( &self, track_id: i64, @@ -2255,7 +2298,15 @@ mod tests { file_size_bytes: Some(1), cover: None, }; - import::upsert_track(lib, &import).unwrap().0 + let id = import::upsert_track(lib, &import).unwrap().0; + let content_id = format!("b3:{}", blake3::hash(import.file_path.as_bytes()).to_hex()); + lib.lock() + .execute( + "UPDATE tracks SET content_id = ?2 WHERE id = ?1", + params![id, content_id], + ) + .unwrap(); + id } #[test] @@ -2414,10 +2465,11 @@ mod tests { .unwrap(); assert_eq!(lib.playlist(playlist.id).unwrap().tracks.len(), 1); - assert!(lib.toggle_like(track_id).unwrap()); - assert_eq!(lib.likes().unwrap(), vec![track_id]); + let content_id = lib.track_content_id_by_id(track_id).unwrap().unwrap(); + assert!(lib.toggle_like_by_content_id(&content_id).unwrap()); + assert_eq!(lib.liked_content_ids().unwrap(), vec![content_id.clone()]); assert_eq!(lib.playlist(LIKES_PLAYLIST_ID).unwrap().tracks.len(), 1); - assert!(!lib.toggle_like(track_id).unwrap()); + assert!(!lib.toggle_like_by_content_id(&content_id).unwrap()); lib.remove_tracks_from_playlist(playlist.id, &[track_id]) .unwrap(); @@ -2432,6 +2484,8 @@ mod tests { let lib = test_library(); let old_id = add_track(&lib, "Old Local", "Artist", "Album"); let new_id = add_track(&lib, "New Local", "Artist", "Album"); + let old_content_id = lib.track_content_id_by_id(old_id).unwrap().unwrap(); + let new_content_id = lib.track_content_id_by_id(new_id).unwrap().unwrap(); let content_id = format!("b3:{}", "c".repeat(64)); let fed = crate::federation::FedTrack { item_id: "fed_item_order".to_string(), @@ -2448,8 +2502,8 @@ mod tests { disc_number: Some(1), }; - assert!(lib.toggle_like(old_id).unwrap()); - assert!(lib.toggle_like(new_id).unwrap()); + assert!(lib.toggle_like_by_content_id(&old_content_id).unwrap()); + assert!(lib.toggle_like_by_content_id(&new_content_id).unwrap()); assert!(lib.toggle_fed_like(&fed).unwrap()); { let conn = lib.lock(); @@ -2479,8 +2533,8 @@ mod tests { .collect(); assert_eq!(titles, vec!["Middle Fed", "New Local", "Old Local"]); - assert!(!lib.toggle_like(old_id).unwrap()); - assert!(lib.toggle_like(old_id).unwrap()); + assert!(!lib.toggle_like_by_content_id(&old_content_id).unwrap()); + assert!(lib.toggle_like_by_content_id(&old_content_id).unwrap()); { let conn = lib.lock(); conn.execute( diff --git a/src/main.rs b/src/main.rs index 17c7ccc..6824dc3 100644 --- a/src/main.rs +++ b/src/main.rs @@ -7,6 +7,7 @@ mod library; mod media; mod player; mod share; +mod streaming; mod ui; mod visualizer; diff --git a/src/player/mod.rs b/src/player/mod.rs index fce6c7e..399e203 100644 --- a/src/player/mod.rs +++ b/src/player/mod.rs @@ -15,8 +15,12 @@ use rodio::{Decoder, DeviceSinkBuilder, Player, stream::MixerDeviceSink}; pub use analyzer::AudioAnalysisSnapshot; -/// Local audio files are read straight from disk. -pub type TrackReader = std::io::BufReader; +pub trait TrackReadSeek: std::io::Read + std::io::Seek + Send + Sync {} + +impl TrackReadSeek for T where T: std::io::Read + std::io::Seek + Send + Sync {} + +/// Playback sources are either ordinary files or growing streaming cache files. +pub type TrackReader = Box; /// Perceptual volume: cubic mapping from percent to linear amplitude, so /// equal percent steps sound like equal loudness steps and low percentages @@ -40,14 +44,16 @@ pub enum PlayerEvent { enum Command { Play { - reader: Box, + reader: TrackReader, byte_len: Option, + mime_type: Option, + seekable: bool, volume: f32, }, /// Append the next track behind the current one without interrupting /// playback — rodio switches sources back to back (gapless-ish). Enqueue { - reader: Box, + reader: TrackReader, byte_len: Option, }, Pause, @@ -88,17 +94,26 @@ pub struct Controller { impl Controller { pub fn play(&self, reader: TrackReader, byte_len: Option, volume: f32) { let _ = self.tx.send(Command::Play { - reader: Box::new(reader), + reader, byte_len, + mime_type: None, + seekable: true, + volume, + }); + } + + pub fn play_stream(&self, reader: TrackReader, mime_type: Option, volume: f32) { + let _ = self.tx.send(Command::Play { + reader, + byte_len: None, + mime_type, + seekable: false, volume, }); } pub fn enqueue(&self, reader: TrackReader, byte_len: Option) { - let _ = self.tx.send(Command::Enqueue { - reader: Box::new(reader), - byte_len, - }); + let _ = self.tx.send(Command::Enqueue { reader, byte_len }); } pub fn pause(&self) { @@ -188,6 +203,8 @@ fn handle( Command::Play { reader, byte_len, + mime_type, + seekable, volume, } => { // The device is opened lazily on first playback so the app works @@ -211,11 +228,14 @@ fn handle( let mut builder = Decoder::builder() .with_data(reader) - .with_seekable(true) + .with_seekable(seekable) .with_gapless(true); if let Some(len) = byte_len { builder = builder.with_byte_len(len); } + if let Some(mime_type) = mime_type.as_deref() { + builder = builder.with_mime_type(mime_type); + } match builder.build() { Ok(decoder) => { shared.analysis.clear(); diff --git a/src/streaming.rs b/src/streaming.rs new file mode 100644 index 0000000..29c9122 --- /dev/null +++ b/src/streaming.rs @@ -0,0 +1,141 @@ +use std::fs::File; +use std::io::{self, Read, Seek, SeekFrom}; +use std::path::Path; +use std::sync::{Arc, Condvar, Mutex}; + +#[derive(Debug, Default)] +struct State { + available: u64, + complete: bool, +} + +#[derive(Debug, Default)] +struct Shared { + state: Mutex, + changed: Condvar, +} + +#[derive(Debug)] +pub struct GrowingFileReader { + file: File, + pos: u64, + shared: Arc, +} + +#[derive(Debug)] +pub struct GrowingFileWriter { + shared: Arc, +} + +pub fn growing_file(path: &Path) -> io::Result<(GrowingFileReader, GrowingFileWriter)> { + let file = File::open(path)?; + let shared = Arc::new(Shared::default()); + Ok(( + GrowingFileReader { + file, + pos: 0, + shared: Arc::clone(&shared), + }, + GrowingFileWriter { shared }, + )) +} + +impl GrowingFileWriter { + pub fn add_available(&self, bytes: u64) { + let mut state = self + .shared + .state + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + state.available = state.available.saturating_add(bytes); + self.shared.changed.notify_all(); + } + + pub fn finish(&self) { + let mut state = self + .shared + .state + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + state.complete = true; + self.shared.changed.notify_all(); + } +} + +impl Drop for GrowingFileWriter { + fn drop(&mut self) { + self.finish(); + } +} + +impl Read for GrowingFileReader { + fn read(&mut self, buf: &mut [u8]) -> io::Result { + if buf.is_empty() { + return Ok(0); + } + loop { + let mut state = self + .shared + .state + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + while self.pos >= state.available && !state.complete { + state = self + .shared + .changed + .wait(state) + .unwrap_or_else(std::sync::PoisonError::into_inner); + } + if self.pos >= state.available && state.complete { + return Ok(0); + } + let readable = (state.available - self.pos).min(buf.len() as u64) as usize; + drop(state); + let read = self.file.read(&mut buf[..readable])?; + self.pos = self.pos.saturating_add(read as u64); + if read > 0 { + return Ok(read); + } + } + } +} + +impl Seek for GrowingFileReader { + fn seek(&mut self, _: SeekFrom) -> io::Result { + Err(io::Error::new( + io::ErrorKind::Unsupported, + "streaming playback is not seekable yet", + )) + } +} + +#[cfg(test)] +mod tests { + use super::*; + use std::io::{Read, Write}; + + #[test] + fn growing_reader_reads_available_bytes_then_eof() { + let path = + std::env::temp_dir().join(format!("furumi-growing-reader-{}", std::process::id())); + let _ = std::fs::remove_file(&path); + std::fs::File::create(&path).unwrap(); + let (mut reader, writer) = growing_file(&path).unwrap(); + let mut output = std::fs::OpenOptions::new().write(true).open(&path).unwrap(); + + output.write_all(b"fur").unwrap(); + writer.add_available(3); + let mut first = [0u8; 3]; + reader.read_exact(&mut first).unwrap(); + assert_eq!(&first, b"fur"); + + output.write_all(b"umi").unwrap(); + writer.add_available(3); + writer.finish(); + let mut rest = Vec::new(); + reader.read_to_end(&mut rest).unwrap(); + assert_eq!(&rest, b"umi"); + drop(output); + let _ = std::fs::remove_file(&path); + } +} diff --git a/src/ui/federation.rs b/src/ui/federation.rs index 8a9077e..34a8534 100644 --- a/src/ui/federation.rs +++ b/src/ui/federation.rs @@ -16,6 +16,21 @@ pub fn draw(frame: &mut Frame, area: Rect, state: &AppState) { let inner = block.inner(area); frame.render_widget(block, area); + if inner.width >= 132 { + let desired_settings_width = ((inner.width as usize * 44) / 100).clamp(68, 92) as u16; + let settings_width = desired_settings_width.min(inner.width.saturating_sub(56)); + let [rows_area, _, status_area] = Layout::horizontal([ + Constraint::Length(settings_width), + Constraint::Length(2), + Constraint::Min(0), + ]) + .areas(inner); + + draw_settings_rows(frame, rows_area, state); + draw_status(frame, status_area, state); + return; + } + let rows_height = (settings_rows(state).len() + 5 + device_presence_sections(state).len()) as u16; let [rows_area, _, status_area] = Layout::vertical([ @@ -340,7 +355,7 @@ fn draw_row_enabled( height: 1, }; let marker = if selected { "▶ " } else { " " }; - let label_width = 48usize; + let label_width = settings_label_width(area.width); let line = Line::from(vec![ Span::styled( marker, @@ -366,27 +381,142 @@ fn draw_row_enabled( *y = (*y).saturating_add(1); } +fn settings_label_width(width: u16) -> usize { + let width = width as usize; + if width >= 88 { + 48 + } else { + width.saturating_sub(28).clamp(24, 48) + } +} + fn status_line(label: &str, value: String) -> Line<'static> { Line::from(vec![ - Span::styled(format!("{label:<22}"), theme::dim()), + Span::styled(format!("{label:<18}"), theme::dim()), Span::raw(value), ]) } -fn bytes_label(bytes: u64) -> String { +fn short_bytes_label(bytes: u64) -> String { if bytes >= 1024 * 1024 { - format!("{bytes} B ({:.1} MiB)", bytes as f64 / 1024.0 / 1024.0) + format!("{:.1} MiB", bytes as f64 / 1024.0 / 1024.0) } else if bytes >= 1024 { - format!("{bytes} B ({:.1} KiB)", bytes as f64 / 1024.0) + format!("{:.1} KiB", bytes as f64 / 1024.0) } else { format!("{bytes} B") } } +fn rtt_label(ms: Option) -> String { + ms.map(|ms| format!("{ms} ms")) + .unwrap_or_else(|| "rtt n/a".to_string()) +} + fn short_id(id: &str) -> String { id.chars().take(12).collect::() + "…" } +fn push_transport_status(lines: &mut Vec>, status: &crate::federation::FedStatus) { + lines.push(Line::default()); + lines.push(Line::styled("Iroh Transport", theme::header())); + let transport = &status.transport; + if transport.total_samples == 0 { + lines.push(status_line("Streams", "no samples yet".to_string())); + return; + } + let runtime_total = transport + .runtime_tx_bytes + .saturating_add(transport.runtime_rx_bytes); + lines.push(status_line( + "Runtime traffic", + format!( + "{} · tx {} · rx {} · active {}", + short_bytes_label(runtime_total), + short_bytes_label(transport.runtime_tx_bytes), + short_bytes_label(transport.runtime_rx_bytes), + transport.active_streams + ), + )); + if transport.runtime_lost_packets > 0 || transport.runtime_lost_bytes > 0 { + lines.push(status_line( + "Runtime loss", + format!( + "{} pkts · {}", + transport.runtime_lost_packets, + short_bytes_label(transport.runtime_lost_bytes) + ), + )); + } + lines.push(status_line( + "Samples", + format!( + "{} total · direct {} · relay {} · custom {} · unknown {}", + transport.total_samples, + transport.direct_samples, + transport.relay_samples, + transport.custom_samples, + transport.unknown_samples + ), + )); + lines.push(status_line( + "Protocols", + format!( + "audio {} · catalog {} · sync {}", + transport.audio_samples, transport.catalog_samples, transport.sync_samples + ), + )); + if let Some(sample) = transport.last.first() { + lines.push(status_line( + "Last stream", + format!( + "{} {} {} · {} · {}", + sample.protocol, + sample.direction, + sample.phase, + sample.selected_path, + rtt_label(sample.selected_rtt_ms) + ), + )); + lines.push(status_line( + "Last peer", + format!( + "{} · paths d/r/c/open {}/{}/{}/{}", + short_id(&sample.peer_id), + sample.direct_paths, + sample.relay_paths, + sample.custom_paths, + sample.open_paths + ), + )); + lines.push(status_line( + "Last bytes", + format!( + "sel {}/{} · total {}/{} · lost {} / {}", + short_bytes_label(sample.selected_tx_bytes), + short_bytes_label(sample.selected_rx_bytes), + short_bytes_label(sample.total_tx_bytes), + short_bytes_label(sample.total_rx_bytes), + sample.lost_packets, + short_bytes_label(sample.lost_bytes) + ), + )); + } + for sample in transport.last.iter().take(3) { + lines.push(Line::from(vec![ + Span::styled(format!("{:<14}", sample.at), theme::dim()), + Span::raw(format!( + "{} {} {} · {} · tx {} rx {}", + sample.protocol, + sample.direction, + sample.phase, + sample.selected_path, + short_bytes_label(sample.total_tx_bytes), + short_bytes_label(sample.total_rx_bytes) + )), + ])); + } +} + fn draw_status(frame: &mut Frame, area: Rect, state: &AppState) { let mut lines: Vec = vec![Line::styled("Status", theme::header())]; match &state.federation.status { @@ -407,39 +537,53 @@ fn draw_status(frame: &mut Frame, area: Rect, state: &AppState) { )); } Some(status) => { - lines.push(status_line("Node", "running".to_string())); - lines.push(status_line("Network", status.network.clone())); - lines.push(status_line("Endpoint ID", status.endpoint_id.clone())); - lines.push(status_line("DHT node ID", status.dht_node_id.clone())); + lines.push(status_line("Node", format!("running · {}", status.network))); + lines.push(status_line( + "Endpoint", + format!( + "{} · dht {}", + short_id(&status.endpoint_id), + short_id(&status.dht_node_id) + ), + )); let peers = if status.connected_peers.is_empty() { - "none yet".to_string() + format!("none · contacts {}", status.known_contacts) } else { - let names: Vec = - status.connected_peers.iter().map(|p| short_id(p)).collect(); - format!("{} — {}", status.connected_peers.len(), names.join(", ")) + let names: Vec = status + .connected_peers + .iter() + .take(3) + .map(|p| short_id(p)) + .collect(); + let more = status.connected_peers.len().saturating_sub(names.len()); + let more = if more > 0 { + format!(" +{more}") + } else { + String::new() + }; + format!( + "{} connected{} · contacts {} · {}", + status.connected_peers.len(), + more, + status.known_contacts, + names.join(", ") + ) }; - lines.push(status_line("Connected peers", peers)); + lines.push(status_line("Peers", peers)); lines.push(status_line( - "Known contacts", - status.known_contacts.to_string(), - )); - lines.push(status_line( - "Stored DHT records", - status - .stored_dht_records - .map(|count| count.to_string()) - .unwrap_or_else(|| "unavailable".to_string()), - )); - lines.push(status_line( - "Stored DHT bytes", - status - .stored_dht_bytes - .map(bytes_label) - .unwrap_or_else(|| "unavailable".to_string()), - )); - lines.push(status_line( - "Published items", - status.published_items.to_string(), + "DHT", + format!( + "{} records · {} · {} published", + status + .stored_dht_records + .map(|count| count.to_string()) + .unwrap_or_else(|| "unavailable".to_string()), + status + .stored_dht_bytes + .map(short_bytes_label) + .unwrap_or_else(|| "unavailable".to_string()), + status.published_items + ), )); lines.push(status_line( "Last sync", @@ -451,6 +595,7 @@ fn draw_status(frame: &mut Frame, area: Rect, state: &AppState) { if let Some(error) = &status.last_error { lines.push(status_line("Error", error.clone())); } + push_transport_status(&mut lines, status); } } lines.push(Line::default()); @@ -458,29 +603,32 @@ fn draw_status(frame: &mut Frame, area: Rect, state: &AppState) { match &state.federation.devices { None => lines.push(Line::styled("loading…", theme::dim())), Some(status) => { - lines.push(status_line("This device", status.this_device_id.clone())); - lines.push(status_line("Sync group", status.group_id.clone())); lines.push(status_line( - "Active devices", - status.active_devices.to_string(), - )); - lines.push(status_line( - "Revoked devices", - status.revoked_devices.to_string(), - )); - lines.push(status_line( - "Pending requests", - status.pending_requests.to_string(), - )); - lines.push(status_line("Ops in log", status.ops_total.to_string())); - lines.push(status_line( - "Tombstones", + "This device", format!( - "{} ({} compactable)", - status.tombstone_ops, status.compactable_tombstones + "{} · {}", + status.this_device_name, + short_id(&status.this_device_id) + ), + )); + lines.push(status_line("Sync group", short_id(&status.group_id))); + lines.push(status_line( + "Devices", + format!( + "{} active · {} revoked · {} pending", + status.active_devices, status.revoked_devices, status.pending_requests + ), + )); + lines.push(status_line( + "Sync log", + format!( + "{} ops · {} outbox · {} tombstones ({} gc)", + status.ops_total, + status.outbox_ops, + status.tombstone_ops, + status.compactable_tombstones ), )); - lines.push(status_line("Outbox ops", status.outbox_ops.to_string())); lines.push(status_line( "Snapshot", format!( diff --git a/src/ui/global.rs b/src/ui/global.rs index 523cf99..c405ccc 100644 --- a/src/ui/global.rs +++ b/src/ui/global.rs @@ -287,7 +287,7 @@ fn draw_grid(frame: &mut Frame, area: Rect, state: &AppState) { let message = if let Some(error) = &global.error { Line::styled(error.clone(), error_style()) } else if global.loading { - Line::styled("loading artists…", theme::dim()) + super::loading_line(state, "loading artists…") } else { Line::styled("no artists in the library", theme::dim()) }; @@ -386,7 +386,7 @@ fn draw_artist(frame: &mut Frame, area: Rect, state: &AppState, id: i64, cursor: Some(Loadable::Failed(error)) => { return centered_line(frame, inner, Line::styled(error.clone(), error_style())); } - _ => return centered_line(frame, inner, Line::styled("loading…", theme::dim())), + _ => return centered_line(frame, inner, super::loading_line(state, "loading…")), }; let header_height = (ART_HEADER_HEIGHT + 1).min(inner.height); @@ -653,7 +653,7 @@ fn draw_release(frame: &mut Frame, area: Rect, state: &AppState, id: i64, cursor Some(Loadable::Failed(error)) => { return centered_line(frame, inner, Line::styled(error.clone(), error_style())); } - _ => return centered_line(frame, inner, Line::styled("loading…", theme::dim())), + _ => return centered_line(frame, inner, super::loading_line(state, "loading…")), }; let header_height = (ART_HEADER_HEIGHT + 1).min(inner.height); @@ -751,7 +751,12 @@ fn draw_search(frame: &mut Frame, area: Rect, state: &AppState, cursor: usize) { } else { "searching…" }; - return centered_line(frame, inner, Line::styled(hint, theme::dim())); + let line = if search.query.is_empty() { + Line::styled(hint, theme::dim()) + } else { + super::loading_line(state, hint) + }; + return centered_line(frame, inner, line); } }; if results.len() == 0 @@ -795,7 +800,7 @@ fn draw_search(frame: &mut Frame, area: Rect, state: &AppState, cursor: usize) { if !results.tracks.is_empty() { rows.push((Line::styled("Tracks", theme::header()), None, None)); for track in &results.tracks { - let heart = if state.likes.contains(&track.id) { + let heart = if state.track_liked(track) { Span::styled("♥ ", theme::accent()) } else { Span::raw(" ") @@ -936,7 +941,7 @@ fn draw_fed_artist(frame: &mut Frame, area: Rect, state: &AppState, cursor: usiz return centered_line( frame, inner, - Line::styled("assembling the card from peers…", theme::dim()), + super::loading_line(state, "assembling the card from peers…"), ); } Loadable::Failed(message) => { diff --git a/src/ui/mod.rs b/src/ui/mod.rs index 680e190..9af2ef8 100644 --- a/src/ui/mod.rs +++ b/src/ui/mod.rs @@ -45,6 +45,26 @@ pub fn draw(frame: &mut Frame, state: &AppState, keymap: &Keymap) { popup::draw(frame, state); } +pub(crate) fn loading_line(state: &AppState, text: impl Into) -> Line<'static> { + Line::from(vec![ + Span::styled(format!("{} ", state.spinner()), theme::accent()), + Span::styled(text.into(), theme::dim()), + ]) +} + +fn is_waiting_message(message: &str) -> bool { + let normalized = message.to_ascii_lowercase(); + normalized.contains("loading") + || normalized.contains("searching") + || normalized.contains("fetching") + || normalized.contains("downloading") + || normalized.contains("locating") + || normalized.contains("importing") + || normalized.contains("waiting") + || normalized.contains("assembling") + || normalized.contains("resolving") +} + fn draw_tabs(frame: &mut Frame, area: Rect, state: &AppState) { let titles = Tab::ALL .iter() @@ -68,11 +88,31 @@ pub(crate) fn track_row( selected: bool, visual_selected: bool, ) { - let fed_liked = track - .fed - .as_ref() - .is_some_and(|fed| state.fed_track_liked(fed)); - let heart = if state.likes.contains(&track.id) || fed_liked { + track_row_with_like_marker( + frame, + area, + state, + track, + index_label, + selected, + visual_selected, + true, + ); +} + +pub(crate) fn track_row_with_like_marker( + frame: &mut Frame, + area: Rect, + state: &AppState, + track: &crate::library::models::TrackItem, + index_label: String, + selected: bool, + visual_selected: bool, + show_like_marker: bool, +) { + let heart = if !show_like_marker { + Span::raw("") + } else if state.track_liked(track) { Span::styled("♥ ", theme::accent()) } else { Span::raw(" ") @@ -311,7 +351,7 @@ fn draw_status(frame: &mut Frame, area: Rect, state: &AppState) { } else { spans.push(Span::styled("▶ ", theme::accent())); } - if state.likes.contains(&track.id) { + if state.track_liked(track) { spans.push(Span::styled("♥ ", theme::accent())); } spans.push(Span::raw(track.title.clone())); @@ -340,6 +380,10 @@ fn draw_status(frame: &mut Frame, area: Rect, state: &AppState) { } let message = match &state.status_message { + Some(message) if is_waiting_message(message) => Line::from(vec![ + Span::styled(format!("{} ", state.spinner()), theme::accent()), + Span::styled(message.clone(), theme::accent()), + ]), Some(message) => Line::styled(message.clone(), theme::accent()), None => match &state.player.current { // Idle line doubles as the current track's tech data display. diff --git a/src/ui/playlists.rs b/src/ui/playlists.rs index da3e1ce..61ff434 100644 --- a/src/ui/playlists.rs +++ b/src/ui/playlists.rs @@ -4,8 +4,8 @@ use ratatui::style::{Color, Style}; use ratatui::text::{Line, Span}; use ratatui::widgets::{Block, Paragraph}; -use super::{theme, track_row}; -use crate::app::state::{AppState, Loadable, TrackSelectionScope}; +use super::{loading_line, theme, track_row_with_like_marker}; +use crate::app::state::{AppState, LIKES_PLAYLIST_ID, Loadable, TrackSelectionScope}; use crate::app::update::playlist_tracks; pub fn draw(frame: &mut Frame, area: Rect, state: &AppState) { @@ -51,11 +51,7 @@ fn draw_list(frame: &mut Frame, area: Rect, state: &AppState) { ); } _ => { - return centered_line( - frame, - inner, - Line::styled("loading playlists…", theme::dim()), - ); + return centered_line(frame, inner, loading_line(state, "loading playlists…")); } }; if list.is_empty() { @@ -112,7 +108,7 @@ fn draw_opened(frame: &mut Frame, area: Rect, state: &AppState, id: i64, cursor: ); } let Some(tracks) = playlist_tracks(state, id) else { - return centered_line(frame, inner, Line::styled("loading…", theme::dim())); + return centered_line(frame, inner, loading_line(state, "loading…")); }; if tracks.is_empty() { return centered_line( @@ -133,7 +129,7 @@ fn draw_opened(frame: &mut Frame, area: Rect, state: &AppState, id: i64, cursor: width: inner.width, height: 1, }; - track_row( + track_row_with_like_marker( frame, row, state, @@ -143,6 +139,7 @@ fn draw_opened(frame: &mut Frame, area: Rect, state: &AppState, id: i64, cursor: state .track_selection .contains(&TrackSelectionScope::Playlist(id), index), + id != LIKES_PLAYLIST_ID, ); } } diff --git a/src/ui/popup.rs b/src/ui/popup.rs index face5b4..ed91f1f 100644 --- a/src/ui/popup.rs +++ b/src/ui/popup.rs @@ -197,7 +197,7 @@ fn draw_connected_devices(frame: &mut Frame, state: &AppState, cursor: usize) { ); frame.render_widget( Paragraph::new(Line::styled( - "enter: activate this device / control selected · esc close", + "enter: activate selected device / control active selected · esc close", theme::dim(), )) .alignment(Alignment::Center),