Compare commits

..
5 Commits
Author SHA1 Message Date
Ultradesu fd766cda24 Added catalog slices support
Build and Publish / Build and Publish Docker Image (push) Successful in 7m6s
2026-07-25 22:39:36 +03:00
Ultradesu 49000f716c Bump
Build and Publish / Build and Publish Docker Image (push) Successful in 4m3s
2026-07-25 00:47:16 +03:00
Ultradesu ee4990e2f1 Fixed control mode tracks art 2026-07-25 00:47:00 +03:00
Ultradesu e64b61c167 Added FED devices
Build and Publish / Build and Publish Docker Image (push) Successful in 4m11s
2026-07-24 17:08:00 +03:00
Ultradesu b737ced3fc Added FED devices
Build and Publish / Build and Publish Docker Image (push) Successful in 4m1s
2026-07-24 16:54:41 +03:00
6 changed files with 626 additions and 128 deletions
Generated
+9 -8
View File
@@ -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#a9012351dcdbdf8dbaa1f5dd71e498b4bc678d99"
dependencies = [
"blake3",
"data-encoding",
@@ -1845,7 +1845,7 @@ checksum = "e6d5a32815ae3f33302d95fdcb2ce17862f8c65363dcfd29360480ba1001fc9c"
[[package]]
name = "furumusic"
version = "0.7.1"
version = "0.8.3"
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#a9012351dcdbdf8dbaa1f5dd71e498b4bc678d99"
dependencies = [
"async-trait",
"blake3",
@@ -3766,12 +3766,13 @@ dependencies = [
[[package]]
name = "netlink-proto"
version = "0.12.0"
version = "0.12.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "b65d130ee111430e47eed7896ea43ca693c387f097dd97376bffafbf25812128"
checksum = "e6f7398dddf5f152d2a91a2921a134c6097056e292c0d4b9906007855e7cece6"
dependencies = [
"bytes",
"futures",
"futures-channel",
"futures-util",
"log",
"netlink-packet-core",
"netlink-sys",
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "furumusic"
version = "0.8.0"
version = "0.8.4"
edition = "2024"
description = "Reusable web-app boilerplate: auth, OIDC/SSO, admin panel, user management, i18n, PostgreSQL"
+265 -53
View File
@@ -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");
@@ -27,6 +29,7 @@ const PAIRING_WAIT_MS: i64 = 5 * 60 * 1000;
const PAIRING_RETRY_DELAY: Duration = Duration::from_secs(1);
const RESPONSE_DRAIN_TIMEOUT: Duration = Duration::from_secs(2);
const DEVICE_SYNC_INTERVAL: Duration = Duration::from_secs(2);
const LOCAL_SEED_RECHECK_MS: i64 = 60 * 1000;
const MAX_LINE: usize = 8 * 1024 * 1024;
const MAX_OPS_PER_BATCH: i64 = 1000;
@@ -456,6 +459,7 @@ pub async fn connect_invite(
pool: &sqlx::PgPool,
service: Arc<MusicDhtService>,
hub: Arc<PlayerDeviceHub>,
transport_stats: Arc<TransportStats>,
user_id: i64,
user_name: &str,
invite_link: &str,
@@ -470,6 +474,7 @@ pub async fn connect_invite(
pool,
Arc::clone(&service),
Arc::clone(&hub),
Arc::clone(&transport_stats),
user_id,
user_name,
&invite,
@@ -596,6 +601,7 @@ pub async fn revoke_device(pool: &sqlx::PgPool, user_id: i64, device_id: &str) -
pub async fn status(pool: &sqlx::PgPool, user_id: i64, user_name: &str) -> Result<FedDeviceStatus> {
let identity = ensure_identity(pool, user_id, user_name).await?;
maybe_seed_local_user_state(pool, user_id).await?;
let active_devices: i64 = sqlx::query_scalar(
"SELECT COUNT(*) FROM furumusic__fed_device
WHERE user_id = $1 AND trusted_at_ms IS NOT NULL AND revoked_at_ms IS NULL",
@@ -746,13 +752,18 @@ pub async fn record_track_like(
liked: bool,
) -> Result<()> {
if let Some(content_id) = track_content_id(pool, track_id).await? {
let fed = if liked {
synced_fed_track_for_track(pool, track_id, &content_id).await?
} else {
None
};
record_local_op(
pool,
user_id,
SyncOpPayload::TrackLikeSet {
content_id,
liked,
fed: None,
fed,
},
)
.await?;
@@ -828,7 +839,8 @@ pub async fn record_playlist_tracks_added(
let title = playlist_title(pool, playlist_id).await?.unwrap_or_default();
let sync_id = ensure_local_playlist_sync_id(pool, user_id, playlist_id, &title).await?;
let rows = playlist_track_content_positions(pool, playlist_id, track_ids).await?;
for (content_id, position) in rows {
for (track_id, content_id, position) in rows {
let fed = synced_fed_track_for_track(pool, track_id, &content_id).await?;
record_local_op(
pool,
user_id,
@@ -836,7 +848,7 @@ pub async fn record_playlist_tracks_added(
playlist_id: sync_id.clone(),
content_id,
position,
fed: None,
fed,
},
)
.await?;
@@ -876,14 +888,16 @@ pub async fn serve_peers(
pool: sqlx::PgPool,
service: Arc<MusicDhtService>,
hub: Arc<PlayerDeviceHub>,
transport_stats: Arc<TransportStats>,
) {
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:#}");
}
});
@@ -894,11 +908,19 @@ pub async fn sync_loop(
pool: sqlx::PgPool,
service: Arc<MusicDhtService>,
hub: Arc<PlayerDeviceHub>,
transport_stats: Arc<TransportStats>,
) {
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:#}");
}
}
@@ -908,6 +930,7 @@ pub async fn sync_once_all(
pool: &sqlx::PgPool,
service: Arc<MusicDhtService>,
hub: Arc<PlayerDeviceHub>,
transport_stats: Arc<TransportStats>,
) -> Result<()> {
let rows = sqlx::query(
"SELECT DISTINCT user_id FROM furumusic__fed_device
@@ -917,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?;
}
}
@@ -928,6 +959,7 @@ pub async fn sync_once(
pool: &sqlx::PgPool,
service: Arc<MusicDhtService>,
hub: Arc<PlayerDeviceHub>,
transport_stats: Arc<TransportStats>,
user_id: i64,
) -> Result<()> {
let devices = active_remote_devices(pool, user_id).await?;
@@ -939,6 +971,7 @@ pub async fn sync_once(
pool,
Arc::clone(&service),
Arc::clone(&hub),
Arc::clone(&transport_stats),
user_id,
&device,
)
@@ -960,6 +993,7 @@ async fn try_connect_invite(
pool: &sqlx::PgPool,
service: Arc<MusicDhtService>,
hub: Arc<PlayerDeviceHub>,
transport_stats: Arc<TransportStats>,
user_id: i64,
user_name: &str,
invite: &InviteWire,
@@ -976,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 {
@@ -993,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),
@@ -1060,7 +1096,9 @@ async fn serve_one(
pool: &sqlx::PgPool,
service: Arc<MusicDhtService>,
hub: Arc<PlayerDeviceHub>,
transport_stats: Arc<TransportStats>,
) -> Result<()> {
record_stream_transport(&transport_stats, "device-sync", "inbound", "open", &stream);
match read_msg(&mut stream).await? {
WireMessage::PairRequest {
invite_id,
@@ -1079,6 +1117,7 @@ async fn serve_one(
pool,
service,
hub,
transport_stats,
invite_id,
secret,
profile,
@@ -1102,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
@@ -1117,6 +1166,7 @@ async fn handle_pair_request(
pool: &sqlx::PgPool,
service: Arc<MusicDhtService>,
hub: Arc<PlayerDeviceHub>,
transport_stats: Arc<TransportStats>,
invite_id: String,
secret: String,
mut profile: DeviceProfileWire,
@@ -1207,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") => {}
@@ -1257,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(())
}
@@ -1266,6 +1318,7 @@ async fn handle_hello(
pool: &sqlx::PgPool,
service: Arc<MusicDhtService>,
hub: Arc<PlayerDeviceHub>,
transport_stats: Arc<TransportStats>,
group_id: String,
mut profile: DeviceProfileWire,
devices: Vec<DeviceProfileWire>,
@@ -1326,6 +1379,7 @@ async fn handle_hello(
)
.await?;
finish_response(&mut stream).await?;
record_stream_transport(&transport_stats, "device-sync", "inbound", "done", &stream);
Ok(())
}
@@ -1333,6 +1387,7 @@ async fn sync_device(
pool: &sqlx::PgPool,
service: Arc<MusicDhtService>,
hub: Arc<PlayerDeviceHub>,
transport_stats: Arc<TransportStats>,
user_id: i64,
device: &StoredDevice,
) -> Result<()> {
@@ -1347,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 {
@@ -1361,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,
@@ -1467,6 +1524,48 @@ async fn ensure_identity(pool: &sqlx::PgPool, user_id: i64, user_name: &str) ->
})
}
async fn maybe_seed_local_user_state(pool: &sqlx::PgPool, user_id: i64) -> Result<()> {
let last_seeded_at_ms: Option<i64> = sqlx::query_scalar(
"SELECT local_seeded_at_ms
FROM furumusic__fed_device_identity
WHERE user_id = $1",
)
.bind(user_id)
.fetch_optional(pool)
.await?;
let Some(last_seeded_at_ms) = last_seeded_at_ms else {
return Ok(());
};
if now_ms().saturating_sub(last_seeded_at_ms) < LOCAL_SEED_RECHECK_MS
&& !local_seed_needs_metadata_backfill(pool, user_id).await?
{
return Ok(());
}
seed_local_user_state(pool, user_id).await
}
async fn local_seed_needs_metadata_backfill(pool: &sqlx::PgPool, user_id: i64) -> Result<bool> {
sqlx::query_scalar(
"SELECT EXISTS (
SELECT 1 FROM furumusic__fed_state_like
WHERE user_id = $1
AND liked = true
AND fed_json IS NULL
AND local_track_id IS NOT NULL
) OR EXISTS (
SELECT 1 FROM furumusic__fed_state_playlist_item
WHERE user_id = $1
AND present = true
AND fed_json IS NULL
AND local_track_id IS NOT NULL
)",
)
.bind(user_id)
.fetch_one(pool)
.await
.map_err(Into::into)
}
async fn seed_local_user_state(pool: &sqlx::PgPool, user_id: i64) -> Result<()> {
let now = now_ms();
let like_rows = sqlx::query(
@@ -1487,17 +1586,24 @@ async fn seed_local_user_state(pool: &sqlx::PgPool, user_id: i64) -> Result<()>
continue;
};
let track_id: i64 = row.get("track_id");
let fed_json = synced_fed_track_for_track(pool, track_id, &content_id)
.await?
.map(serde_json::to_value)
.transpose()?;
sqlx::query(
"INSERT INTO furumusic__fed_state_like
(user_id, content_id, liked, hlc_ms, op_id, local_track_id, fed_json)
VALUES ($1, $2, true, $3, $4, $5, NULL)
ON CONFLICT (user_id, content_id) DO NOTHING",
VALUES ($1, $2, true, $3, $4, $5, $6)
ON CONFLICT (user_id, content_id) DO UPDATE SET
local_track_id = COALESCE(furumusic__fed_state_like.local_track_id, EXCLUDED.local_track_id),
fed_json = COALESCE(furumusic__fed_state_like.fed_json, EXCLUDED.fed_json)",
)
.bind(user_id)
.bind(&content_id)
.bind(now)
.bind(format!("local_seed:like:{track_id}"))
.bind(track_id)
.bind(fed_json)
.execute(pool)
.await?;
}
@@ -1512,22 +1618,8 @@ async fn seed_local_user_state(pool: &sqlx::PgPool, user_id: i64) -> Result<()>
.await?;
for playlist in playlists {
let playlist_id: i64 = playlist.get("id");
let sync_id = format!("webpl_{user_id}_{playlist_id}");
let title: String = playlist.get("title");
sqlx::query(
"INSERT INTO furumusic__fed_state_playlist
(user_id, playlist_id, local_playlist_id, title, deleted, hlc_ms, op_id)
VALUES ($1, $2, $3, $4, false, $5, $6)
ON CONFLICT (user_id, playlist_id) DO NOTHING",
)
.bind(user_id)
.bind(&sync_id)
.bind(playlist_id)
.bind(&title)
.bind(now)
.bind(format!("local_seed:playlist:{playlist_id}"))
.execute(pool)
.await?;
let sync_id = ensure_seed_playlist_sync_id(pool, user_id, playlist_id, &title, now).await?;
let items = sqlx::query(
"SELECT pt.track_id, pt.position, c.content_id
@@ -1549,12 +1641,18 @@ async fn seed_local_user_state(pool: &sqlx::PgPool, user_id: i64) -> Result<()>
};
let track_id: i64 = item.get("track_id");
let position: i32 = item.get("position");
let fed_json = synced_fed_track_for_track(pool, track_id, &content_id)
.await?
.map(serde_json::to_value)
.transpose()?;
sqlx::query(
"INSERT INTO furumusic__fed_state_playlist_item
(user_id, playlist_id, content_id, present, position, hlc_ms, op_id,
local_track_id, fed_json)
VALUES ($1, $2, $3, true, $4, $5, $6, $7, NULL)
ON CONFLICT (user_id, playlist_id, content_id) DO NOTHING",
VALUES ($1, $2, $3, true, $4, $5, $6, $7, $8)
ON CONFLICT (user_id, playlist_id, content_id) DO UPDATE SET
local_track_id = COALESCE(furumusic__fed_state_playlist_item.local_track_id, EXCLUDED.local_track_id),
fed_json = COALESCE(furumusic__fed_state_playlist_item.fed_json, EXCLUDED.fed_json)",
)
.bind(user_id)
.bind(&sync_id)
@@ -1563,6 +1661,7 @@ async fn seed_local_user_state(pool: &sqlx::PgPool, user_id: i64) -> Result<()>
.bind(now)
.bind(format!("local_seed:playlist_item:{playlist_id}:{track_id}"))
.bind(track_id)
.bind(fed_json)
.execute(pool)
.await?;
}
@@ -1593,17 +1692,16 @@ async fn seed_local_user_state(pool: &sqlx::PgPool, user_id: i64) -> Result<()>
.bind(user_id)
.fetch_one(pool)
.await?;
if missing_likes + missing_playlist_items == 0 {
sqlx::query(
"UPDATE furumusic__fed_device_identity
SET local_seeded_at_ms = $2
WHERE user_id = $1 AND local_seeded_at_ms = 0",
)
.bind(user_id)
.bind(now)
.execute(pool)
.await?;
} else {
sqlx::query(
"UPDATE furumusic__fed_device_identity
SET local_seeded_at_ms = $2
WHERE user_id = $1",
)
.bind(user_id)
.bind(now)
.execute(pool)
.await?;
if missing_likes + missing_playlist_items != 0 {
tracing::debug!(
user_id,
missing_likes,
@@ -2985,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<i64>, variant: &str) -> Option<String> {
file_id.map(|id| format!("/api/player/cover/{id}/{variant}"))
}
async fn web_playback_payload(
pool: &sqlx::PgPool,
state: &PlaybackStateWire,
@@ -3216,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
@@ -3241,7 +3344,7 @@ async fn web_track_json_by_id(
"release_id": row.get::<i64, _>("release_id"),
"release_title": row.get::<String, _>("release_title"),
"release_year": row.get::<Option<i32>, _>("release_year"),
"cover_url": serde_json::Value::Null,
"cover_url": cover_variant_url(row.get::<Option<i64>, _>("cover_file_id"), "medium"),
"stream_url": format!("/api/player/stream/{track_id}"),
"uploader_name": "Fed",
"audio_format": row.get::<Option<String>, _>("audio_format"),
@@ -3274,6 +3377,21 @@ async fn track_artist_json(
.collect())
}
async fn track_artist_names(pool: &sqlx::PgPool, track_id: i64, role: &str) -> Result<Vec<String>> {
let rows = sqlx::query(
"SELECT a.name::text AS name
FROM furumusic__track_artist ta
JOIN furumusic__artist a ON a.id = ta.artist_id
WHERE ta.track_id = $1 AND ta.role = $2
ORDER BY ta.position",
)
.bind(track_id)
.bind(role)
.fetch_all(pool)
.await?;
Ok(rows.into_iter().map(|row| row.get("name")).collect())
}
fn placeholder_playback_track(
track: &PlaybackTrack,
index: usize,
@@ -3387,6 +3505,45 @@ async fn track_content_id(pool: &sqlx::PgPool, track_id: i64) -> Result<Option<S
Ok(value.and_then(|value| music_dht::normalize_content_id(&value)))
}
async fn synced_fed_track_for_track(
pool: &sqlx::PgPool,
track_id: i64,
content_id: &str,
) -> Result<Option<SyncedFedTrack>> {
let Some(content_id) = music_dht::normalize_content_id(content_id) else {
return Ok(None);
};
let row = sqlx::query(
"SELECT t.title::text AS title, t.track_number, t.disc_number,
t.duration_seconds, r.title::text AS release_title, r.year AS release_year
FROM furumusic__track t
JOIN furumusic__release r ON r.id = t.release_id
WHERE t.id = $1",
)
.bind(track_id)
.fetch_optional(pool)
.await?;
let Some(row) = row else {
return Ok(None);
};
let duration = row.get::<f64, _>("duration_seconds");
Ok(Some(SyncedFedTrack {
item_id: format!("web:{track_id}"),
owner: "WEB".to_string(),
title: row.get("title"),
artist_names: track_artist_names(pool, track_id, "main").await?,
featured_artist_names: track_artist_names(pool, track_id, "featuring").await?,
year: row.get("release_year"),
duration_seconds: duration
.is_finite()
.then(|| duration.round().max(0.0) as i64),
content_id,
release_title: row.get("release_title"),
track_number: row.get("track_number"),
disc_number: row.get("disc_number"),
}))
}
async fn ensure_materialized_playlist(
pool: &sqlx::PgPool,
user_id: i64,
@@ -3414,17 +3571,36 @@ async fn ensure_materialized_playlist(
.bind(&now)
.fetch_one(pool)
.await?;
upsert_materialized_playlist_link(pool, user_id, playlist_sync_id, id, title).await?;
Ok(id)
}
async fn upsert_materialized_playlist_link(
pool: &sqlx::PgPool,
user_id: i64,
playlist_sync_id: &str,
local_playlist_id: i64,
title: &str,
) -> Result<()> {
sqlx::query(
"UPDATE furumusic__fed_state_playlist
SET local_playlist_id = $3
WHERE user_id = $1 AND playlist_id = $2",
"INSERT INTO furumusic__fed_state_playlist
(user_id, playlist_id, local_playlist_id, title, deleted, hlc_ms, op_id)
VALUES ($1, $2, $3, $4, false, 0, $5)
ON CONFLICT (user_id, playlist_id) DO UPDATE SET
local_playlist_id = COALESCE(furumusic__fed_state_playlist.local_playlist_id, EXCLUDED.local_playlist_id),
title = CASE
WHEN furumusic__fed_state_playlist.hlc_ms = 0 THEN EXCLUDED.title
ELSE furumusic__fed_state_playlist.title
END",
)
.bind(user_id)
.bind(playlist_sync_id)
.bind(id)
.bind(local_playlist_id)
.bind(title)
.bind(format!("local_materialized:{local_playlist_id}"))
.execute(pool)
.await?;
Ok(id)
Ok(())
}
async fn ensure_playlist_for_item(
@@ -3449,6 +3625,40 @@ async fn ensure_playlist_for_item(
ensure_materialized_playlist(pool, user_id, playlist_sync_id, None, &title).await
}
async fn ensure_seed_playlist_sync_id(
pool: &sqlx::PgPool,
user_id: i64,
playlist_id: i64,
title: &str,
hlc_ms: i64,
) -> Result<String> {
if let Some(sync_id) = local_playlist_sync_id(pool, user_id, playlist_id).await? {
return Ok(sync_id);
}
let sync_id = format!("webpl_{user_id}_{playlist_id}");
sqlx::query(
"INSERT INTO furumusic__fed_state_playlist
(user_id, playlist_id, local_playlist_id, title, deleted, hlc_ms, op_id)
VALUES ($1, $2, $3, $4, false, $5, $6)
ON CONFLICT (user_id, playlist_id) DO UPDATE SET
local_playlist_id = COALESCE(furumusic__fed_state_playlist.local_playlist_id, EXCLUDED.local_playlist_id),
title = CASE
WHEN furumusic__fed_state_playlist.hlc_ms = 0 THEN EXCLUDED.title
ELSE furumusic__fed_state_playlist.title
END,
deleted = false",
)
.bind(user_id)
.bind(&sync_id)
.bind(playlist_id)
.bind(title)
.bind(hlc_ms)
.bind(format!("local_seed:playlist:{playlist_id}"))
.execute(pool)
.await?;
Ok(sync_id)
}
async fn ensure_local_playlist_sync_id(
pool: &sqlx::PgPool,
user_id: i64,
@@ -3524,9 +3734,9 @@ async fn playlist_track_content_positions(
pool: &sqlx::PgPool,
playlist_id: i64,
track_ids: &[i64],
) -> Result<Vec<(String, i64)>> {
) -> Result<Vec<(i64, String, i64)>> {
let rows = sqlx::query(
"SELECT c.content_id, pt.position
"SELECT pt.track_id, c.content_id, pt.position
FROM furumusic__playlist_track pt
JOIN furumusic__track t ON t.id = pt.track_id
JOIN furumusic__media_file m ON m.id = t.audio_file_id
@@ -3542,9 +3752,10 @@ async fn playlist_track_content_positions(
Ok(rows
.into_iter()
.filter_map(|row| {
let track_id: i64 = row.get("track_id");
let content_id: String = row.get("content_id");
music_dht::normalize_content_id(&content_id)
.map(|content_id| (content_id, row.get::<i32, _>("position") as i64))
.map(|content_id| (track_id, content_id, row.get::<i32, _>("position") as i64))
})
.collect())
}
@@ -3707,6 +3918,7 @@ async fn ops_for_peer(
}
async fn snapshot(pool: &sqlx::PgPool, user_id: i64) -> Result<SyncSnapshot> {
maybe_seed_local_user_state(pool, user_id).await?;
let like_rows = sqlx::query(
"SELECT content_id, liked, hlc_ms, op_id, fed_json
FROM furumusic__fed_state_like
+170 -5
View File
@@ -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<MusicDhtService>,
@@ -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<u64>,
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<TransportSample>,
}
#[derive(Debug, Default)]
pub struct TransportStats {
inner: std::sync::Mutex<TransportStatsState>,
}
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::<Vec<_>>(),
})
}
}
pub fn record_stream_transport(
stats: &Arc<TransportStats>,
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<Option<Running>>,
last_sync: std::sync::Mutex<Option<String>>,
last_error: std::sync::Mutex<Option<String>>,
transport_stats: Arc<TransportStats>,
}
fn now_iso() -> String {
@@ -87,6 +240,7 @@ pub fn handle() -> Arc<Federation> {
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
+143 -61
View File
@@ -3,18 +3,24 @@
//! 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};
pub use music_dht::catalog::CATALOG_ALPN;
use music_dht::catalog::{
CatalogArtist, CatalogArtistPreview, CatalogImageHeader as ImageHeader, CatalogRelease,
CatalogRequest, CatalogResponse, CatalogTrack,
};
use music_dht::{ByteStream, EndpointId, ItemId, ItemKind, StreamAcceptor, normalize_name};
use serde::{Deserialize, Serialize};
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.
pub const CATALOG_ALPN: &[u8] = b"furumi-fd/catalog/1";
/// Maximum size of a JSON protocol line (request or response header).
const MAX_PROTOCOL_LINE: usize = 4096;
@@ -63,58 +69,6 @@ struct TrackMetadata {
disc_number: Option<i32>,
}
#[derive(Debug, Deserialize)]
struct CatalogRequest {
artist: String,
#[serde(default)]
want: Option<String>,
#[serde(default)]
release: Option<String>,
}
#[derive(Debug, Default, Serialize)]
struct CatalogResponse {
ok: bool,
#[serde(skip_serializing_if = "Option::is_none")]
error: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
artist: Option<CatalogArtist>,
}
#[derive(Debug, Default, Serialize)]
struct CatalogArtist {
name: String,
releases: Vec<CatalogRelease>,
}
#[derive(Debug, Default, Serialize)]
struct CatalogRelease {
title: String,
release_type: String,
year: Option<i32>,
tracks: Vec<CatalogTrack>,
}
#[derive(Debug, Default, Serialize)]
struct CatalogTrack {
title: String,
track_number: Option<i32>,
disc_number: Option<i32>,
duration_seconds: Option<f64>,
#[serde(skip_serializing_if = "Option::is_none")]
content_id: Option<String>,
item_id: String,
}
#[derive(Debug, Default, Serialize)]
struct ImageHeader {
ok: bool,
#[serde(skip_serializing_if = "Option::is_none")]
error: Option<String>,
mime_type: String,
size: u64,
}
// ---------------------------------------------------------------------------
// Framing helpers
// ---------------------------------------------------------------------------
@@ -274,6 +228,31 @@ async fn track_artist_image_file(pool: &PgPool, track_id: i64) -> Result<Option<
Ok(row.map(|row| (row.get(0), row.get(1))))
}
async fn track_catalog_artist_names(
pool: &PgPool,
track_id: i64,
) -> Result<(Vec<String>, Vec<String>)> {
let mut artists = Vec::new();
let mut featured = Vec::new();
let rows = sqlx::query(
"SELECT a.name, ta.role FROM furumusic__track_artist ta
JOIN furumusic__artist a ON a.id = ta.artist_id
WHERE ta.track_id = $1 ORDER BY ta.position",
)
.bind(track_id)
.fetch_all(pool)
.await?;
for row in rows {
let name: String = row.get(0);
match row.get::<String, _>(1).as_str() {
"featuring" => featured.push(name),
"main" => artists.push(name),
_ => {}
}
}
Ok((artists, featured))
}
async fn track_metadata(pool: &PgPool, track_id: i64) -> Result<Option<TrackMetadata>> {
let Some(track) = sqlx::query(
"SELECT t.title, t.track_number, t.disc_number, COALESCE(t.year, r.year),
@@ -344,13 +323,16 @@ pub async fn serve_audio(
pool: PgPool,
storage_dir: String,
own: EndpointId,
transport_stats: Arc<TransportStats>,
) {
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 +344,9 @@ async fn serve_audio_one(
pool: PgPool,
storage_dir: String,
own: EndpointId,
transport_stats: Arc<TransportStats>,
) -> 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 +449,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 +478,17 @@ pub async fn serve_catalog(
pool: PgPool,
storage_dir: String,
own: EndpointId,
transport_stats: Arc<TransportStats>,
) {
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 +500,9 @@ async fn serve_catalog_one(
pool: PgPool,
storage_dir: String,
own: EndpointId,
transport_stats: Arc<TransportStats>,
) -> 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,
@@ -527,7 +518,23 @@ async fn serve_catalog_one(
Err(err) => CatalogResponse {
ok: false,
error: Some(format!("catalog lookup failed: {err:#}")),
artist: None,
..CatalogResponse::default()
},
};
stream
.send
.write_all(&serde_json::to_vec(&response)?)
.await?;
}
Some("artists") => {
let cursor = request.cursor.clone();
let limit = request.limit.unwrap_or(64).clamp(1, 200);
let response = match build_artist_slice(&pool, cursor, limit).await {
Ok(response) => response,
Err(err) => CatalogResponse {
ok: false,
error: Some(format!("artist slice lookup failed: {err:#}")),
..CatalogResponse::default()
},
};
stream
@@ -569,7 +576,7 @@ async fn serve_catalog_one(
let response = CatalogResponse {
ok: false,
error: Some(format!("unknown request kind '{other}'")),
artist: None,
..CatalogResponse::default()
};
stream
.send
@@ -579,6 +586,7 @@ async fn serve_catalog_one(
}
stream.send.finish()?;
let _ = stream.send.stopped().await;
record_stream_transport(&transport_stats, "catalog", "inbound", "done", &stream);
Ok(())
}
@@ -595,7 +603,7 @@ async fn build_catalog(pool: &PgPool, own: &EndpointId, artist: &str) -> Result<
return Ok(CatalogResponse {
ok: false,
error: Some("artist not found in the library".to_string()),
artist: None,
..CatalogResponse::default()
});
};
let artist_id: i64 = artist_row.get(0);
@@ -630,8 +638,11 @@ async fn build_catalog(pool: &PgPool, own: &EndpointId, artist: &str) -> Result<
for row in track_rows {
let track_id: i64 = row.get(0);
let duration: f64 = row.get(4);
let (artists, featured_artists) = track_catalog_artist_names(pool, track_id).await?;
tracks.push(CatalogTrack {
title: row.get(1),
artists,
featured_artists,
track_number: row.get(2),
disc_number: row.get(3),
duration_seconds: (duration > 0.0).then_some(duration),
@@ -649,11 +660,82 @@ async fn build_catalog(pool: &PgPool, own: &EndpointId, artist: &str) -> Result<
Ok(CatalogResponse {
ok: true,
error: None,
artist: Some(CatalogArtist {
name: artist_row.get(1),
releases,
appears_on: Vec::new(),
}),
..CatalogResponse::default()
})
}
async fn build_artist_slice(
pool: &PgPool,
cursor: Option<String>,
limit: usize,
) -> Result<CatalogResponse> {
let offset = cursor
.as_deref()
.and_then(|value| value.parse::<i64>().ok())
.unwrap_or(0)
.max(0);
let rows = sqlx::query(
r#"SELECT a.name::text AS name,
mf.file_path::text AS image_path,
COALESCE(s.release_count, 0)::bigint AS release_count,
COALESCE(s.track_count, 0)::bigint AS track_count
FROM furumusic__artist a
LEFT JOIN furumusic__media_file mf ON mf.id = a.image_file_id
LEFT JOIN (
SELECT appearance.artist_id,
COUNT(DISTINCT appearance.release_id) FILTER (WHERE appearance.is_primary_release_artist) AS release_count,
COUNT(DISTINCT appearance.track_id) AS track_count
FROM (
SELECT ta.artist_id,
t.id AS track_id,
r.id AS release_id,
primary_release.artist_id IS NOT NULL AS is_primary_release_artist
FROM furumusic__track_artist ta
JOIN furumusic__track t ON t.id = ta.track_id AND t.is_hidden = false
JOIN furumusic__release r ON r.id = t.release_id AND r.is_hidden = false
LEFT JOIN furumusic__release_artist primary_release
ON primary_release.release_id = r.id
AND primary_release.artist_id = ta.artist_id
AND primary_release.position = 0
) appearance
GROUP BY appearance.artist_id
) s ON s.artist_id = a.id
WHERE a.is_hidden = false
AND COALESCE(s.track_count, 0) > 0
ORDER BY (COALESCE(s.release_count, 0) > 0) DESC,
COALESCE(s.release_count, 0) DESC,
COALESCE(s.track_count, 0) DESC,
a.name_sort
LIMIT $1 OFFSET $2"#,
)
.bind(limit as i64 + 1)
.bind(offset)
.fetch_all(pool)
.await?;
let mut artists = Vec::with_capacity(rows.len().min(limit));
let has_more = rows.len() > limit;
for row in rows.into_iter().take(limit) {
let name: String = row.get(0);
artists.push(CatalogArtistPreview {
artist_key: normalize_name(&name),
name,
image_path: row.get(1),
release_count: row.get(2),
track_count: row.get(3),
});
}
let next_cursor = has_more.then(|| (offset + artists.len() as i64).to_string());
Ok(CatalogResponse {
ok: true,
artists,
next_cursor,
..CatalogResponse::default()
})
}
+38
View File
@@ -2303,6 +2303,29 @@ tbody tr:hover {
<div class="probe-row"><span>Published items</span><strong x-text="(federationStatus.node && federationStatus.node.published_items) ?? '-'"></strong></div>
<div class="probe-row"><span>Last sync</span><strong x-text="federationStatus.last_sync || 'not yet'"></strong></div>
</div>
<div class="probe-table" x-show="fedTransport().total_samples > 0" style="margin-top:10px">
<div class="probe-row">
<span>Transport path</span>
<strong>
<span class="badge" :class="fedPathBadge(fedTransport().last_path)" x-text="fedTransport().last_path || 'unknown'"></span>
</strong>
</div>
<div class="probe-row"><span>RTT</span><strong x-text="fedRtt(fedTransport().last_rtt_ms)"></strong></div>
<div class="probe-row"><span>Path samples</span><strong x-text="`${fedTransport().direct_samples || 0} direct · ${fedTransport().relay_samples || 0} relay · ${fedTransport().custom_samples || 0} custom · ${fedTransport().unknown_samples || 0} unknown`"></strong></div>
<div class="probe-row"><span>Protocols</span><strong x-text="`${fedTransport().audio_samples || 0} audio · ${fedTransport().catalog_samples || 0} catalog · ${fedTransport().sync_samples || 0} sync`"></strong></div>
<div class="probe-row"><span>Last peer</span><strong x-text="fedShort(fedTransport().last_peer)"></strong></div>
</div>
<div class="probe-table" x-show="fedTransport().last && fedTransport().last.length" style="margin-top:10px">
<template x-for="(sample, index) in fedTransport().last.slice(0, 5)" :key="`${sample.at}-${sample.protocol}-${sample.direction}-${sample.phase}-${sample.peer_id}-${index}`">
<div class="probe-row">
<span x-text="`${sample.protocol} · ${sample.direction} · ${sample.phase}`"></span>
<strong>
<span class="badge" :class="fedPathBadge(sample.selected_path)" x-text="sample.selected_path || 'unknown'"></span>
<span x-text="` ${fedRtt(sample.selected_rtt_ms)} · tx ${formatBytes(sample.total_tx_bytes || 0)} · rx ${formatBytes(sample.total_rx_bytes || 0)} · lost ${formatBytes(sample.lost_bytes || 0)}`"></span>
</strong>
</div>
</template>
</div>
<p class="probe-intro muted" x-show="federationStatus.last_error" x-text="federationStatus.last_error"></p>
<div class="toolbar" style="margin-top:14px; flex-wrap:wrap; gap:8px">
<button class="btn" type="button" @click="loadFederation()" :disabled="federationLoading">
@@ -3268,6 +3291,21 @@ function adminV2() {
return id ? `${id.slice(0, 12)}` : '-';
},
fedTransport() {
return (this.federationStatus.node && this.federationStatus.node.transport) || {};
},
fedPathBadge(path) {
if (path === 'direct') return 'ok';
if (path === 'relay') return 'pending';
if (path === 'custom') return 'running';
return 'disabled';
},
fedRtt(ms) {
return ms != null ? `${Math.round(Number(ms))} ms` : '-';
},
async loadSettingsProbe(showErrors = true) {
this.settingsProbeLoading = true;
try {