Fixed invite logic to be lightweight
This commit is contained in:
+1
-1
@@ -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"
|
||||
|
||||
+35
-18
@@ -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::<rusqlite::Result<Vec<_>>>()?)
|
||||
}
|
||||
|
||||
|
||||
@@ -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();
|
||||
|
||||
Reference in New Issue
Block a user