From 52c2f7f7fe239df84187f7135f942245ed48a246 Mon Sep 17 00:00:00 2001 From: Ultradesu Date: Wed, 29 Jul 2026 14:21:01 +0100 Subject: [PATCH] Fixed invite logic to be lightweight --- Cargo.toml | 2 +- src/devices.rs | 53 +++++++++++++++++++++++++++++--------------- src/devices/tests.rs | 52 +++++++++++++++++++++++++++++++++++++++++++ 3 files changed, 88 insertions(+), 19 deletions(-) diff --git a/Cargo.toml b/Cargo.toml index e198398..f0b51d1 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "furumi_tui" -version = "0.2.3" +version = "0.2.4" edition = "2024" rust-version = "1.97" description = "A federated P2P player for personal music libraries" diff --git a/src/devices.rs b/src/devices.rs index fe86c6f..19afb8f 100644 --- a/src/devices.rs +++ b/src/devices.rs @@ -31,6 +31,9 @@ const PAIRING_RETRY_DELAY: Duration = Duration::from_secs(1); const RESPONSE_DRAIN_TIMEOUT: Duration = Duration::from_secs(2); const SYNC_INTERVAL: Duration = Duration::from_secs(2); const MAX_LINE: usize = 8 * 1024 * 1024; +/// Playback control is ephemeral. Keeping old full-queue commands in a new +/// peer's catch-up batch can make the initial pairing frame arbitrarily large. +const PLAYBACK_COMMAND_TTL_MS: i64 = 5 * 60 * 1_000; const MAX_OPS_PER_BATCH: usize = 1000; #[derive(Debug, Clone, PartialEq, Eq)] @@ -2173,25 +2176,39 @@ impl DeviceSync { LEFT JOIN sync_peer_acks a ON a.peer_device_id = ?1 AND a.origin_device_id = o.origin_device_id WHERE o.seq > COALESCE(a.max_seq, 0) - ORDER BY o.hlc_ms, o.op_id - LIMIT ?2", - )?; - let rows = stmt.query_map(params![peer_device_id, MAX_OPS_PER_BATCH as i64], |row| { - let payload_json: String = row.get(4)?; - Ok(SyncOpWire { - op_id: row.get(0)?, - origin_device_id: row.get(1)?, - seq: row.get(2)?, - hlc_ms: row.get(3)?, - payload: serde_json::from_str(&payload_json).map_err(|err| { - rusqlite::Error::FromSqlConversionFailure( - 4, - rusqlite::types::Type::Text, - Box::new(err), + AND ( + o.kind != 'playback_command' + OR ( + o.hlc_ms >= ?2 + AND json_extract(o.payload_json, '$.target_device_id') = ?1 ) - })?, - }) - })?; + ) + ORDER BY o.hlc_ms, o.op_id + LIMIT ?3", + )?; + let rows = stmt.query_map( + params![ + peer_device_id, + now_ms().saturating_sub(PLAYBACK_COMMAND_TTL_MS), + MAX_OPS_PER_BATCH as i64 + ], + |row| { + let payload_json: String = row.get(4)?; + Ok(SyncOpWire { + op_id: row.get(0)?, + origin_device_id: row.get(1)?, + seq: row.get(2)?, + hlc_ms: row.get(3)?, + payload: serde_json::from_str(&payload_json).map_err(|err| { + rusqlite::Error::FromSqlConversionFailure( + 4, + rusqlite::types::Type::Text, + Box::new(err), + ) + })?, + }) + }, + )?; Ok(rows.collect::>>()?) } diff --git a/src/devices/tests.rs b/src/devices/tests.rs index 001f25d..be77b8c 100644 --- a/src/devices/tests.rs +++ b/src/devices/tests.rs @@ -228,6 +228,58 @@ fn playback_command_is_targeted_and_deduplicated() { assert!(rx.try_recv().is_err()); } +#[test] +fn playback_commands_are_caught_up_only_by_their_target_while_fresh() { + let sync = test_sync(); + let command = PlaybackCommand::SetState { + state: PlaybackStateWire { + queue: Vec::new(), + queue_pos: 0, + playing: false, + paused: false, + idle_since_ms: None, + position_secs: 0.0, + volume: 42, + shuffle: false, + repeat: PlaybackRepeat::Off, + }, + seek: false, + }; + sync.record_playback_command("dev_target", command.clone()) + .unwrap(); + sync.record_playback_command("dev_other", command).unwrap(); + + let target_ops = sync.ops_for_peer("dev_target").unwrap(); + assert_eq!( + target_ops + .iter() + .filter(|op| matches!(op.payload, SyncOpPayload::PlaybackCommand { .. })) + .count(), + 1 + ); + assert!( + sync.ops_for_peer("dev_unknown") + .unwrap() + .iter() + .all(|op| !matches!(op.payload, SyncOpPayload::PlaybackCommand { .. })) + ); + + lock(&sync.conn) + .execute( + "UPDATE sync_ops + SET hlc_ms = ?1 + WHERE kind = 'playback_command'", + [now_ms().saturating_sub(PLAYBACK_COMMAND_TTL_MS + 1)], + ) + .unwrap(); + assert!( + sync.ops_for_peer("dev_target") + .unwrap() + .iter() + .all(|op| !matches!(op.payload, SyncOpPayload::PlaybackCommand { .. })) + ); +} + #[test] fn newer_device_trust_reactivates_revoked_device() { let sync = test_sync();