From ba1c565bdce590fa583dae80157d6ce1b3415058 Mon Sep 17 00:00:00 2001 From: AB Date: Thu, 10 Sep 2026 17:08:34 +0300 Subject: [PATCH] Integrate shared playback coordination and release 0.10.7 --- Cargo.lock | 6 +- Cargo.toml | 4 +- src/federation/devices.rs | 355 +++++++++++++++++++++++++++++----- src/federation/mod.rs | 5 + src/music/mod.rs | 2 + src/player/mod.rs | 341 ++++++++++++++++++++++++++------ templates/player/scripts.html | 5 +- 7 files changed, 601 insertions(+), 117 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 3c6d85d..d86d348 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1988,7 +1988,7 @@ checksum = "e6d5a32815ae3f33302d95fdcb2ce17862f8c65363dcfd29360480ba1001fc9c" [[package]] name = "furumusic" -version = "0.10.6" +version = "0.10.7" dependencies = [ "anyhow", "async-stream", @@ -3900,9 +3900,9 @@ dependencies = [ [[package]] name = "music-dht" -version = "0.4.1" +version = "0.5.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "7bfe5cee89fa891b00e84738f947c924c6845fd84fedb4697cb073e3fe63bb81" +checksum = "4c3f75d49d3a742a6777a1c4cc53f397399dfc9ddb70034c57aeb44d75d2ee4f" dependencies = [ "async-trait", "blake3", diff --git a/Cargo.toml b/Cargo.toml index 0658b78..dc4dbf7 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "furumusic" -version = "0.10.6" +version = "0.10.7" edition = "2024" description = "Reusable web-app boilerplate: auth, OIDC/SSO, admin panel, user management, i18n, PostgreSQL" @@ -45,4 +45,4 @@ uuid = "1" librqbit = { version = "8.1.1", features = ["disable-upload"] } # P2P federation: publishes the library into a shared DHT and serves audio / # catalogs to furumi peers (TUI clients) over the frid stack. -music-dht = "0.4.1" +music-dht = "0.5.0" diff --git a/src/federation/devices.rs b/src/federation/devices.rs index 0223c1b..2dab5a3 100644 --- a/src/federation/devices.rs +++ b/src/federation/devices.rs @@ -5,6 +5,9 @@ //! protocol as the TUI clients on `furumi/sync/1` and maps operations into //! user-scoped Postgres state. +use music_dht::playback::{ + Announcement, Checkpoint, CommandStamp, Config as PlaybackConfig, Engine, +}; use std::collections::BTreeMap; use std::str::FromStr; use std::sync::Arc; @@ -24,7 +27,7 @@ use super::{TransportStats, record_stream_transport}; pub const SYNC_ALPN: &[u8] = b"furumi/sync/2"; const CLIENT_VERSION: &str = env!("CARGO_PKG_VERSION"); -const PROTOCOL_VERSION: u16 = 2; +const PROTOCOL_VERSION: u16 = 3; const INVITE_TTL_MS: i64 = 10 * 60 * 1000; const PAIRING_WAIT_MS: i64 = 5 * 60 * 1000; const PAIRING_RETRY_DELAY: Duration = Duration::from_secs(1); @@ -193,29 +196,8 @@ struct PlaybackStateWire { repeat: PlaybackRepeat, } -#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] -struct PlaybackSnapshot { - device_id: String, - device_name: String, - active: bool, - updated_at_ms: i64, - state: PlaybackStateWire, -} - -#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] -#[serde(tag = "kind", rename_all = "snake_case")] -enum PlaybackCommand { - SetState { - state: PlaybackStateWire, - #[serde(default)] - seek: bool, - }, - ActiveChanged { - active_device_id: String, - active_device_name: String, - state: PlaybackStateWire, - }, -} +type PlaybackSnapshot = music_dht::playback::Snapshot; +type PlaybackCommand = music_dht::playback::Command; #[derive(Debug, Clone, Serialize, Deserialize)] struct SyncOpWire { @@ -273,6 +255,8 @@ enum SyncOpPayload { PlaybackCommand { target_device_id: String, command: PlaybackCommand, + #[serde(default)] + authority: Option, }, ListenRecorded { event: ListenEvent, @@ -1235,7 +1219,14 @@ async fn try_connect_invite( enforce_single_user_binding(pool, user_id, &profile.device_id).await?; apply_device_profile(pool, user_id, &profile, true).await?; if let Some(playback) = playback { - apply_playback_snapshot(Arc::clone(&hub), pool, user_id, playback).await?; + apply_playback_snapshot( + Arc::clone(&hub), + pool, + user_id, + &profile.device_id, + playback, + ) + .await?; } } apply_device_profiles(pool, user_id, &devices).await?; @@ -1454,7 +1445,14 @@ async fn handle_pair_request( let own_profile = own_profile(pool, user_id, "", &own_ticket).await?; apply_device_profile(pool, user_id, &profile, true).await?; if let Some(playback) = playback { - apply_playback_snapshot(Arc::clone(&hub), pool, user_id, playback).await?; + apply_playback_snapshot( + Arc::clone(&hub), + pool, + user_id, + &profile.device_id, + playback, + ) + .await?; } apply_snapshot(pool, user_id, incoming_snapshot).await?; apply_ops(pool, Arc::clone(&hub), user_id, ops).await?; @@ -1524,7 +1522,14 @@ async fn handle_hello( apply_device_profile(pool, user_id, &profile, false).await?; apply_device_profiles(pool, user_id, &devices).await?; if let Some(playback) = playback { - apply_playback_snapshot(Arc::clone(&hub), pool, user_id, playback).await?; + apply_playback_snapshot( + Arc::clone(&hub), + pool, + user_id, + &profile.device_id, + playback, + ) + .await?; } apply_snapshot(pool, user_id, incoming_snapshot).await?; apply_ops(pool, Arc::clone(&hub), user_id, ops).await?; @@ -1618,7 +1623,14 @@ async fn sync_device( } => { apply_device_profiles(pool, user_id, &devices).await?; if let Some(playback) = playback { - apply_playback_snapshot(Arc::clone(&hub), pool, user_id, playback).await?; + apply_playback_snapshot( + Arc::clone(&hub), + pool, + user_id, + &device.device_id, + playback, + ) + .await?; } apply_snapshot(pool, user_id, snapshot).await?; apply_ops(pool, hub, user_id, ops).await?; @@ -1952,8 +1964,33 @@ async fn own_profile( }) } -async fn record_local_op(pool: &sqlx::PgPool, user_id: i64, payload: SyncOpPayload) -> Result<()> { +async fn record_local_op( + pool: &sqlx::PgPool, + user_id: i64, + mut payload: SyncOpPayload, +) -> Result<()> { let identity = ensure_identity(pool, user_id, "").await?; + if let SyncOpPayload::PlaybackCommand { + command, authority, .. + } = &mut payload + { + *authority = Some( + with_playback_engine(pool, user_id, &identity, |engine| { + if let PlaybackCommand::ActiveChanged { + active_device_id, .. + } = command + { + if engine.owner() != Some(active_device_id.as_str()) { + engine.transfer(active_device_id, playback_clock()); + } + } + engine.stamp() + }) + .await? + .context("no playback owner; select an output first")?, + ); + } + let now = now_ms(); let row = sqlx::query( "UPDATE furumusic__fed_device_identity @@ -2193,12 +2230,8 @@ async fn apply_op( ) .await } - SyncOpPayload::PlaybackCommand { - target_device_id, - command, - } => { - apply_playback_command(pool, hub, user_id, target_device_id, command, &op.op_id) - .await?; + SyncOpPayload::PlaybackCommand { .. } => { + apply_playback_command(pool, hub, user_id, op).await?; Ok(false) } SyncOpPayload::ListenRecorded { event } => { @@ -2798,12 +2831,41 @@ async fn apply_playback_command( pool: &sqlx::PgPool, hub: Arc, user_id: i64, - target_device_id: &str, - command: &PlaybackCommand, - op_id: &str, + op: &SyncOpWire, ) -> Result<()> { + let SyncOpPayload::PlaybackCommand { + target_device_id, + command, + authority, + } = &op.payload + else { + return Ok(()); + }; + let authority = authority.as_ref(); + let origin = &op.origin_device_id; + let op_id = &op.op_id; let identity = ensure_identity(pool, user_id, "").await?; - if target_device_id != identity.device_id { + let Some(authority) = authority else { + return Ok(()); + }; + let handoff = matches!(command, PlaybackCommand::ActiveChanged { .. }); + if !handoff && target_device_id != &identity.device_id { + return Ok(()); + } + if let PlaybackCommand::ActiveChanged { + active_device_id, .. + } = command + { + if active_device_id != &authority.claim.owner { + return Ok(()); + } + } + let accepted = with_playback_engine(pool, user_id, &identity, |engine| { + engine.accept_command(origin, authority, handoff, playback_clock()) + && (handoff || engine.is_owner()) + }) + .await?; + if !accepted || target_device_id != &identity.device_id { return Ok(()); } let inserted = sqlx::query( @@ -2822,17 +2884,36 @@ async fn apply_playback_command( } match command { PlaybackCommand::SetState { state, .. } => { - enqueue_web_transfer(pool, hub, user_id, state).await?; + enqueue_web_transfer(pool, hub, user_id, state, op).await?; } PlaybackCommand::ActiveChanged { active_device_id, state, .. } if active_device_id == &identity.device_id => { - enqueue_web_transfer(pool, hub, user_id, state).await?; + enqueue_web_transfer(pool, hub, user_id, state, op).await?; } - PlaybackCommand::ActiveChanged { .. } => { - let _ = hub.enqueue_fed_command(user_id, "pause", serde_json::json!({})); + PlaybackCommand::ActiveChanged { + active_device_id, + active_device_name, + state, + } => { + let payload = web_playback_payload(pool, state).await?; + if !with_playback_engine(pool, user_id, &identity, |engine| { + engine.command_is_current(origin, authority) + }) + .await? + { + return Ok(()); + } + hub.apply_fed_playback_state_json( + user_id, + active_device_id, + active_device_name, + true, + payload, + ) + .map_err(|message| anyhow::anyhow!(message))?; } } Ok(()) @@ -2842,14 +2923,36 @@ async fn apply_playback_snapshot( hub: Arc, pool: &sqlx::PgPool, user_id: i64, + sender: &str, snapshot: PlaybackSnapshot, ) -> Result<()> { + if snapshot.device_id != sender { + return Ok(()); + } + let identity = ensure_identity(pool, user_id, "").await?; + let Some(coordination) = &snapshot.coordination else { + return Ok(()); + }; + let owner = with_playback_engine(pool, user_id, &identity, |engine| { + if !engine.observe(&snapshot.device_id, coordination, playback_clock()) { + return None; + } + engine + .announcement_is_current(&snapshot.device_id, coordination) + .then(|| snapshot.device_id.clone()) + }) + .await?; let payload = web_playback_payload(pool, &snapshot.state).await?; + let active = owner.as_deref() == Some(snapshot.device_id.as_str()) + && with_playback_engine(pool, user_id, &identity, |engine| { + engine.announcement_is_current(&snapshot.device_id, coordination) + }) + .await?; hub.apply_fed_playback_state_json( user_id, &snapshot.device_id, &snapshot.device_name, - snapshot.active, + active, payload, ) .map_err(|message| anyhow::anyhow!(message))?; @@ -2861,8 +2964,24 @@ async fn enqueue_web_transfer( hub: Arc, user_id: i64, state: &PlaybackStateWire, + op: &SyncOpWire, ) -> Result<()> { + let identity = ensure_identity(pool, user_id, "").await?; let payload = web_playback_payload(pool, state).await?; + let SyncOpPayload::PlaybackCommand { + authority: Some(authority), + .. + } = &op.payload + else { + return Ok(()); + }; + if !with_playback_engine(pool, user_id, &identity, |engine| { + engine.command_is_current(&op.origin_device_id, authority) + }) + .await? + { + return Ok(()); + } if payload .get("tracks") .and_then(serde_json::Value::as_array) @@ -2891,6 +3010,7 @@ pub async fn record_web_playback_command( pool, user_id, SyncOpPayload::PlaybackCommand { + authority: None, target_device_id: target_device_id.to_string(), command: PlaybackCommand::SetState { state: wire, @@ -2915,14 +3035,22 @@ pub async fn record_web_active_transfer( ensure_web_playback_target(pool, user_id, target_device_id).await?; let wire = playback_state_from_browser_json(pool, state).await?; let target_name = web_playback_target_name(pool, user_id, target_device_id).await?; + let identity = ensure_identity(pool, user_id, "").await?; + with_playback_engine(pool, user_id, &identity, |engine| { + engine.transfer(target_device_id, playback_clock()) + }) + .await?; + record_local_op( pool, user_id, SyncOpPayload::PlaybackCommand { + authority: None, target_device_id: target_device_id.to_string(), - command: PlaybackCommand::SetState { + command: PlaybackCommand::ActiveChanged { + active_device_id: target_device_id.to_string(), + active_device_name: target_name.clone(), state: wire.clone(), - seek: true, }, }, ) @@ -2940,6 +3068,7 @@ pub async fn record_web_active_transfer( pool, user_id, SyncOpPayload::PlaybackCommand { + authority: None, target_device_id: previous_device_id.to_string(), command, }, @@ -2958,10 +3087,15 @@ pub async fn record_web_active_takeover( ensure_web_playback_target(pool, user_id, previous_device_id).await?; let identity = ensure_identity(pool, user_id, "").await?; let wire = playback_state_from_browser_json(pool, state).await?; + with_playback_engine(pool, user_id, &identity, |engine| { + engine.transfer(&identity.device_id, playback_clock()) + }) + .await?; record_local_op( pool, user_id, SyncOpPayload::PlaybackCommand { + authority: None, target_device_id: previous_device_id.to_string(), command: PlaybackCommand::ActiveChanged { active_device_id: identity.device_id, @@ -3423,11 +3557,16 @@ async fn local_playback_snapshot( user_id: i64, identity: &Identity, ) -> Option { - // Keep publishing an inactive snapshot after a handoff. Omitting the - // snapshot left the last `active: true` value alive on trusted peers until - // its TTL elapsed, allowing the always-on web peer to reclaim playback. - let active = hub.federation_playback_is_local(user_id); - let state = hub.playback_state_json_for_commands(user_id)?; + let coordination = coordinate_web_output(pool, Arc::clone(&hub), user_id, identity) + .await + .ok()?; + let active = coordination + .claim + .as_ref() + .is_some_and(|claim| claim.owner == identity.device_id); + let state = hub + .playback_state_json_for_commands(user_id) + .unwrap_or_else(|| serde_json::json!({})); let wire = playback_state_from_browser_json(pool, state).await.ok()?; Some(PlaybackSnapshot { device_id: identity.device_id.clone(), @@ -3435,9 +3574,123 @@ async fn local_playback_snapshot( active, updated_at_ms: now_ms(), state: wire, + coordination: Some(coordination), }) } +/// Drive coordination from browser HTTP traffic even before the first peer +/// connects. This does not require a running federation transport. +pub async fn refresh_web_output( + pool: &sqlx::PgPool, + hub: Arc, + user_id: i64, +) -> Result<()> { + let identity = ensure_identity(pool, user_id, "").await?; + coordinate_web_output(pool, hub, user_id, &identity).await?; + Ok(()) +} + +async fn coordinate_web_output( + pool: &sqlx::PgPool, + hub: Arc, + user_id: i64, + identity: &Identity, +) -> Result { + let (available, playing, report, startup) = hub.federation_output_report(user_id); + let coordination = with_playback_engine(pool, user_id, identity, |engine| { + engine.set_output(available, playing); + if startup { + engine.request_startup(); + } + // A server poll cannot refresh the heartbeat of a suspended browser. + engine.output_report(report, playback_clock()); + engine.tick(playback_clock()); + engine.announcement() + }) + .await?; + if let Some(owner) = &coordination.claim { + hub.enforce_federation_owner(user_id, &identity.device_id, &owner.owner); + } + Ok(coordination) +} + +// One serialized coordinator per account. The lock covers checkpoint commit so +// a later request cannot publish a term before its predecessor is durable. +struct PlaybackSession { + engine: Option, + group_id: String, +} +type PlaybackSessions = std::sync::Mutex>>>; + +async fn with_playback_engine( + pool: &sqlx::PgPool, + user_id: i64, + identity: &Identity, + f: impl FnOnce(&mut Engine) -> R, +) -> Result { + static SESSIONS: std::sync::OnceLock = std::sync::OnceLock::new(); + let session = { + let mut sessions = SESSIONS + .get_or_init(Default::default) + .lock() + .expect("playback sessions"); + sessions + .entry(user_id) + .or_insert_with(|| { + Arc::new(tokio::sync::Mutex::new(PlaybackSession { + engine: None, + group_id: String::new(), + })) + }) + .clone() + }; + let mut session = session.lock().await; + if session.engine.is_none() || session.group_id != identity.group_id { + session.group_id = identity.group_id.clone(); + let row = sqlx::query("SELECT playback_coordination_json, playback_config_json FROM furumusic__fed_device_identity WHERE user_id = $1") + .bind(user_id).fetch_one(pool).await?; + let durable = row + .get::, _>("playback_coordination_json") + .map(serde_json::from_value::) + .transpose()? + .filter(|checkpoint| checkpoint.scope == identity.group_id) + .map(|checkpoint| checkpoint.state) + .unwrap_or_default(); + let config = row + .get::, _>("playback_config_json") + .map(serde_json::from_value::) + .transpose()? + .unwrap_or_else(PlaybackConfig::passive); + session.engine = Some(Engine::new( + identity.device_id.clone(), + config, + durable, + playback_clock(), + )); + } + let engine = session + .engine + .as_mut() + .expect("initialized playback session"); + let previous = engine.clone(); + let result = f(engine); + if previous.durable() != engine.durable() { + if let Err(error) = sqlx::query("UPDATE furumusic__fed_device_identity SET playback_coordination_json = $2 WHERE user_id = $1") + .bind(user_id).bind(serde_json::to_value(&Checkpoint { scope: identity.group_id.clone(), state: engine.durable().clone() })?).execute(pool).await { + *engine = previous; return Err(error.into()); + } + } + Ok(result) +} + +fn playback_clock() -> u64 { + static START: std::sync::OnceLock = std::sync::OnceLock::new(); + START + .get_or_init(std::time::Instant::now) + .elapsed() + .as_millis() as u64 +} + async fn playback_state_from_browser_json( pool: &sqlx::PgPool, state: serde_json::Value, diff --git a/src/federation/mod.rs b/src/federation/mod.rs index 1d2f021..cf70ddb 100644 --- a/src/federation/mod.rs +++ b/src/federation/mod.rs @@ -1123,6 +1123,11 @@ impl Federation { .await } + pub async fn fed_device_web_refresh(&self, user_id: i64) -> Result<()> { + let pool = self.pool().await?; + devices::refresh_web_output(&pool, crate::player::PlayerDeviceHub::shared(), user_id).await + } + pub async fn fed_device_web_command( &self, user_id: i64, diff --git a/src/music/mod.rs b/src/music/mod.rs index e6d1130..8f425fa 100644 --- a/src/music/mod.rs +++ b/src/music/mod.rs @@ -1984,6 +1984,8 @@ pub mod db_migrations { )", ) .await?; + ctx.db.raw("ALTER TABLE furumusic__fed_device_identity ADD COLUMN IF NOT EXISTS playback_coordination_json JSONB").await?; + ctx.db.raw("ALTER TABLE furumusic__fed_device_identity ADD COLUMN IF NOT EXISTS playback_config_json JSONB").await?; ctx.db .raw( "CREATE TABLE IF NOT EXISTS furumusic__fed_device ( diff --git a/src/player/mod.rs b/src/player/mod.rs index 523f7ea..88c9431 100644 --- a/src/player/mod.rs +++ b/src/player/mod.rs @@ -108,7 +108,7 @@ struct LocalUploadResponse { upload: LocalUploadDto, } -const PLAYER_DEVICE_TTL_MS: i64 = 30_000; +const PLAYER_DEVICE_TTL_MS: i64 = 120_000; const PLAYER_DEVICE_RETURN_TAKEOVER_MS: i64 = 30 * 60 * 1_000; const PLAYER_DEVICE_COMMAND_TTL_MS: i64 = 20_000; const PLAYER_DEVICE_MAX_COMMANDS: usize = 32; @@ -124,6 +124,7 @@ struct PlayerDevice { id: String, name: String, kind: String, + report_sequence: u64, last_seen_ms: i64, } @@ -165,6 +166,8 @@ struct PlayerDeviceHubState { commands_by_device: HashMap<(i64, String), VecDeque>, playback_state_by_user: HashMap, jams_by_id: HashMap, + playback_startup_by_user: HashMap, + output_report_sequence: u64, } #[derive(Debug, Default)] @@ -187,6 +190,9 @@ impl PlayerDeviceHub { let now = current_millis(); let mut state = self.state.lock().expect("player device hub lock"); self.prune_locked(&mut state, now); + if self.user_has_joined_jam_locked(&state, user_id) { + return Ok(()); + } let devices = state .devices_by_user .get(&user_id) @@ -205,17 +211,105 @@ impl PlayerDeviceHub { .map(|device| device.id.clone()) }) .ok_or("no browser playback device")?; - state.active_device_by_user.insert(user_id, target.clone()); + if command == "transfer_state" { + state.active_device_by_user.insert(user_id, target.clone()); + } self.enqueue_command_locked(&mut state, user_id, &target, command, payload, now); Ok(()) } - pub(crate) fn federation_playback_is_local(&self, user_id: i64) -> bool { - let state = self.state.lock().expect("player device hub lock"); - !state + pub(crate) fn federation_output_report(&self, user_id: i64) -> (bool, bool, u64, bool) { + let mut state = self.state.lock().expect("player device hub lock"); + if self.user_has_joined_jam_locked(&state, user_id) { + return (false, false, 0, false); + } + let now = current_millis(); + let active = state.active_device_by_user.get(&user_id); + let active_browser = active + .filter(|id| !is_fed_virtual_device_id(id)) + .and_then(|id| state.devices_by_user.get(&user_id)?.get(id)) + .filter(|device| now.saturating_sub(device.last_seen_ms) < PLAYER_DEVICE_TTL_MS); + let candidate = state.devices_by_user.get(&user_id).and_then(|devices| { + devices + .values() + .filter(|device| { + !is_fed_virtual_device_id(&device.id) + && now.saturating_sub(device.last_seen_ms) < PLAYER_DEVICE_TTL_MS + }) + .max_by_key(|device| (device.report_sequence, &device.id)) + }); + let available = candidate.is_some(); + // Only the selected browser can renew the gateway's owned output. + let report = active_browser.map_or(0, |device| device.report_sequence); + let playing = active_browser.is_some() + && state + .playback_state_by_user + .get(&user_id) + .is_some_and(|playback| playback.track.is_some() && !playback.paused); + let startup = state + .playback_startup_by_user + .remove(&user_id) + .is_some_and(|started| { + started.elapsed().as_millis() <= PLAYER_DEVICE_COMMAND_TTL_MS as u128 + }); + (available, playing, report, startup) + } + + pub(crate) fn enforce_federation_owner(&self, user_id: i64, local: &str, owner: &str) { + let mut state = self.state.lock().expect("player device hub lock"); + if self.user_has_joined_jam_locked(&state, user_id) { + return; + } + if owner == local { + let now = current_millis(); + let active_local = state + .active_device_by_user + .get(&user_id) + .is_some_and(|id| !is_fed_virtual_device_id(id)); + if !active_local { + let candidate = state.devices_by_user.get(&user_id).and_then(|devices| { + devices + .values() + .filter(|device| { + !is_fed_virtual_device_id(&device.id) + && now.saturating_sub(device.last_seen_ms) < PLAYER_DEVICE_TTL_MS + }) + .max_by_key(|device| (device.report_sequence, &device.id)) + .map(|device| device.id.clone()) + }); + if let Some(candidate) = candidate { + state + .active_device_by_user + .insert(user_id, candidate.clone()); + if let Some(playback) = + state + .playback_state_by_user + .get(&user_id) + .and_then(|playback| { + serde_json::to_value(playback_state_at(playback.clone(), now)).ok() + }) + { + self.enqueue_command_locked( + &mut state, + user_id, + &candidate, + "transfer_state", + playback, + now, + ); + } + } + } + return; + } + state .active_device_by_user - .get(&user_id) - .is_some_and(|id| is_fed_virtual_device_id(id)) + .insert(user_id, fed_virtual_device_id(owner)); + // Poll responses also identify the winner. Purge delayed play/transfer + // commands so reconnecting browsers cannot resume an obsolete session. + state + .commands_by_device + .retain(|(user, _), _| *user != user_id); } pub(crate) fn playback_state_json_for_commands( @@ -269,39 +363,27 @@ impl PlayerDeviceHub { state.devices_by_user.entry(user_id).or_default().insert( virtual_id.clone(), PlayerDevice { + report_sequence: 0, id: virtual_id.clone(), name: fed_device_name.to_string(), kind: "fed".to_string(), last_seen_ms: now, }, ); - // Match the trusted-device playback contract used by the TUI: a - // background/stale active snapshot must not steal playback from a - // browser that is actively playing. An explicit web handoff changes - // `active_device_by_user` to the federated virtual device before the - // snapshot arrives, so it still passes through here. - let local_playback_is_protected = state - .active_device_by_user - .get(&user_id) - .is_some_and(|active_id| !is_fed_virtual_device_id(active_id)) - && state - .playback_state_by_user - .get(&user_id) - .is_some_and(|playback| playback.track.is_some() && !playback.paused); - if active && local_playback_is_protected { + // The shared coordinator has already resolved ownership. A local + // playing flag is not permission to reject its winning claim. + if self.user_has_joined_jam_locked(&state, user_id) { return Ok(()); } - let should_update_playback = active - || state - .active_device_by_user - .get(&user_id) - .is_some_and(|active_id| active_id == &virtual_id); if active { + state + .commands_by_device + .retain(|(user, _), _| *user != user_id); state .active_device_by_user .insert(user_id, virtual_id.clone()); } - if should_update_playback { + if active { state.playback_state_by_user.insert(user_id, playback_state); } Ok(()) @@ -330,9 +412,32 @@ impl PlayerDeviceHub { .playback_state_by_user .get(&user_id) .is_some_and(|playback| playback.track.is_some() && !playback.paused); - let should_claim_idle_playback = is_new_or_returning - && previous_active_id.as_deref() != Some(device_id) - && !active_is_playing; + let mut policy = music_dht::playback::Config::default(); + // Local browser failover is permitted. The always-on gateway is not + // an automatic candidate against another federated output. + if previous_active_id + .as_deref() + .is_some_and(is_fed_virtual_device_id) + { + policy.automatic_failover = false; + } + let owner = previous_active_id.as_ref().map(|id| { + let age = state + .devices_by_user + .get(&user_id) + .and_then(|devices| devices.get(id)) + .map_or(u64::MAX, |device| { + now.saturating_sub(device.last_seen_ms).max(0) as u64 + }); + (active_is_playing, age) + }); + let should_claim_idle_playback = previous_active_id.as_deref() != Some(device_id) + && policy.should_claim(is_new_or_returning, owner); + if is_new_or_returning { + state + .playback_startup_by_user + .insert(user_id, std::time::Instant::now()); + } if should_claim_idle_playback { let transfer_state = state .playback_state_by_user @@ -379,9 +484,57 @@ impl PlayerDeviceHub { let now = current_millis(); let mut state = self.state.lock().expect("player device hub lock"); self.prune_locked(&mut state, now); + let previous = state.active_device_by_user.get(&user_id).cloned(); + if let Some(previous) = + previous.filter(|id| id != device_id && !is_fed_virtual_device_id(id)) + { + let age = state + .devices_by_user + .get(&user_id) + .and_then(|devices| devices.get(&previous)) + .map_or(u64::MAX, |device| { + now.saturating_sub(device.last_seen_ms).max(0) as u64 + }); + let playing = state + .playback_state_by_user + .get(&user_id) + .is_some_and(|playback| playback.track.is_some() && !playback.paused); + if music_dht::playback::Config::default().should_claim(false, Some((playing, age))) { + state + .active_device_by_user + .insert(user_id, device_id.to_string()); + if let Some(payload) = + state + .playback_state_by_user + .get(&user_id) + .and_then(|playback| { + serde_json::to_value(playback_state_at(playback.clone(), now)).ok() + }) + { + self.enqueue_command_locked( + &mut state, + user_id, + device_id, + "transfer_state", + payload, + now, + ); + } + } + } self.touch_locked(&mut state, user_id, device_id, user_agent, now); self.update_playback_state_locked(&mut state, user_id, device_id, playback_state, now); self.touch_jam_locked(&mut state, user_id, device_id, current_jam_id, now); + if current_jam_id.is_none() + && state + .active_device_by_user + .get(&user_id) + .is_some_and(|active| active != device_id) + { + state + .commands_by_device + .remove(&(user_id, device_id.to_string())); + } let commands = state .commands_by_device .remove(&(user_id, device_id.to_string())) @@ -534,8 +687,11 @@ impl PlayerDeviceHub { user_agent: Option<&str>, now: i64, ) { + state.output_report_sequence = state.output_report_sequence.saturating_add(1); + let report_sequence = state.output_report_sequence; let devices = state.devices_by_user.entry(user_id).or_default(); let device = PlayerDevice { + report_sequence, id: device_id.to_string(), name: device_name_from_user_agent(user_agent), kind: device_kind_from_user_agent(user_agent).to_string(), @@ -546,15 +702,7 @@ impl PlayerDeviceHub { .device_last_seen_ms .insert((user_id, device_id.to_string()), now); - let active_online = state - .active_device_by_user - .get(&user_id) - .is_some_and(|active_id| devices.contains_key(active_id)); - if !active_online { - state - .active_device_by_user - .insert(user_id, device_id.to_string()); - } + // Discovery only registers devices; startup/select decides ownership. } fn update_playback_state_locked( @@ -976,25 +1124,11 @@ impl PlayerDeviceHub { devices.retain(|_, device| { now.saturating_sub(device.last_seen_ms) <= PLAYER_DEVICE_TTL_MS }); - let active_valid = state - .active_device_by_user - .get(user_id) - .is_some_and(|active_id| devices.contains_key(active_id)); - if !active_valid { - if let Some(first_device_id) = devices.keys().next().cloned() { - state - .active_device_by_user - .insert(*user_id, first_device_id); - } else { - state.active_device_by_user.remove(user_id); - state.playback_state_by_user.remove(user_id); - } - } + // Keep ownership and queue when presence expires. The shared + // protocol decides failover; HashMap order must never choose audio. + let _ = user_id; !devices.is_empty() }); - state - .playback_state_by_user - .retain(|user_id, _| state.devices_by_user.contains_key(user_id)); state .commands_by_device @@ -1170,6 +1304,83 @@ fn device_kind_from_user_agent(user_agent: Option<&str>) -> &'static str { mod device_tests { use super::*; + #[test] + fn gateway_timer_does_not_manufacture_browser_reports() { + let hub = PlayerDeviceHub::default(); + assert_eq!(hub.federation_output_report(1), (false, false, 0, false)); + hub.heartbeat(1, "browser", None, None, None); + let first = hub.federation_output_report(1); + let repeated = hub.federation_output_report(1); + assert!(first.0); + assert!(first.3); + assert_eq!(first.2, repeated.2); + assert!(!repeated.3); + } + + #[test] + fn an_old_browser_startup_does_not_claim_a_late_peer() { + let hub = PlayerDeviceHub::default(); + hub.heartbeat(1, "browser", None, None, None); + hub.state.lock().unwrap().playback_startup_by_user.insert( + 1, + std::time::Instant::now() + - std::time::Duration::from_millis(PLAYER_DEVICE_COMMAND_TTL_MS as u64 + 1), + ); + assert!(!hub.federation_output_report(1).3); + } + + #[test] + fn expired_presence_does_not_replace_federated_owner() { + let hub = PlayerDeviceHub::default(); + hub.state + .lock() + .unwrap() + .active_device_by_user + .insert(1, "fed:remote".into()); + let response = hub.poll(1, "browser", None, None, None); + assert_eq!(response.active_device_id.as_deref(), Some("fed:remote")); + assert!(response.commands.is_empty()); + } + + #[test] + fn browser_poll_fails_over_only_after_local_owner_timeout() { + let hub = PlayerDeviceHub::default(); + hub.heartbeat(1, "first", None, None, None); + assert_eq!( + hub.poll(1, "second", None, None, None) + .active_device_id + .as_deref(), + Some("first") + ); + hub.state + .lock() + .unwrap() + .devices_by_user + .get_mut(&1) + .unwrap() + .get_mut("first") + .unwrap() + .last_seen_ms = current_millis() - PLAYER_DEVICE_TTL_MS - 1; + assert_eq!( + hub.poll(1, "second", None, None, None) + .active_device_id + .as_deref(), + Some("second") + ); + } + + #[test] + fn pruning_keeps_the_owner_when_every_device_is_offline() { + let hub = PlayerDeviceHub::default(); + hub.heartbeat(1, "browser", None, None, None); + let mut state = hub.state.lock().unwrap(); + hub.prune_locked(&mut state, current_millis() + PLAYER_DEVICE_TTL_MS + 1); + assert_eq!( + state.active_device_by_user.get(&1).map(String::as_str), + Some("browser") + ); + } + #[test] fn detects_furumi_android_native_client() { let user_agent = Some("FurumiAndroid/1.0 Android Mobile"); @@ -1217,7 +1428,7 @@ mod device_tests { } #[test] - fn federated_snapshot_does_not_steal_active_browser_playback() { + fn resolved_federation_owner_overrides_a_playing_browser() { let hub = PlayerDeviceHub::default(); let user_id = 7; { @@ -1225,6 +1436,7 @@ mod device_tests { state.devices_by_user.entry(user_id).or_default().insert( "browser".to_string(), PlayerDevice { + report_sequence: 0, id: "browser".to_string(), name: "Browser".to_string(), kind: "computer".to_string(), @@ -1276,7 +1488,7 @@ mod device_tests { .active_device_by_user .get(&user_id) .map(String::as_str), - Some("browser") + Some("fed:remote") ); assert!( state @@ -1372,6 +1584,7 @@ mod device_tests { state.devices_by_user.entry(user_id).or_default().insert( "browser".to_string(), PlayerDevice { + report_sequence: 0, id: "browser".to_string(), name: "Browser".to_string(), kind: "computer".to_string(), @@ -5347,6 +5560,12 @@ async fn devices_heartbeat_handler( return Ok(json_error(StatusCode::BAD_REQUEST, &format!("{err}"))); } } + if let Err(error) = crate::federation::handle() + .fed_device_web_refresh(user.id) + .await + { + tracing::warn!(user_id = user.id, %error, "playback coordination startup failed"); + } Json(response).into_response() } @@ -5364,6 +5583,12 @@ async fn devices_poll_handler( return Ok(json_error(StatusCode::BAD_REQUEST, "invalid device id")); }; + if let Err(error) = crate::federation::handle() + .fed_device_web_refresh(user.id) + .await + { + tracing::warn!(user_id = user.id, %error, "playback coordination refresh failed"); + } let response = hub.poll( user.id, &device_id, diff --git a/templates/player/scripts.html b/templates/player/scripts.html index e40c39d..036498a 100644 --- a/templates/player/scripts.html +++ b/templates/player/scripts.html @@ -1688,7 +1688,7 @@ document.addEventListener('alpine:init', () => { } const player = Alpine.store('player'); - if (player && Array.isArray(data.commands)) { + if (player && (this.isActive() || this.shouldPlayJamLocally()) && Array.isArray(data.commands)) { data.commands.forEach(command => player._executeRemoteCommand(command)); } if (player && !this.isActive()) { @@ -1706,7 +1706,6 @@ document.addEventListener('alpine:init', () => { }, _apply(data) { - const wasActive = this.isActive(); const previousJamId = this.currentJamId; this.activeDeviceId = data.active_device_id || null; this.devices = Array.isArray(data.devices) ? data.devices : []; @@ -1722,7 +1721,7 @@ document.addEventListener('alpine:init', () => { if (previousJamId !== this.currentJamId || !this.canPlayJamLocally()) { this._setJamLocalPlayback(false, { pauseLocal: true }); } - if (wasActive && !this.isActive()) { + if (!this.isActive() && !this.shouldPlayJamLocally()) { Alpine.store('player')?._pauseLocal(); } this._maybeShowRemoteHint();