diff --git a/Cargo.lock b/Cargo.lock index 0562bd4..7cbbf5e 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1575,9 +1575,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" dependencies = [ "serde", ] @@ -1718,7 +1718,7 @@ dependencies = [ [[package]] name = "federation-net" version = "0.1.0" -source = "git+https://gt.hexor.cy/ab/frid.git#e5353fa9b93d78be6cd811b8430d7dd5e725e6e5" +source = "git+https://gt.hexor.cy/ab/frid.git#085a4752da25a8d3fe7eac673af081a4a73c08bd" dependencies = [ "blake3", "data-encoding", @@ -1845,7 +1845,7 @@ checksum = "e6d5a32815ae3f33302d95fdcb2ce17862f8c65363dcfd29360480ba1001fc9c" [[package]] name = "furumusic" -version = "0.8.1" +version = "0.8.2" dependencies = [ "anyhow", "async-trait", @@ -3622,7 +3622,7 @@ dependencies = [ [[package]] name = "music-dht" version = "0.1.0" -source = "git+https://gt.hexor.cy/ab/frid.git#e5353fa9b93d78be6cd811b8430d7dd5e725e6e5" +source = "git+https://gt.hexor.cy/ab/frid.git#085a4752da25a8d3fe7eac673af081a4a73c08bd" dependencies = [ "async-trait", "blake3", diff --git a/src/federation/devices.rs b/src/federation/devices.rs index dfb2d66..2d43b68 100644 --- a/src/federation/devices.rs +++ b/src/federation/devices.rs @@ -18,6 +18,8 @@ use tokio::io::{AsyncRead, AsyncReadExt}; use crate::player::PlayerDeviceHub; +use super::{TransportStats, record_stream_transport}; + pub const SYNC_ALPN: &[u8] = b"furumi/sync/1"; const CLIENT_VERSION: &str = env!("CARGO_PKG_VERSION"); @@ -457,6 +459,7 @@ pub async fn connect_invite( pool: &sqlx::PgPool, service: Arc, hub: Arc, + transport_stats: Arc, user_id: i64, user_name: &str, invite_link: &str, @@ -471,6 +474,7 @@ pub async fn connect_invite( pool, Arc::clone(&service), Arc::clone(&hub), + Arc::clone(&transport_stats), user_id, user_name, &invite, @@ -884,14 +888,16 @@ pub async fn serve_peers( pool: sqlx::PgPool, service: Arc, hub: Arc, + transport_stats: Arc, ) { while let Some(stream) = acceptor.accept().await { let pool = pool.clone(); let service = Arc::clone(&service); let hub = Arc::clone(&hub); + let transport_stats = Arc::clone(&transport_stats); tokio::spawn(async move { let peer = stream.peer_id; - if let Err(err) = serve_one(stream, &pool, service, hub).await { + if let Err(err) = serve_one(stream, &pool, service, hub, transport_stats).await { tracing::warn!(peer = %peer, "web fed device sync stream failed: {err:#}"); } }); @@ -902,11 +908,19 @@ pub async fn sync_loop( pool: sqlx::PgPool, service: Arc, hub: Arc, + transport_stats: Arc, ) { let mut interval = tokio::time::interval(DEVICE_SYNC_INTERVAL); loop { interval.tick().await; - if let Err(err) = sync_once_all(&pool, Arc::clone(&service), Arc::clone(&hub)).await { + if let Err(err) = sync_once_all( + &pool, + Arc::clone(&service), + Arc::clone(&hub), + Arc::clone(&transport_stats), + ) + .await + { tracing::debug!("web fed device sync tick failed: {err:#}"); } } @@ -916,6 +930,7 @@ pub async fn sync_once_all( pool: &sqlx::PgPool, service: Arc, hub: Arc, + transport_stats: Arc, ) -> Result<()> { let rows = sqlx::query( "SELECT DISTINCT user_id FROM furumusic__fed_device @@ -925,7 +940,15 @@ pub async fn sync_once_all( .await?; for row in rows { let user_id: i64 = row.get("user_id"); - if let Err(err) = sync_once(pool, Arc::clone(&service), Arc::clone(&hub), user_id).await { + if let Err(err) = sync_once( + pool, + Arc::clone(&service), + Arc::clone(&hub), + Arc::clone(&transport_stats), + user_id, + ) + .await + { set_last_error(pool, user_id, Some(&format!("{err:#}"))).await?; } } @@ -936,6 +959,7 @@ pub async fn sync_once( pool: &sqlx::PgPool, service: Arc, hub: Arc, + transport_stats: Arc, user_id: i64, ) -> Result<()> { let devices = active_remote_devices(pool, user_id).await?; @@ -947,6 +971,7 @@ pub async fn sync_once( pool, Arc::clone(&service), Arc::clone(&hub), + Arc::clone(&transport_stats), user_id, &device, ) @@ -968,6 +993,7 @@ async fn try_connect_invite( pool: &sqlx::PgPool, service: Arc, hub: Arc, + transport_stats: Arc, user_id: i64, user_name: &str, invite: &InviteWire, @@ -984,6 +1010,7 @@ async fn try_connect_invite( let snapshot = snapshot(pool, user_id).await?; let playback = local_playback_snapshot(pool, Arc::clone(&hub), user_id, &identity).await; let mut stream = service.open_stream(peer, SYNC_ALPN).await?; + record_stream_transport(&transport_stats, "device-sync", "outbound", "open", &stream); write_msg( &mut stream, &WireMessage::PairRequest { @@ -1001,10 +1028,11 @@ async fn try_connect_invite( ) .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")?; + record_stream_transport(&transport_stats, "device-sync", "outbound", "done", &stream); + match response { WireMessage::PairResponse { accepted: true, group_id: Some(group_id), @@ -1068,7 +1096,9 @@ async fn serve_one( pool: &sqlx::PgPool, service: Arc, hub: Arc, + transport_stats: Arc, ) -> Result<()> { + record_stream_transport(&transport_stats, "device-sync", "inbound", "open", &stream); match read_msg(&mut stream).await? { WireMessage::PairRequest { invite_id, @@ -1087,6 +1117,7 @@ async fn serve_one( pool, service, hub, + transport_stats, invite_id, secret, profile, @@ -1110,7 +1141,17 @@ async fn serve_one( playback, } => { handle_hello( - stream, pool, service, hub, group_id, profile, devices, vector, ops, snapshot, + stream, + pool, + service, + hub, + transport_stats, + group_id, + profile, + devices, + vector, + ops, + snapshot, playback, ) .await @@ -1125,6 +1166,7 @@ async fn handle_pair_request( pool: &sqlx::PgPool, service: Arc, hub: Arc, + transport_stats: Arc, invite_id: String, secret: String, mut profile: DeviceProfileWire, @@ -1215,6 +1257,7 @@ async fn handle_pair_request( ) .await?; finish_response(&mut stream).await?; + record_stream_transport(&transport_stats, "device-sync", "inbound", "done", &stream); return Ok(()); } Some("accepted") => {} @@ -1265,6 +1308,7 @@ async fn handle_pair_request( ) .await?; finish_response(&mut stream).await?; + record_stream_transport(&transport_stats, "device-sync", "inbound", "done", &stream); Ok(()) } @@ -1274,6 +1318,7 @@ async fn handle_hello( pool: &sqlx::PgPool, service: Arc, hub: Arc, + transport_stats: Arc, group_id: String, mut profile: DeviceProfileWire, devices: Vec, @@ -1334,6 +1379,7 @@ async fn handle_hello( ) .await?; finish_response(&mut stream).await?; + record_stream_transport(&transport_stats, "device-sync", "inbound", "done", &stream); Ok(()) } @@ -1341,6 +1387,7 @@ async fn sync_device( pool: &sqlx::PgPool, service: Arc, hub: Arc, + transport_stats: Arc, user_id: i64, device: &StoredDevice, ) -> Result<()> { @@ -1355,6 +1402,7 @@ async fn sync_device( let snapshot = snapshot(pool, user_id).await?; let playback = local_playback_snapshot(pool, Arc::clone(&hub), user_id, &identity).await; let mut stream = service.open_stream(peer, SYNC_ALPN).await?; + record_stream_transport(&transport_stats, "device-sync", "outbound", "open", &stream); write_msg( &mut stream, &WireMessage::Hello { @@ -1369,10 +1417,11 @@ async fn sync_device( ) .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")?; + record_stream_transport(&transport_stats, "device-sync", "outbound", "done", &stream); + match response { WireMessage::SyncResponse { accepted: true, devices, @@ -3034,6 +3083,10 @@ fn json_f64(value: f64) -> serde_json::Value { .unwrap_or_else(|| serde_json::json!(0.0)) } +fn cover_variant_url(file_id: Option, variant: &str) -> Option { + file_id.map(|id| format!("/api/player/cover/{id}/{variant}")) +} + async fn web_playback_payload( pool: &sqlx::PgPool, state: &PlaybackStateWire, @@ -3265,7 +3318,8 @@ async fn web_track_json_by_id( "SELECT t.id, t.title::text AS title, t.track_number, t.disc_number, t.duration_seconds, r.id AS release_id, r.title::text AS release_title, r.year AS release_year, mf.audio_format, mf.audio_bitrate, - mf.audio_sample_rate, mf.audio_bit_depth, mf.file_size_bytes + mf.audio_sample_rate, mf.audio_bit_depth, mf.file_size_bytes, + COALESCE(t.cover_file_id, r.cover_file_id) AS cover_file_id FROM furumusic__track t JOIN furumusic__release r ON r.id = t.release_id LEFT JOIN furumusic__media_file mf ON mf.id = t.audio_file_id @@ -3290,7 +3344,7 @@ async fn web_track_json_by_id( "release_id": row.get::("release_id"), "release_title": row.get::("release_title"), "release_year": row.get::, _>("release_year"), - "cover_url": serde_json::Value::Null, + "cover_url": cover_variant_url(row.get::, _>("cover_file_id"), "medium"), "stream_url": format!("/api/player/stream/{track_id}"), "uploader_name": "Fed", "audio_format": row.get::, _>("audio_format"), diff --git a/src/federation/mod.rs b/src/federation/mod.rs index 9dbdb8d..ead27de 100644 --- a/src/federation/mod.rs +++ b/src/federation/mod.rs @@ -16,15 +16,15 @@ pub mod devices; mod serve; mod storage; -use std::collections::{HashMap, HashSet}; +use std::collections::{HashMap, HashSet, VecDeque}; use std::path::PathBuf; use std::sync::{Arc, OnceLock}; use std::time::Duration; use anyhow::{Context, Result}; use music_dht::{ - ItemKind, ItemSpec, MusicDhtConfig, MusicDhtService, NetworkId, PeerTicket, PublishStats, - RendezvousConfig, SyncStats, + ByteStream, ByteStreamConnectionStats, ItemKind, ItemSpec, MusicDhtConfig, MusicDhtService, + NetworkId, PeerTicket, PublishStats, RendezvousConfig, SyncStats, }; use serde_json::{Value, json}; use sqlx::PgPool; @@ -37,6 +37,7 @@ pub use serve::{AUDIO_ALPN, CATALOG_ALPN}; /// How often the published library is re-synchronized with the database. const SYNC_INTERVAL: Duration = Duration::from_secs(60); +const TRANSPORT_SAMPLE_LIMIT: usize = 16; struct Running { service: Arc, @@ -50,6 +51,157 @@ struct ContentHashJob { file_path: String, } +#[derive(Debug, Clone)] +struct TransportSample { + at: String, + protocol: &'static str, + direction: &'static str, + phase: &'static str, + peer_id: String, + selected_path: String, + open_paths: usize, + direct_paths: usize, + relay_paths: usize, + custom_paths: usize, + selected_rtt_ms: Option, + selected_tx_bytes: u64, + selected_rx_bytes: u64, + total_tx_bytes: u64, + total_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_iso(), + 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, 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, + last: VecDeque, +} + +#[derive(Debug, Default)] +pub struct TransportStats { + inner: std::sync::Mutex, +} + +impl TransportStats { + fn reset(&self) { + *lock(&self.inner) = TransportStatsState::default(); + } + + fn record( + &self, + 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; + 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) -> Value { + let state = lock(&self.inner); + let latest = state.last.front(); + json!({ + "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, + "last_path": latest.map(|sample| sample.selected_path.clone()), + "last_rtt_ms": latest.and_then(|sample| sample.selected_rtt_ms), + "last_peer": latest.map(|sample| sample.peer_id.clone()), + "last": state.last.iter().map(|sample| json!({ + "at": sample.at, + "protocol": sample.protocol, + "direction": sample.direction, + "phase": sample.phase, + "peer_id": sample.peer_id, + "selected_path": sample.selected_path, + "open_paths": sample.open_paths, + "direct_paths": sample.direct_paths, + "relay_paths": sample.relay_paths, + "custom_paths": sample.custom_paths, + "selected_rtt_ms": sample.selected_rtt_ms, + "selected_tx_bytes": sample.selected_tx_bytes, + "selected_rx_bytes": sample.selected_rx_bytes, + "total_tx_bytes": sample.total_tx_bytes, + "total_rx_bytes": sample.total_rx_bytes, + "lost_packets": sample.lost_packets, + "lost_bytes": sample.lost_bytes, + })).collect::>(), + }) + } +} + +pub fn record_stream_transport( + stats: &Arc, + protocol: &'static str, + direction: &'static str, + phase: &'static str, + stream: &ByteStream, +) { + stats.record(protocol, direction, phase, stream.connection_stats()); +} + pub struct Federation { /// Transport data directory; server-side DHT state and identity live in PostgreSQL. data_dir: PathBuf, @@ -61,6 +213,7 @@ pub struct Federation { running: tokio::sync::Mutex>, last_sync: std::sync::Mutex>, last_error: std::sync::Mutex>, + transport_stats: Arc, } fn now_iso() -> String { @@ -87,6 +240,7 @@ pub fn handle() -> Arc { running: tokio::sync::Mutex::new(None), last_sync: std::sync::Mutex::new(None), last_error: std::sync::Mutex::new(None), + transport_stats: Arc::new(TransportStats::default()), }) })) } @@ -189,6 +343,7 @@ impl Federation { let dht_storage = Arc::new(PostgresFederationStorage::new(pool.clone()).await?); let secret_key = dht_storage.load_or_create_secret_key().await?; + self.transport_stats.reset(); let config = MusicDhtConfig::builder() .data_dir(&self.data_dir) @@ -236,6 +391,7 @@ impl Federation { pool.clone(), storage_dir.clone(), service.endpoint_id(), + Arc::clone(&self.transport_stats), )); let catalog_acceptor = service .stream_acceptor(CATALOG_ALPN) @@ -245,6 +401,7 @@ impl Federation { pool.clone(), storage_dir, service.endpoint_id(), + Arc::clone(&self.transport_stats), )); let device_acceptor = service .stream_acceptor(devices::SYNC_ALPN) @@ -255,9 +412,14 @@ impl Federation { pool.clone(), Arc::clone(&service), Arc::clone(&device_hub), + Arc::clone(&self.transport_stats), + )); + let device_sync_task = tokio::spawn(devices::sync_loop( + pool, + Arc::clone(&service), + device_hub, + Arc::clone(&self.transport_stats), )); - let device_sync_task = - tokio::spawn(devices::sync_loop(pool, Arc::clone(&service), device_hub)); *guard = Some(Running { service, @@ -596,6 +758,7 @@ impl Federation { "connected_peers": peers, "known_contacts": service.known_peers().len(), "published_items": published, + "transport": self.transport_stats.snapshot(), }) } None => json!({ "running": false }), @@ -668,6 +831,7 @@ impl Federation { &pool, service, crate::player::PlayerDeviceHub::shared(), + Arc::clone(&self.transport_stats), user_id, user_name, invite, @@ -698,6 +862,7 @@ impl Federation { &pool, service, crate::player::PlayerDeviceHub::shared(), + Arc::clone(&self.transport_stats), user_id, ) .await diff --git a/src/federation/serve.rs b/src/federation/serve.rs index 7d4b400..1205246 100644 --- a/src/federation/serve.rs +++ b/src/federation/serve.rs @@ -3,6 +3,7 @@ //! with the furumi TUI client and any other furumi peer. use std::path::{Path, PathBuf}; +use std::sync::Arc; use anyhow::Result; use music_dht::{ByteStream, EndpointId, ItemId, ItemKind, StreamAcceptor}; @@ -11,6 +12,8 @@ use sqlx::PgPool; use sqlx::Row as _; use tokio::io::{AsyncRead, AsyncReadExt, AsyncSeekExt, AsyncWriteExt}; +use super::{TransportStats, record_stream_transport}; + /// ALPN of the peer-to-peer audio streaming protocol. pub const AUDIO_ALPN: &[u8] = b"furumi-fd/audio/1"; /// ALPN of the per-artist catalog protocol. @@ -344,13 +347,16 @@ pub async fn serve_audio( pool: PgPool, storage_dir: String, own: EndpointId, + transport_stats: Arc, ) { while let Some(stream) = acceptor.accept().await { let pool = pool.clone(); let storage_dir = storage_dir.clone(); + let transport_stats = Arc::clone(&transport_stats); tokio::spawn(async move { let peer = stream.peer_id; - if let Err(err) = serve_audio_one(stream, pool, storage_dir, own).await { + if let Err(err) = serve_audio_one(stream, pool, storage_dir, own, transport_stats).await + { tracing::warn!(peer = %peer, "federation audio stream failed: {err:#}"); } }); @@ -362,7 +368,9 @@ async fn serve_audio_one( pool: PgPool, storage_dir: String, own: EndpointId, + transport_stats: Arc, ) -> Result<()> { + 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, @@ -465,6 +473,7 @@ async fn serve_audio_one( // Wait until the peer read everything before dropping the stream, // otherwise the tail of the file is lost. let _ = stream.send.stopped().await; + record_stream_transport(&transport_stats, "audio", "inbound", "done", &stream); Ok(()) } @@ -493,13 +502,17 @@ pub async fn serve_catalog( pool: PgPool, storage_dir: String, own: EndpointId, + transport_stats: Arc, ) { while let Some(stream) = acceptor.accept().await { let pool = pool.clone(); let storage_dir = storage_dir.clone(); + let transport_stats = Arc::clone(&transport_stats); tokio::spawn(async move { let peer = stream.peer_id; - if let Err(err) = serve_catalog_one(stream, pool, storage_dir, own).await { + if let Err(err) = + serve_catalog_one(stream, pool, storage_dir, own, transport_stats).await + { tracing::warn!(peer = %peer, "federation catalog request failed: {err:#}"); } }); @@ -511,7 +524,9 @@ async fn serve_catalog_one( pool: PgPool, storage_dir: String, own: EndpointId, + transport_stats: Arc, ) -> Result<()> { + record_stream_transport(&transport_stats, "catalog", "inbound", "open", &stream); let request: CatalogRequest = serde_json::from_slice(&read_line(&mut stream.recv).await?)?; tracing::info!( peer = %stream.peer_id, @@ -579,6 +594,7 @@ async fn serve_catalog_one( } stream.send.finish()?; let _ = stream.send.stopped().await; + record_stream_transport(&transport_stats, "catalog", "inbound", "done", &stream); Ok(()) } diff --git a/templates/admin/v2.html b/templates/admin/v2.html index 46f75dd..c08f53a 100644 --- a/templates/admin/v2.html +++ b/templates/admin/v2.html @@ -2303,6 +2303,29 @@ tbody tr:hover {
Published items
Last sync
+
+
+ Transport path + + + +
+
RTT
+
Path samples
+
Protocols
+
Last peer
+
+
+ +