11 Commits
Author SHA1 Message Date
Ultradesu cced546cd5 Fixed invite logic to be lightweight 2026-07-29 14:24:22 +01:00
Ultradesu 52c2f7f7fe Fixed invite logic to be lightweight 2026-07-29 14:21:01 +01:00
Ultradesu 9b9f5d9ce5 Global is default mode 2026-07-29 10:52:57 +01:00
Ultradesu 1f9539ffd7 Bump protocols. Added proto status 2026-07-28 22:59:35 +01:00
Ultradesu d854a4d4e0 Bump furumi 0.2.2 -> 0.2.3 2026-07-28 22:45:50 +01:00
Ultradesu 6175423b0c Bump protocols. Added proto status 2026-07-28 22:24:17 +01:00
Ultradesu 93dc8f02fb Bump protocols. Added proto status 2026-07-28 22:19:01 +01:00
Ultradesu 89c78dcadc Bump protocols. Added proto status 2026-07-28 22:18:29 +01:00
Ultradesu 34153cca9d Added jam 2026-07-28 18:58:38 +01:00
Ultradesu d4c219837b Added JAM feature 2026-07-28 18:41:08 +01:00
Ultradesu 1c721b68f5 Added JAM feature 2026-07-28 18:35:59 +01:00
24 changed files with 2119 additions and 83 deletions
+14
View File
@@ -89,6 +89,20 @@ joining another user's trusted-device group.
Keeping these layers separate prevents discovery convenience from silently
becoming a synchronization trust decision.
### Federation Jam control
Jam is a third, deliberately narrow authority boundary. A host creates an
opaque `frid://j/...` runtime capability and remains the only node producing
audio. Other TUI peers use a dedicated Jam ALPN to submit the same portable
playback commands used by connected-device control and receive the host's
playback snapshot. They receive queue metadata, not audio.
Jam never exchanges trusted membership, likes, playlists, or listening
history. Volume remains local. Commands carry unique IDs and are retried until
the host acknowledges them, while inactive participants expire from the
runtime session. Regenerating the capability or restarting the host invalidates
the previous link.
## Discovery and direct communication
Furumi separates finding content from transferring it.
Generated
+12 -12
View File
@@ -434,8 +434,6 @@ dependencies = [
[[package]]
name = "block"
version = "0.1.6"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "0d8c1fef690941d3e7788d328517591fecc684c084084702d6ff1641e993699a"
[[package]]
name = "block-buffer"
@@ -1211,13 +1209,13 @@ dependencies = [
[[package]]
name = "displaydoc"
version = "0.2.6"
version = "0.2.7"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "1ac70aa55017e108007fbaf5aa0f54b021c98f92ff8af59d42eda9da96e3dd4f"
checksum = "c6232dd377dcc64799954cbd3a9bb882e9cdc1308ccd87b1c098f1fb2eaf82a8"
dependencies = [
"proc-macro2",
"quote",
"syn 2.0.119",
"syn 3.0.3",
]
[[package]]
@@ -1454,8 +1452,9 @@ dependencies = [
[[package]]
name = "federation-net"
version = "0.1.0"
source = "git+https://gt.hexor.cy/ab/frid.git?rev=8de7d1292708fa0b225e5a4a9d5ab4f0676202d3#8de7d1292708fa0b225e5a4a9d5ab4f0676202d3"
version = "0.2.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "15a8707baeccb46b5935138f9cb3df3c988c0730b807a2634d901f26b39250d6"
dependencies = [
"blake3",
"data-encoding",
@@ -1566,7 +1565,7 @@ dependencies = [
[[package]]
name = "furumi_tui"
version = "0.2.1"
version = "0.2.4"
dependencies = [
"anyhow",
"blake3",
@@ -2988,8 +2987,9 @@ dependencies = [
[[package]]
name = "music-dht"
version = "0.2.0"
source = "git+https://gt.hexor.cy/ab/frid.git?rev=8de7d1292708fa0b225e5a4a9d5ab4f0676202d3#8de7d1292708fa0b225e5a4a9d5ab4f0676202d3"
version = "0.3.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "5f32daa9edf769fb5686e92ae6884e9fda6ea082452e4208309fc14301e26aef"
dependencies = [
"async-trait",
"blake3",
@@ -5578,9 +5578,9 @@ dependencies = [
[[package]]
name = "toml"
version = "1.1.3+spec-1.1.0"
version = "1.1.4+spec-1.1.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "53c96ecdfa941c8fc4fcaed14f99ada8ebed502eef533015095a07e3301d4c3c"
checksum = "3aace63f4bbcdfc2c965b059de67119c89c4017a70d633be6c104910f67056f5"
dependencies = [
"indexmap",
"serde_core",
+7 -2
View File
@@ -1,6 +1,6 @@
[package]
name = "furumi_tui"
version = "0.2.1"
version = "0.2.4"
edition = "2024"
rust-version = "1.97"
description = "A federated P2P player for personal music libraries"
@@ -21,7 +21,7 @@ image = { version = "0.25.10", default-features = false, features = ["jpeg", "pn
lofty = "0.22"
# P2P federation: library index in a shared DHT + audio streaming between
# peers (same protocol as furumi-fd).
music-dht = { git = "https://gt.hexor.cy/ab/frid.git", rev = "8de7d1292708fa0b225e5a4a9d5ab4f0676202d3" }
music-dht = "0.3"
ratatui = "0.30.1"
rhai = { version = "1", features = ["sync"] }
rodio = { version = "0.22.2", default-features = false, features = ["playback", "mp3", "flac", "vorbis", "wav", "symphonia-aac", "symphonia-isomp4", "symphonia-alac"] }
@@ -36,6 +36,11 @@ tracing = "0.1.44"
tracing-subscriber = { version = "0.3.23", features = ["env-filter"] }
unicode-width = "0.2.2"
[patch.crates-io]
# souvlaki 0.8.3 still depends on the unmaintained block 0.1.6. Keep its API
# intact while using an opaque inhabited FFI type accepted by current Rust.
block = { path = "vendor/block" }
[target.'cfg(target_os="macos")'.dependencies]
core-foundation = "0.10.1"
+10
View File
@@ -166,4 +166,14 @@ pub enum AppEvent {
DevicePlayback(crate::devices::PlaybackSnapshot),
/// Playback command addressed to this device.
PlaybackCommand(crate::devices::PlaybackCommand),
/// Current lifecycle/status of the federation Jam.
JamStatus(crate::jam::JamStatus),
/// Host playback snapshot received by a Jam participant.
JamPlayback(crate::devices::PlaybackSnapshot),
/// Participant command accepted by this Jam host.
JamCommand(crate::devices::PlaybackCommand),
/// Result of creating/regenerating a Jam capability.
JamInvite(Result<String, String>),
/// Result of joining a Jam capability.
JamJoined(Result<String, String>),
}
+177 -13
View File
@@ -38,6 +38,7 @@ pub struct Runtime {
pub event_tx: mpsc::UnboundedSender<AppEvent>,
pub library: Arc<Library>,
pub devices: Arc<crate::devices::DeviceSync>,
pub jam: Arc<crate::jam::JamManager>,
pub federation: Arc<crate::federation::Federation>,
/// When the last Federation-tab status snapshot was requested.
pub fed_status_at: Option<std::time::Instant>,
@@ -249,7 +250,12 @@ pub async fn run(
let devices = crate::devices::DeviceSync::new(Arc::clone(&library))?;
devices.set_event_tx(event_tx.clone());
let federation = crate::federation::Federation::new(Arc::clone(&library), Arc::clone(&devices));
let jam = crate::jam::JamManager::new(event_tx.clone());
let federation = crate::federation::Federation::new(
Arc::clone(&library),
Arc::clone(&devices),
Arc::clone(&jam),
);
state.federation.settings = federation.settings();
state.federation.devices = Some(devices.status());
if let Ok((device_id, device_name)) = devices.identity_summary() {
@@ -257,19 +263,21 @@ pub async fn run(
state.device_playback.self_device_name = device_name.clone();
state.device_playback.active_device_id = Some(device_id);
state.device_playback.active_device_name = Some(device_name);
state.device_playback.startup_takeover_pending = true;
}
let player_events = event_tx.clone();
let mut runtime = Runtime {
event_tx,
library,
devices,
jam,
federation,
fed_status_at: None,
library_network_refresh_at: None,
library_network_refreshing: Arc::new(std::sync::atomic::AtomicBool::new(false)),
library_network_cursors: Arc::new(std::sync::Mutex::new(std::collections::HashMap::new())),
library_network_done: Arc::new(std::sync::Mutex::new(std::collections::HashSet::new())),
library_network_mode: crate::config::settings::LibrarySourceMode::Local,
library_network_mode: state.global.filters.source_mode,
library_network_art_fetching: Arc::new(std::sync::atomic::AtomicBool::new(false)),
library_network_art_attempted: Arc::new(std::sync::Mutex::new(
std::collections::HashSet::new(),
@@ -360,7 +368,7 @@ fn sync_player_shared(state: &mut AppState, runtime: &Runtime) {
} else {
runtime.player.shared.audio_analysis()
};
if state.device_playback.role == state::DevicePlaybackRole::Active {
if state.device_playback.is_audio_owner() {
publish_playback_snapshot(state, runtime);
}
}
@@ -474,7 +482,7 @@ fn apply_playback_state_to_ui(
}
fn publish_playback_snapshot(state: &mut AppState, runtime: &Runtime) {
if state.device_playback.role != state::DevicePlaybackRole::Active {
if !state.device_playback.is_audio_owner() {
return;
}
publish_playback_snapshot_with_active(state, runtime, true);
@@ -503,6 +511,19 @@ fn publish_playback_snapshot_with_active(state: &mut AppState, runtime: &Runtime
state: playback_state_from_ui(state),
};
runtime.devices.publish_playback(snapshot);
if state.device_playback.role == state::DevicePlaybackRole::Jam
&& state.device_playback.jam_host
{
runtime
.jam
.publish_host_playback(crate::devices::PlaybackSnapshot {
device_id: state.device_playback.self_device_id.clone(),
device_name: state.device_playback.self_device_name.clone(),
active: true,
updated_at_ms: unix_time_ms(),
state: playback_state_from_ui(state),
});
}
}
fn update_local_idle_since(state: &mut AppState) {
@@ -532,7 +553,7 @@ fn active_idle_lease_expired(snapshot: &crate::devices::PlaybackSnapshot, now: i
}
fn local_active_lease_protected(state: &mut AppState, now: i64) -> bool {
if state.device_playback.role != state::DevicePlaybackRole::Active || !state.player.playing {
if !state.device_playback.is_audio_owner() || !state.player.playing {
return false;
}
if !state.player.paused {
@@ -576,7 +597,7 @@ pub(crate) fn become_control_device(
runtime: &Runtime,
snapshot: crate::devices::PlaybackSnapshot,
) {
if state.device_playback.role == state::DevicePlaybackRole::Active {
if state.device_playback.is_audio_owner() {
runtime.player.stop();
publish_inactive_playback_snapshot(state, runtime);
}
@@ -596,6 +617,7 @@ pub(crate) fn become_control_device(
pub(crate) fn become_active_device(state: &mut AppState, runtime: &mut Runtime, start_audio: bool) {
let was_control = state.device_playback.role == state::DevicePlaybackRole::Control;
state.device_playback.role = state::DevicePlaybackRole::Active;
state.device_playback.jam_host = false;
let Ok((device_id, device_name)) = runtime.devices.identity_summary() else {
return;
};
@@ -617,7 +639,7 @@ pub(crate) fn become_active_device(state: &mut AppState, runtime: &mut Runtime,
}
pub(crate) fn transfer_active_to_this_device(state: &mut AppState, runtime: &mut Runtime) {
if state.device_playback.role == state::DevicePlaybackRole::Active {
if state.device_playback.is_audio_owner() {
publish_playback_snapshot(state, runtime);
request_urgent_device_sync(runtime);
return;
@@ -722,6 +744,12 @@ fn record_control_playback_state(state: &mut AppState, runtime: &Runtime, seek:
state: playback_state_from_ui(state),
seek,
};
if state.device_playback.role == state::DevicePlaybackRole::Jam {
if let Err(err) = runtime.jam.submit_command(command) {
state.status_message = Some(format!("Jam command failed: {err:#}"));
}
return;
}
record_playback_command_async(runtime, target, command, "device command");
}
@@ -1479,10 +1507,41 @@ fn perform_effect(state: &mut AppState, runtime: &mut Runtime, effect: Effect) {
| Effect::DeviceSetName(_)
| Effect::DeviceRevoke(_)
| Effect::DeviceLeaveGroup
| Effect::JamCreate
| Effect::JamJoin(_)
if !state.connected_devices_enabled() =>
{
state.status_message = Some("enable federation before using connected devices".into());
}
Effect::JamCreate => {
let federation = Arc::clone(&runtime.federation);
let tx = runtime.event_tx.clone();
tokio::spawn(async move {
let result = federation
.create_jam()
.await
.map_err(|err| format!("{err:#}"));
let _ = tx.send(AppEvent::JamInvite(result));
});
}
Effect::JamJoin(invite) => {
let result = runtime
.federation
.join_jam(&invite)
.map(|()| "joined Jam; waiting for host state".to_string())
.map_err(|err| format!("{err:#}"));
let _ = runtime.event_tx.send(AppEvent::JamJoined(result));
let _ = runtime
.event_tx
.send(AppEvent::JamStatus(runtime.jam.status()));
}
Effect::JamLeave => {
runtime.jam.leave();
state.jam = runtime.jam.status();
state.device_playback.role = state::DevicePlaybackRole::Active;
state.device_playback.jam_host = false;
state.status_message = Some("left Jam".into());
}
Effect::DeviceShowInvite => {
let fed = Arc::clone(&runtime.federation);
let devices = Arc::clone(&runtime.devices);
@@ -1679,6 +1738,8 @@ fn is_controlled_playback_effect(effect: &Effect) -> bool {
fn perform_control_playback_effect(state: &mut AppState, runtime: &mut Runtime, effect: Effect) {
let mut seek = false;
let local_only_volume = state.device_playback.role == state::DevicePlaybackRole::Jam
&& matches!(effect, Effect::SetVolume(_));
match effect {
Effect::PlayCurrent => {
state.player.current = state.player.queue.get(state.player.queue_pos).cloned();
@@ -1708,6 +1769,9 @@ fn perform_control_playback_effect(state: &mut AppState, runtime: &mut Runtime,
| Effect::LoadListenHistory => {}
_ => {}
}
if local_only_volume {
return;
}
record_control_playback_state(state, runtime, seek);
}
@@ -2721,6 +2785,12 @@ fn handle_device_playback_snapshot(
runtime: &mut Runtime,
snapshot: crate::devices::PlaybackSnapshot,
) {
// Personal-device reconciliation must never change Jam ownership. Jam
// has its own authority and lifecycle even when the same TUI also belongs
// to a trusted-device group.
if state.device_playback.role == state::DevicePlaybackRole::Jam {
return;
}
if snapshot.device_id == state.device_playback.self_device_id {
return;
}
@@ -2742,6 +2812,18 @@ fn handle_device_playback_snapshot(
if !snapshot.active {
return;
}
// Starting a player is an explicit claim of the active role. Import the
// current queue/position from the previously active peer, then announce a
// normal handoff so that the old owner becomes a control device. This is
// intentionally one-shot: subsequent snapshots use the regular lease and
// explicit-transfer rules.
if state.device_playback.startup_takeover_pending {
state.device_playback.startup_takeover_pending = false;
become_control_device(state, runtime, snapshot);
transfer_active_to_this_device(state, runtime);
state.status_message = Some("playback moved to this newly started player".to_string());
return;
}
if local_active_lease_protected(state, now) {
tracing::debug!(
remote = %snapshot.device_id,
@@ -2750,10 +2832,10 @@ fn handle_device_playback_snapshot(
return;
}
let lease_expired = active_idle_lease_expired(&snapshot, now);
let already_controls_this_device = state.device_playback.is_control()
let already_controls_this_device = state.device_playback.is_personal_control()
&& state.device_playback.active_device_id.as_deref() == Some(snapshot.device_id.as_str());
if !lease_expired || already_controls_this_device {
let was_active = state.device_playback.role == state::DevicePlaybackRole::Active;
let was_active = state.device_playback.is_audio_owner();
let was_paused = state.player.playing && state.player.paused;
become_control_device(state, runtime, snapshot.clone());
if was_active && was_paused {
@@ -2765,7 +2847,7 @@ fn handle_device_playback_snapshot(
return;
}
if state.device_playback.is_control() {
if state.device_playback.is_personal_control() {
return;
}
become_active_device(state, runtime, false);
@@ -2829,7 +2911,7 @@ fn handle_playback_command(
if active_device_id == state.device_playback.self_device_id {
return;
}
let was_active = state.device_playback.role == state::DevicePlaybackRole::Active;
let was_active = state.device_playback.is_audio_owner();
let snapshot = crate::devices::PlaybackSnapshot {
device_id: active_device_id,
device_name: active_device_name,
@@ -2879,7 +2961,7 @@ fn handle_app_event(state: &mut AppState, runtime: &mut Runtime, event: AppEvent
})
.unwrap_or(1)
.max(1);
let active_revoked = state.device_playback.is_control()
let active_revoked = state.device_playback.is_personal_control()
&& state
.device_playback
.active_device_id
@@ -2897,7 +2979,7 @@ fn handle_app_event(state: &mut AppState, runtime: &mut Runtime, event: AppEvent
})
.is_some_and(|device| device.revoked)
});
let active_missing = state.device_playback.is_control()
let active_missing = state.device_playback.is_personal_control()
&& state
.device_playback
.active_device_id
@@ -2959,9 +3041,91 @@ fn handle_app_event(state: &mut AppState, runtime: &mut Runtime, event: AppEvent
AppEvent::DevicePlayback(snapshot) => {
handle_device_playback_snapshot(state, runtime, snapshot);
}
AppEvent::PlaybackCommand(_)
if state.device_playback.role == state::DevicePlaybackRole::Jam =>
{
tracing::debug!("ignored personal-device playback command while Jam is active");
}
AppEvent::PlaybackCommand(command) => {
handle_playback_command(state, runtime, command);
}
AppEvent::JamStatus(status) => {
state.jam = status.clone();
match status.role {
crate::jam::JamRole::Host => {
state.device_playback.role = state::DevicePlaybackRole::Jam;
state.device_playback.jam_host = true;
publish_playback_snapshot(state, runtime);
}
crate::jam::JamRole::Participant => {
if state.device_playback.is_audio_owner() {
runtime.player.stop();
runtime.player_start_pending = false;
}
state.device_playback.role = state::DevicePlaybackRole::Jam;
state.device_playback.jam_host = false;
}
crate::jam::JamRole::None => {
if state.device_playback.role == state::DevicePlaybackRole::Jam {
state.device_playback.role = state::DevicePlaybackRole::Active;
state.device_playback.jam_host = false;
}
}
}
}
AppEvent::JamPlayback(snapshot) => {
let local_volume = state.player.volume;
become_control_device(state, runtime, snapshot);
state.device_playback.role = state::DevicePlaybackRole::Jam;
state.device_playback.jam_host = false;
state.player.volume = local_volume;
state.status_message = Some(format!(
"Jam · controlling {}",
state.device_playback.active_label()
));
}
AppEvent::JamCommand(command) => {
let command = match command {
crate::devices::PlaybackCommand::SetState {
state: mut wire,
seek,
} => {
wire.volume = state.player.volume;
crate::devices::PlaybackCommand::SetState { state: wire, seek }
}
crate::devices::PlaybackCommand::ActiveChanged { .. } => {
state.status_message =
Some("Jam cannot transfer audio away from its host".into());
return;
}
};
handle_playback_command(state, runtime, command);
state.device_playback.role = state::DevicePlaybackRole::Jam;
state.device_playback.jam_host = true;
publish_playback_snapshot(state, runtime);
}
AppEvent::JamInvite(result) => match result {
Ok(invite) => {
if !state.device_playback.is_audio_owner() {
transfer_active_to_this_device(state, runtime);
}
state.jam = runtime.jam.status();
state.device_playback.role = state::DevicePlaybackRole::Jam;
state.device_playback.jam_host = true;
publish_playback_snapshot(state, runtime);
state.popup = Some(state::Popup::FedCopyText {
title: "Jam invite".to_string(),
text: invite,
help: "Copied capability lets federation peers control this host player until restart or regeneration.".to_string(),
});
state.status_message = Some("Jam started".into());
}
Err(error) => state.status_message = Some(format!("Jam: {error}")),
},
AppEvent::JamJoined(result) => match result {
Ok(message) => state.status_message = Some(message),
Err(error) => state.status_message = Some(format!("Jam: {error}")),
},
AppEvent::FedSearchLoaded { seq, result } => {
if runtime.search_seq.load(std::sync::atomic::Ordering::SeqCst) != seq {
return;
+44 -1
View File
@@ -82,7 +82,7 @@ pub(crate) fn connected_device_rows(state: &AppState) -> Vec<ConnectedDevicePopu
is_self: true,
online: true,
section: DevicePresenceSection::Online,
active: state.device_playback.role == crate::app::state::DevicePlaybackRole::Active,
active: state.device_playback.is_audio_owner(),
playing: state.player.playing,
paused: state.player.paused,
queue_len: state.player.queue.len(),
@@ -386,6 +386,32 @@ fn handle_connected_devices(
let cursor = cursor.min(last);
match key.code {
KeyCode::Esc | KeyCode::Char('q') => {}
KeyCode::Char('h') => {
super::perform_effect(state, runtime, crate::app::update::Effect::JamCreate);
}
KeyCode::Char('J') => {
state.popup = Some(Popup::FedInput {
field: FedInputField::JamInvite,
input: crate::app::input::LineEdit::default(),
});
}
KeyCode::Char('l') if state.jam.role != crate::jam::JamRole::None => {
super::perform_effect(state, runtime, crate::app::update::Effect::JamLeave);
}
KeyCode::Char('c') => {
if let Some(invite) = state.jam.invite.as_deref() {
match copy_to_clipboard(invite) {
Ok(()) => state.status_message = Some("Jam invite copied".into()),
Err(error) => {
state.status_message = Some(format!("clipboard: {error}"));
state.popup = Some(Popup::ConnectedDevices { cursor });
}
}
} else {
state.status_message = Some("only the Jam host has an invite".into());
state.popup = Some(Popup::ConnectedDevices { cursor });
}
}
KeyCode::Up | KeyCode::Char('k') => {
state.popup = Some(Popup::ConnectedDevices {
cursor: cursor.saturating_sub(1),
@@ -397,6 +423,12 @@ fn handle_connected_devices(
});
}
KeyCode::Enter => {
if state.device_playback.role == crate::app::state::DevicePlaybackRole::Jam {
state.status_message =
Some("leave the current Jam before switching personal devices".into());
state.popup = Some(Popup::ConnectedDevices { cursor });
return;
}
if cursor == 0 {
super::transfer_active_to_this_device(state, runtime);
state.status_message = Some("active playback moved to this device".into());
@@ -499,6 +531,17 @@ fn handle_fed_input(
);
}
}
FedInputField::JamInvite => {
if value.is_empty() {
state.status_message = Some("Jam invite is empty".into());
} else {
super::perform_effect(
state,
runtime,
crate::app::update::Effect::JamJoin(value),
);
}
}
}
}
_ => {
+22
View File
@@ -805,6 +805,7 @@ pub enum FedInputField {
ConnectTicket,
DeviceName,
ConnectInvite,
JamInvite,
}
impl FedInputField {
@@ -814,6 +815,7 @@ impl FedInputField {
FedInputField::ConnectTicket => "Connect to peer (paste ticket)",
FedInputField::DeviceName => "Device name",
FedInputField::ConnectInvite => "Connect device (paste frid://i invite)",
FedInputField::JamInvite => "Join Jam (paste frid://j invite)",
}
}
@@ -831,6 +833,9 @@ impl FedInputField {
FedInputField::ConnectInvite => {
"Paste a frid:// invite generated by another client to add this device to its sync group."
}
FedInputField::JamInvite => {
"Paste a frid://j capability to control the host player for this Jam only."
}
}
}
}
@@ -1210,6 +1215,7 @@ pub enum DevicePlaybackRole {
#[default]
Active,
Control,
Jam,
}
impl DevicePlaybackRole {
@@ -1217,6 +1223,7 @@ impl DevicePlaybackRole {
match self {
DevicePlaybackRole::Active => "active",
DevicePlaybackRole::Control => "control",
DevicePlaybackRole::Jam => "jam",
}
}
}
@@ -1232,11 +1239,25 @@ pub struct DevicePlaybackState {
pub local_idle_since_ms: Option<i64>,
pub remote: BTreeMap<String, crate::devices::PlaybackSnapshot>,
pub last_remote_snapshot: Option<crate::devices::PlaybackSnapshot>,
pub jam_host: bool,
/// A freshly started TUI owns playback by protocol. The first active
/// snapshot discovered during startup is imported and handed off here.
pub startup_takeover_pending: bool,
}
impl DevicePlaybackState {
pub fn is_control(&self) -> bool {
self.role == DevicePlaybackRole::Control
|| (self.role == DevicePlaybackRole::Jam && !self.jam_host)
}
pub fn is_personal_control(&self) -> bool {
self.role == DevicePlaybackRole::Control
}
pub fn is_audio_owner(&self) -> bool {
self.role == DevicePlaybackRole::Active
|| (self.role == DevicePlaybackRole::Jam && self.jam_host)
}
pub fn active_label(&self) -> String {
@@ -1264,6 +1285,7 @@ pub struct AppState {
pub settings_cursor: usize,
pub player: PlayerBar,
pub device_playback: DevicePlaybackState,
pub jam: crate::jam::JamStatus,
pub visualizer: crate::visualizer::VisualizerState,
pub global: GlobalTab,
pub artist_views: HashMap<i64, Loadable<ArtistDetail>>,
+6
View File
@@ -68,6 +68,12 @@ pub enum Effect {
DeviceRevoke(String),
/// Leave the current personal-device group after publishing self-revoke.
DeviceLeaveGroup,
/// Create or regenerate the runtime Jam capability.
JamCreate,
/// Join a federation Jam by capability.
JamJoin(String),
/// Leave the current Jam without changing personal-device state.
JamLeave,
/// Assemble the federated artist card (fan-out to the owning peers).
FedOpenArtist(String),
/// Download federated tracks into the local library, one by one.
+4 -4
View File
@@ -158,25 +158,25 @@ fn source_mode_cycles_on_library_playlists_and_queue_tabs() {
update(&mut state, Action::CycleSourceMode),
Some(Effect::SourceModeChanged)
);
assert_eq!(state.global.filters.source_mode, LibrarySourceMode::My);
assert_eq!(state.global.filters.source_mode, LibrarySourceMode::Local);
state.active_tab = Tab::Playlists;
assert_eq!(
update(&mut state, Action::CycleSourceMode),
Some(Effect::SourceModeChanged)
);
assert_eq!(state.global.filters.source_mode, LibrarySourceMode::Global);
assert_eq!(state.global.filters.source_mode, LibrarySourceMode::My);
state.active_tab = Tab::Queue;
assert_eq!(
update(&mut state, Action::CycleSourceMode),
Some(Effect::SourceModeChanged)
);
assert_eq!(state.global.filters.source_mode, LibrarySourceMode::Local);
assert_eq!(state.global.filters.source_mode, LibrarySourceMode::Global);
state.active_tab = Tab::Federation;
assert_eq!(update(&mut state, Action::CycleSourceMode), None);
assert_eq!(state.global.filters.source_mode, LibrarySourceMode::Local);
assert_eq!(state.global.filters.source_mode, LibrarySourceMode::Global);
}
#[test]
+2 -2
View File
@@ -4,9 +4,9 @@ use serde::{Deserialize, Serialize};
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum LibrarySourceMode {
#[default]
Local,
My,
#[default]
Global,
}
@@ -128,6 +128,6 @@ hide_featured_only = true
assert_eq!(settings.volume, 100);
assert!(settings.library.hide_featured_only);
assert_eq!(settings.library.source_mode, LibrarySourceMode::Local);
assert_eq!(settings.library.source_mode, LibrarySourceMode::Global);
}
}
+60 -23
View File
@@ -24,13 +24,16 @@ use crate::library::models::{ArtistRef, TrackItem};
pub const SYNC_ALPN: &[u8] = b"furumi/sync/2";
const CLIENT_VERSION: &str = env!("CARGO_PKG_VERSION");
const PROTOCOL_VERSION: u16 = 2;
pub const PROTOCOL_VERSION: u16 = 2;
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);
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)]
@@ -1249,8 +1252,7 @@ impl DeviceSync {
device: &StoredDevice,
transport_stats: Arc<crate::federation::TransportStats>,
) -> Result<()> {
let ticket: PeerTicket = device.endpoint_ticket.parse()?;
let peer = service.connect(ticket).await?;
let peer = resolve_device_peer(&service, device).await?;
let own_ticket = service.ticket().await?.to_string();
let identity = self.ensure_identity()?;
let profile = self.own_profile(&own_ticket)?;
@@ -2174,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<_>>>()?)
}
@@ -2380,7 +2396,7 @@ impl DeviceSync {
let own = self.ensure_identity()?.device_id;
let conn = lock(&self.conn);
let mut stmt = conn.prepare(
"SELECT device_id, endpoint_ticket
"SELECT device_id, endpoint_id, endpoint_ticket
FROM sync_devices
WHERE trusted_at_ms IS NOT NULL
AND revoked_at_ms IS NULL
@@ -2390,7 +2406,8 @@ impl DeviceSync {
let rows = stmt.query_map([own], |row| {
Ok(StoredDevice {
device_id: row.get(0)?,
endpoint_ticket: row.get(1)?,
endpoint_id: row.get(1)?,
endpoint_ticket: row.get(2)?,
})
})?;
Ok(rows.collect::<rusqlite::Result<Vec<_>>>()?)
@@ -2642,9 +2659,29 @@ impl DeviceSync {
#[derive(Debug)]
struct StoredDevice {
device_id: String,
endpoint_id: String,
endpoint_ticket: String,
}
async fn resolve_device_peer(
service: &MusicDhtService,
device: &StoredDevice,
) -> Result<music_dht::EndpointId> {
if let Ok(peer) = device.endpoint_id.parse::<music_dht::EndpointId>()
&& (service.connected_peers().contains(&peer)
|| service
.known_peers()
.iter()
.any(|contact| contact.peer_id == peer))
{
// The live DHT contact carries a ticket for the current schema. The
// persisted device ticket may predate a schema upgrade.
return Ok(peer);
}
let ticket: PeerTicket = device.endpoint_ticket.parse()?;
service.connect(ticket).await.map_err(Into::into)
}
pub async fn serve_peers(
mut acceptor: StreamAcceptor,
sync: Arc<DeviceSync>,
+52
View File
@@ -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();
+2
View File
@@ -18,6 +18,8 @@ use crate::library::Library;
/// ALPN of the audio streaming protocol (shared with furumi-fd).
pub const AUDIO_ALPN: &[u8] = b"furumi-fd/audio/1";
/// Version of the audio transfer stream protocol.
pub const AUDIO_PROTOCOL_VERSION: u16 = 1;
/// Maximum size of a JSON protocol line (request or response header).
const MAX_PROTOCOL_LINE: usize = 4096;
+190
View File
@@ -0,0 +1,190 @@
//! Informational publication and observation of protocol versions.
use std::collections::BTreeMap;
use std::sync::{Arc, Mutex, MutexGuard};
use std::time::Duration;
use anyhow::{Context, Result};
pub use music_dht::capabilities::CAPABILITIES_ALPN;
use music_dht::capabilities::{
CAPABILITIES_PROTOCOL_VERSION, CapabilityManifest, CapabilityMessage, read_message,
write_message,
};
use music_dht::{ByteStream, EndpointId, MusicDhtService, StreamAcceptor};
const PROBE_INTERVAL: Duration = Duration::from_secs(30);
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct ProtocolVersions {
pub local: BTreeMap<String, u16>,
pub observed: BTreeMap<String, u16>,
pub observed_peers: usize,
pub newer: Vec<NewerProtocol>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct NewerProtocol {
pub id: String,
pub local: u16,
pub observed: u16,
}
impl ProtocolVersions {
pub fn snapshot(observed: &ObservedVersions) -> Self {
let local = local_manifest().protocols;
let observed_versions = lock(&observed.versions).clone();
let newer = observed_versions
.iter()
.filter_map(|(id, remote)| {
let local_version = local.get(id)?;
(*remote > *local_version).then(|| NewerProtocol {
id: id.clone(),
local: *local_version,
observed: *remote,
})
})
.collect();
Self {
local,
observed: observed_versions,
observed_peers: lock(&observed.peers).len(),
newer,
}
}
}
#[derive(Default)]
pub struct ObservedVersions {
versions: Mutex<BTreeMap<String, u16>>,
peers: Mutex<BTreeMap<String, String>>,
}
fn local_manifest() -> CapabilityManifest {
CapabilityManifest::frid("furumi", env!("CARGO_PKG_VERSION"))
.with_protocol("audio", super::audio::AUDIO_PROTOCOL_VERSION)
}
pub async fn serve(mut acceptor: StreamAcceptor) {
while let Some(stream) = acceptor.accept().await {
tokio::spawn(async move {
if let Err(error) = serve_one(stream).await {
tracing::debug!("capability stream failed: {error:#}");
}
});
}
}
async fn serve_one(mut stream: ByteStream) -> Result<()> {
let response = match read_message(&mut stream).await? {
CapabilityMessage::Get {
version: CAPABILITIES_PROTOCOL_VERSION,
} => CapabilityMessage::Manifest {
manifest: local_manifest(),
},
CapabilityMessage::Get { version } => CapabilityMessage::Error {
message: format!("unsupported capability protocol {version}"),
},
_ => CapabilityMessage::Error {
message: "expected capability request".to_string(),
},
};
write_message(&mut stream, &response).await?;
stream.send.finish()?;
let _ = tokio::time::timeout(Duration::from_secs(2), stream.send.stopped()).await;
Ok(())
}
pub async fn probe_loop(service: Arc<MusicDhtService>, observed: Arc<ObservedVersions>) {
let mut interval = tokio::time::interval(PROBE_INTERVAL);
loop {
interval.tick().await;
let peers = service
.connected_peers()
.into_iter()
.chain(
service
.known_peers()
.into_iter()
.map(|contact| contact.peer_id),
)
.collect::<std::collections::BTreeSet<_>>();
for peer in peers {
if let Err(error) = probe_peer(&service, peer, &observed).await {
tracing::trace!(%peer, "peer capability probe unavailable: {error:#}");
}
}
}
}
async fn probe_peer(
service: &MusicDhtService,
peer: EndpointId,
observed: &ObservedVersions,
) -> Result<()> {
let mut stream = service.open_stream(peer, CAPABILITIES_ALPN).await?;
write_message(
&mut stream,
&CapabilityMessage::Get {
version: CAPABILITIES_PROTOCOL_VERSION,
},
)
.await?;
stream.send.finish()?;
let response = tokio::time::timeout(Duration::from_secs(5), read_message(&mut stream))
.await
.context("capability request timed out")??;
let CapabilityMessage::Manifest { manifest } = response else {
anyhow::bail!("peer did not return a capability manifest");
};
manifest.validate()?;
{
let mut versions = lock(&observed.versions);
for (id, version) in manifest.protocols {
versions
.entry(id)
.and_modify(|current| *current = (*current).max(version))
.or_insert(version);
}
}
lock(&observed.peers).insert(peer.to_string(), manifest.application_version);
Ok(())
}
fn lock<T>(mutex: &Mutex<T>) -> MutexGuard<'_, T> {
mutex
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn local_manifest_lists_every_player_protocol() {
let manifest = local_manifest();
for id in [
"federation_net",
"ticket",
"rendezvous",
"music_dht",
"catalog",
"audio",
"device_sync",
"jam",
] {
assert!(manifest.protocols.contains_key(id), "missing {id}");
}
manifest.validate().unwrap();
}
#[test]
fn snapshot_reports_only_strictly_newer_versions() {
let observed = ObservedVersions::default();
lock(&observed.versions).insert("music_dht".to_string(), 99);
lock(&observed.versions).insert("jam".to_string(), crate::jam::PROTOCOL_VERSION);
let snapshot = ProtocolVersions::snapshot(&observed);
assert_eq!(snapshot.newer.len(), 1);
assert_eq!(snapshot.newer[0].id, "music_dht");
}
}
+59 -1
View File
@@ -13,6 +13,7 @@
//! the network too).
mod audio;
mod capabilities;
pub mod catalog;
use std::collections::{HashMap, VecDeque};
@@ -36,6 +37,7 @@ use crate::library::NetworkArtistPreview;
use crate::library::models::{ArtistRef, TrackItem};
pub use audio::{AUDIO_ALPN, DownloadProgress, StreamingStart, TrackMetadata};
pub use capabilities::ProtocolVersions;
pub use catalog::{CATALOG_ALPN, FedAppearsOn, FedArtistCard, FedCardTrack, FedRelease};
/// How often the published library is re-synchronized with the local index.
@@ -386,6 +388,7 @@ pub struct FedStatus {
pub last_sync: Option<String>,
pub last_error: Option<String>,
pub transport: TransportStatsSnapshot,
pub protocols: ProtocolVersions,
}
/// Outcome of preparing a federated track for playback.
@@ -416,6 +419,7 @@ struct Running {
pub struct Federation {
library: Arc<Library>,
devices: Arc<crate::devices::DeviceSync>,
jam: Arc<crate::jam::JamManager>,
data_dir: PathBuf,
cache_dir: PathBuf,
media_dir: PathBuf,
@@ -425,6 +429,7 @@ pub struct Federation {
last_sync: std::sync::Mutex<Option<String>>,
last_error: std::sync::Mutex<Option<String>>,
transport_stats: Arc<TransportStats>,
observed_protocols: Arc<capabilities::ObservedVersions>,
}
#[derive(Debug, Clone)]
@@ -515,7 +520,11 @@ async fn dht_record_payload_bytes(data_dir: PathBuf, now_ms: u64) -> Result<u64>
}
impl Federation {
pub fn new(library: Arc<Library>, devices: Arc<crate::devices::DeviceSync>) -> Arc<Self> {
pub fn new(
library: Arc<Library>,
devices: Arc<crate::devices::DeviceSync>,
jam: Arc<crate::jam::JamManager>,
) -> Arc<Self> {
let dirs = crate::config::project_dirs();
let data_dir = dirs
.as_ref()
@@ -539,6 +548,7 @@ impl Federation {
Arc::new(Self {
library,
devices,
jam,
data_dir,
cache_dir,
media_dir,
@@ -548,6 +558,7 @@ impl Federation {
last_sync: std::sync::Mutex::new(None),
last_error: std::sync::Mutex::new(initial_error),
transport_stats: Arc::new(TransportStats::default()),
observed_protocols: Arc::new(capabilities::ObservedVersions::default()),
})
}
@@ -555,6 +566,27 @@ impl Federation {
lock(&self.settings).clone()
}
pub async fn create_jam(&self) -> Result<String> {
let service = {
let running = self.running.lock().await;
Arc::clone(
&running
.as_ref()
.context("enable federation before creating a Jam")?
.service,
)
};
let (device_id, device_name) = self.devices.identity_summary()?;
self.jam
.create_or_regenerate(service, device_id, device_name)
.await
}
pub fn join_jam(&self, invite: &str) -> Result<()> {
let (device_id, device_name) = self.devices.identity_summary()?;
self.jam.join(invite, device_id, device_name)
}
fn cached_metadata_snapshot(&self) -> Vec<CachedTrackMetadata> {
lock(&self.metadata_cache).values().cloned().collect()
}
@@ -627,6 +659,10 @@ impl Federation {
.stream_protocol(CATALOG_ALPN)
// Personal-device sync (likes, playlists, trusted devices).
.stream_protocol(crate::devices::SYNC_ALPN)
// Capability-scoped shared playback control.
.stream_protocol(crate::jam::JAM_ALPN)
// Informational application/protocol versions.
.schema_independent_stream_protocol(capabilities::CAPABILITIES_ALPN)
.build()
.map_err(|err| anyhow::anyhow!("invalid federation config: {err}"))?;
let (service, mut events) = MusicDhtService::start(config)
@@ -690,6 +726,23 @@ impl Federation {
let device_tick_task = tokio::spawn(async move {
crate::devices::sync_loop(device_sync, device_service, device_transport).await;
});
let jam_acceptor = service
.stream_acceptor(crate::jam::JAM_ALPN)
.map_err(|err| anyhow::anyhow!("failed to take the Jam acceptor: {err}"))?;
let jam_serve_task =
tokio::spawn(crate::jam::serve_peers(jam_acceptor, Arc::clone(&self.jam)));
let jam_poll_task = tokio::spawn(crate::jam::poll_loop(
Arc::clone(&self.jam),
Arc::clone(&service),
));
let capabilities_acceptor = service
.stream_acceptor(capabilities::CAPABILITIES_ALPN)
.map_err(|err| anyhow::anyhow!("failed to take capabilities acceptor: {err}"))?;
let capabilities_serve_task = tokio::spawn(capabilities::serve(capabilities_acceptor));
let capabilities_probe_task = tokio::spawn(capabilities::probe_loop(
Arc::clone(&service),
Arc::clone(&self.observed_protocols),
));
*guard = Some(Running {
service,
@@ -702,6 +755,10 @@ impl Federation {
catalog_task,
device_sync_task,
device_tick_task,
jam_serve_task,
jam_poll_task,
capabilities_serve_task,
capabilities_probe_task,
],
});
self.set_error(None);
@@ -840,6 +897,7 @@ impl Federation {
network: settings.network_id,
last_sync: lock(&self.last_sync).clone(),
last_error: lock(&self.last_error).clone(),
protocols: ProtocolVersions::snapshot(&self.observed_protocols),
..FedStatus::default()
};
if let Some(running) = guard.as_ref() {
+728
View File
@@ -0,0 +1,728 @@
//! Capability-based Jam playback control for independent federation peers.
use std::collections::{HashMap, HashSet, VecDeque};
use std::sync::{Arc, Mutex, MutexGuard};
use std::time::Duration;
use anyhow::{Context, Result};
use music_dht::{ByteStream, MusicDhtService, PeerTicket, StreamAcceptor};
use serde::{Deserialize, Serialize};
use tokio::sync::mpsc;
use crate::app::event::AppEvent;
use crate::devices::{PlaybackCommand, PlaybackSnapshot};
pub const JAM_ALPN: &[u8] = b"furumi/jam/1";
pub const PROTOCOL_VERSION: u16 = 1;
const MAX_LINE: usize = 8 * 1024 * 1024;
const MAX_COMMANDS: usize = 128;
const PARTICIPANT_TTL_MS: i64 = 30 * 60 * 1_000;
const POLL_INTERVAL: Duration = Duration::from_millis(500);
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum JamRole {
None,
Host,
Participant,
}
#[derive(Debug, Clone)]
pub struct JamStatus {
pub role: JamRole,
pub host_name: Option<String>,
pub invite: Option<String>,
pub participants: Vec<JamParticipant>,
pub connected: bool,
pub last_error: Option<String>,
}
impl Default for JamStatus {
fn default() -> Self {
Self {
role: JamRole::None,
host_name: None,
invite: None,
participants: Vec::new(),
connected: false,
last_error: None,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct JamParticipant {
pub participant_id: String,
pub name: String,
pub last_seen_ms: i64,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
struct JamInvite {
v: u16,
#[serde(rename = "t")]
ticket: String,
#[serde(rename = "j")]
jam_id: String,
#[serde(rename = "s")]
secret: String,
#[serde(rename = "d")]
host_device_id: String,
#[serde(rename = "n")]
host_name: String,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
struct JamCommand {
command_id: String,
participant_id: String,
command: PlaybackCommand,
sent_at_ms: i64,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(tag = "type", rename_all = "snake_case")]
enum WireMessage {
Poll {
version: u16,
jam_id: String,
secret: String,
participant: JamParticipant,
#[serde(default)]
commands: Vec<JamCommand>,
},
Snapshot {
accepted: bool,
#[serde(default)]
error: Option<String>,
#[serde(default)]
acknowledged_command_ids: Vec<String>,
#[serde(default)]
playback: Option<PlaybackSnapshot>,
#[serde(default)]
participants: Vec<JamParticipant>,
host_time_ms: i64,
},
Leave {
version: u16,
jam_id: String,
secret: String,
participant_id: String,
},
}
#[derive(Default)]
struct HostState {
invite: Option<JamInvite>,
invite_uri: Option<String>,
playback: Option<PlaybackSnapshot>,
participants: HashMap<String, JamParticipant>,
seen_commands: HashSet<String>,
seen_order: VecDeque<String>,
}
struct JoinedState {
invite: JamInvite,
participant: JamParticipant,
pending: VecDeque<JamCommand>,
participants: Vec<JamParticipant>,
connected: bool,
last_error: Option<String>,
last_activity_ms: i64,
}
#[derive(Default)]
struct State {
host: HostState,
joined: Option<JoinedState>,
}
pub struct JamManager {
state: Mutex<State>,
event_tx: mpsc::UnboundedSender<AppEvent>,
}
impl JamManager {
pub fn new(event_tx: mpsc::UnboundedSender<AppEvent>) -> Arc<Self> {
Arc::new(Self {
state: Mutex::new(State::default()),
event_tx,
})
}
pub async fn create_or_regenerate(
&self,
service: Arc<MusicDhtService>,
host_device_id: String,
host_name: String,
) -> Result<String> {
let invite = JamInvite {
v: PROTOCOL_VERSION,
ticket: service.ticket().await?.to_string(),
jam_id: format!("jam_{}", random_hex(12)),
secret: random_hex(32),
host_device_id,
host_name,
};
let uri = encode_invite(&invite)?;
let mut state = lock(&self.state);
state.joined = None;
state.host = HostState {
invite: Some(invite),
invite_uri: Some(uri.clone()),
..HostState::default()
};
Ok(uri)
}
pub fn join(&self, uri: &str, participant_id: String, name: String) -> Result<()> {
let invite = parse_invite(uri)?;
let mut state = lock(&self.state);
state.host = HostState::default();
state.joined = Some(JoinedState {
invite,
participant: JamParticipant {
participant_id,
name,
last_seen_ms: now_ms(),
},
pending: VecDeque::new(),
participants: Vec::new(),
connected: false,
last_error: None,
last_activity_ms: now_ms(),
});
Ok(())
}
pub fn leave(&self) {
let mut state = lock(&self.state);
state.joined = None;
state.host = HostState::default();
}
pub fn status(&self) -> JamStatus {
let state = lock(&self.state);
if let Some(joined) = &state.joined {
return JamStatus {
role: JamRole::Participant,
host_name: Some(joined.invite.host_name.clone()),
invite: None,
participants: joined.participants.clone(),
connected: joined.connected,
last_error: joined.last_error.clone(),
};
}
if let Some(invite) = &state.host.invite {
return JamStatus {
role: JamRole::Host,
host_name: Some(invite.host_name.clone()),
invite: state.host.invite_uri.clone(),
participants: state.host.participants.values().cloned().collect(),
connected: true,
last_error: None,
};
}
JamStatus::default()
}
pub fn publish_host_playback(&self, snapshot: PlaybackSnapshot) {
let mut state = lock(&self.state);
if state.host.invite.is_some() {
state.host.playback = Some(snapshot);
}
}
pub fn submit_command(&self, command: PlaybackCommand) -> Result<()> {
let mut state = lock(&self.state);
let joined = state
.joined
.as_mut()
.context("this player is not controlling a Jam")?;
if joined.pending.len() >= MAX_COMMANDS {
joined.pending.pop_front();
}
joined.pending.push_back(JamCommand {
command_id: format!("cmd_{}", random_hex(16)),
participant_id: joined.participant.participant_id.clone(),
command,
sent_at_ms: now_ms(),
});
joined.last_activity_ms = now_ms();
Ok(())
}
async fn poll_once(&self, service: Arc<MusicDhtService>) -> Result<()> {
let (invite, participant, commands) = {
let state = lock(&self.state);
let joined = state.joined.as_ref().context("not joined")?;
(
joined.invite.clone(),
joined.participant.clone(),
joined.pending.iter().cloned().collect::<Vec<_>>(),
)
};
let ticket: PeerTicket = invite.ticket.parse().context("invalid Jam host ticket")?;
let peer = service.connect(ticket).await?;
let mut stream = service.open_stream(peer, JAM_ALPN).await?;
write_message(
&mut stream,
&WireMessage::Poll {
version: PROTOCOL_VERSION,
jam_id: invite.jam_id.clone(),
secret: invite.secret.clone(),
participant,
commands,
},
)
.await?;
stream.send.finish()?;
let response = read_message(&mut stream).await?;
let WireMessage::Snapshot {
accepted,
error,
acknowledged_command_ids,
playback,
participants,
..
} = response
else {
anyhow::bail!("unexpected Jam response");
};
anyhow::ensure!(
accepted,
"{}",
error.unwrap_or_else(|| "Jam refused".into())
);
{
let mut state = lock(&self.state);
let Some(joined) = state.joined.as_mut() else {
return Ok(());
};
if joined.invite.jam_id != invite.jam_id {
return Ok(());
}
let acknowledged = acknowledged_command_ids.into_iter().collect::<HashSet<_>>();
joined
.pending
.retain(|command| !acknowledged.contains(&command.command_id));
joined.participants = participants;
joined.connected = true;
joined.last_error = None;
joined.participant.last_seen_ms = now_ms();
if playback
.as_ref()
.is_some_and(|snapshot| snapshot.state.playing && !snapshot.state.paused)
{
joined.last_activity_ms = now_ms();
}
}
if let Some(playback) = playback {
let _ = self.event_tx.send(AppEvent::JamPlayback(playback));
}
let _ = self.event_tx.send(AppEvent::JamStatus(self.status()));
Ok(())
}
}
pub async fn serve_peers(mut acceptor: StreamAcceptor, manager: Arc<JamManager>) {
while let Some(stream) = acceptor.accept().await {
let manager = Arc::clone(&manager);
tokio::spawn(async move {
if let Err(err) = serve_one(stream, manager).await {
tracing::debug!("Jam stream failed: {err:#}");
}
});
}
}
async fn serve_one(mut stream: ByteStream, manager: Arc<JamManager>) -> Result<()> {
match read_message(&mut stream).await? {
WireMessage::Poll {
version,
jam_id,
secret,
mut participant,
commands,
} => {
let (accepted, error, acknowledged, playback, participants) = {
let mut state = lock(&manager.state);
let valid =
version == PROTOCOL_VERSION
&& state.host.invite.as_ref().is_some_and(|invite| {
invite.jam_id == jam_id && invite.secret == secret
});
if !valid {
(
false,
Some("invalid Jam capability".to_string()),
vec![],
None,
vec![],
)
} else {
let now = now_ms();
state.host.participants.retain(|_, row| {
now.saturating_sub(row.last_seen_ms) <= PARTICIPANT_TTL_MS
});
participant.last_seen_ms = now;
state
.host
.participants
.insert(participant.participant_id.clone(), participant);
let mut acknowledged = Vec::new();
for command in commands.into_iter().take(MAX_COMMANDS) {
acknowledged.push(command.command_id.clone());
if state.host.seen_commands.insert(command.command_id.clone()) {
state.host.seen_order.push_back(command.command_id.clone());
let _ = manager.event_tx.send(AppEvent::JamCommand(command.command));
}
}
while state.host.seen_order.len() > 4096 {
if let Some(id) = state.host.seen_order.pop_front() {
state.host.seen_commands.remove(&id);
}
}
(
true,
None,
acknowledged,
state.host.playback.clone(),
state.host.participants.values().cloned().collect(),
)
}
};
write_message(
&mut stream,
&WireMessage::Snapshot {
accepted,
error,
acknowledged_command_ids: acknowledged,
playback,
participants,
host_time_ms: now_ms(),
},
)
.await?;
stream.send.finish()?;
let _ = tokio::time::timeout(Duration::from_secs(2), stream.send.stopped()).await;
let _ = manager.event_tx.send(AppEvent::JamStatus(manager.status()));
}
WireMessage::Leave {
version,
jam_id,
secret,
participant_id,
} => {
let mut state = lock(&manager.state);
if version == PROTOCOL_VERSION
&& state
.host
.invite
.as_ref()
.is_some_and(|invite| invite.jam_id == jam_id && invite.secret == secret)
{
state.host.participants.remove(&participant_id);
}
drop(state);
let _ = manager.event_tx.send(AppEvent::JamStatus(manager.status()));
}
WireMessage::Snapshot { .. } => anyhow::bail!("unexpected Jam snapshot"),
}
Ok(())
}
pub async fn poll_loop(manager: Arc<JamManager>, service: Arc<MusicDhtService>) {
let mut interval = tokio::time::interval(POLL_INTERVAL);
loop {
interval.tick().await;
if manager.status().role != JamRole::Participant {
continue;
}
let expired = {
let state = lock(&manager.state);
state.joined.as_ref().is_some_and(|joined| {
now_ms().saturating_sub(joined.last_activity_ms) > PARTICIPANT_TTL_MS
})
};
if expired {
manager.leave();
let _ = manager.event_tx.send(AppEvent::JamStatus(manager.status()));
let _ = manager.event_tx.send(AppEvent::StatusMessage(
"Jam ended after 30 minutes without playback or commands".into(),
));
continue;
}
if let Err(err) = manager.poll_once(Arc::clone(&service)).await {
{
let mut state = lock(&manager.state);
if let Some(joined) = state.joined.as_mut() {
joined.connected = false;
joined.last_error = Some(format!("{err:#}"));
}
}
let _ = manager.event_tx.send(AppEvent::JamStatus(manager.status()));
}
}
}
async fn write_message(stream: &mut ByteStream, message: &WireMessage) -> Result<()> {
let mut bytes = serde_json::to_vec(message)?;
anyhow::ensure!(bytes.len() <= MAX_LINE, "Jam message is too large");
bytes.push(b'\n');
stream.send.write_all(&bytes).await?;
Ok(())
}
async fn read_message(stream: &mut ByteStream) -> Result<WireMessage> {
let mut bytes = Vec::new();
let mut byte = [0_u8; 1];
loop {
let read = stream.recv.read(&mut byte).await?;
if read.unwrap_or(0) == 0 || byte[0] == b'\n' {
break;
}
bytes.push(byte[0]);
anyhow::ensure!(bytes.len() <= MAX_LINE, "Jam message is too large");
}
anyhow::ensure!(!bytes.is_empty(), "empty Jam message");
Ok(serde_json::from_slice(&bytes)?)
}
fn encode_invite(invite: &JamInvite) -> Result<String> {
Ok(format!(
"frid://j/{}",
base64url_encode(&serde_json::to_vec(invite)?)
))
}
fn parse_invite(uri: &str) -> Result<JamInvite> {
let encoded = uri
.trim()
.strip_prefix("frid://j/")
.context("expected frid://j invite")?;
let invite: JamInvite = serde_json::from_slice(&base64url_decode(encoded)?)?;
anyhow::ensure!(
invite.v == PROTOCOL_VERSION,
"unsupported Jam invite version"
);
anyhow::ensure!(
!invite.ticket.is_empty()
&& !invite.jam_id.is_empty()
&& invite.secret.len() >= 16
&& !invite.host_device_id.is_empty(),
"incomplete Jam invite"
);
Ok(invite)
}
fn now_ms() -> i64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|duration| duration.as_millis() as i64)
.unwrap_or(0)
}
fn random_hex(bytes: usize) -> String {
use std::sync::atomic::{AtomicU64, Ordering};
static COUNTER: AtomicU64 = AtomicU64::new(1);
let seed = format!(
"{}:{}:{}",
now_ms(),
std::process::id(),
COUNTER.fetch_add(1, Ordering::Relaxed)
);
let hash = blake3::hash(seed.as_bytes()).to_hex().to_string();
hash[..(bytes * 2).min(hash.len())].to_string()
}
fn lock<T>(mutex: &Mutex<T>) -> MutexGuard<'_, T> {
mutex
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
}
fn base64url_encode(bytes: &[u8]) -> String {
const TABLE: &[u8; 64] = b"ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz0123456789-_";
let mut out = String::new();
for chunk in bytes.chunks(3) {
let b0 = chunk[0];
let b1 = chunk.get(1).copied().unwrap_or(0);
let b2 = chunk.get(2).copied().unwrap_or(0);
out.push(TABLE[(b0 >> 2) as usize] as char);
out.push(TABLE[(((b0 & 3) << 4) | (b1 >> 4)) as usize] as char);
if chunk.len() > 1 {
out.push(TABLE[(((b1 & 15) << 2) | (b2 >> 6)) as usize] as char);
}
if chunk.len() > 2 {
out.push(TABLE[(b2 & 63) as usize] as char);
}
}
out
}
fn base64url_decode(value: &str) -> Result<Vec<u8>> {
fn decode(byte: u8) -> Option<u8> {
match byte {
b'A'..=b'Z' => Some(byte - b'A'),
b'a'..=b'z' => Some(byte - b'a' + 26),
b'0'..=b'9' => Some(byte - b'0' + 52),
b'-' => Some(62),
b'_' => Some(63),
_ => None,
}
}
let bytes = value.as_bytes();
anyhow::ensure!(bytes.len() % 4 != 1, "invalid base64url Jam invite");
let mut out = Vec::with_capacity(bytes.len() * 3 / 4);
let mut i = 0;
while i < bytes.len() {
let a = decode(bytes[i]).context("invalid base64url Jam invite")?;
let b = decode(*bytes.get(i + 1).context("truncated Jam invite")?)
.context("invalid base64url Jam invite")?;
let c = bytes.get(i + 2).and_then(|byte| decode(*byte));
let d = bytes.get(i + 3).and_then(|byte| decode(*byte));
out.push((a << 2) | (b >> 4));
if let Some(c) = c {
out.push((b << 4) | (c >> 2));
if let Some(d) = d {
out.push((c << 6) | d);
}
}
i += 4;
}
Ok(out)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn invite_round_trips_and_is_separate_from_pairing() {
let invite = JamInvite {
v: PROTOCOL_VERSION,
ticket: "ticket".into(),
jam_id: "jam_test".into(),
secret: "0123456789abcdef".into(),
host_device_id: "dev_host".into(),
host_name: "Host".into(),
};
let uri = encode_invite(&invite).unwrap();
assert!(uri.starts_with("frid://j/"));
assert_eq!(parse_invite(&uri).unwrap().jam_id, "jam_test");
assert!(parse_invite("frid://i/abcd").is_err());
}
#[tokio::test]
async fn participant_receives_host_state_and_host_receives_command() {
let unique = random_hex(8);
let host_dir = std::env::temp_dir().join(format!("furumi-jam-host-{unique}"));
let guest_dir = std::env::temp_dir().join(format!("furumi-jam-guest-{unique}"));
std::fs::create_dir_all(&host_dir).unwrap();
std::fs::create_dir_all(&guest_dir).unwrap();
let network = music_dht::NetworkId::from_name(&format!("jam-test-{unique}"));
let host_config = music_dht::MusicDhtConfig::builder()
.data_dir(&host_dir)
.network_id(network)
.stream_protocol(JAM_ALPN)
.build()
.unwrap();
let guest_config = music_dht::MusicDhtConfig::builder()
.data_dir(&guest_dir)
.network_id(network)
.stream_protocol(JAM_ALPN)
.build()
.unwrap();
let (host_service, mut host_events) = music_dht::MusicDhtService::start(host_config)
.await
.unwrap();
let (guest_service, mut guest_events) = music_dht::MusicDhtService::start(guest_config)
.await
.unwrap();
let host_service = Arc::new(host_service);
let guest_service = Arc::new(guest_service);
let host_drain = tokio::spawn(async move { while host_events.recv().await.is_some() {} });
let guest_drain = tokio::spawn(async move { while guest_events.recv().await.is_some() {} });
let (host_tx, mut host_rx) = mpsc::unbounded_channel();
let (guest_tx, mut guest_rx) = mpsc::unbounded_channel();
let host = JamManager::new(host_tx);
let guest = JamManager::new(guest_tx);
let invite = host
.create_or_regenerate(Arc::clone(&host_service), "dev_host".into(), "Host".into())
.await
.unwrap();
guest
.join(&invite, "dev_guest".into(), "Guest".into())
.unwrap();
let playback_state = crate::devices::PlaybackStateWire {
queue: Vec::new(),
queue_pos: 0,
playing: false,
paused: false,
idle_since_ms: Some(now_ms()),
position_secs: 0.0,
volume: 73,
shuffle: false,
repeat: crate::devices::PlaybackRepeat::Off,
};
host.publish_host_playback(PlaybackSnapshot {
device_id: "dev_host".into(),
device_name: "Host".into(),
active: true,
updated_at_ms: now_ms(),
state: playback_state.clone(),
});
let acceptor = host_service.stream_acceptor(JAM_ALPN).unwrap();
let server = tokio::spawn(serve_peers(acceptor, Arc::clone(&host)));
guest.poll_once(Arc::clone(&guest_service)).await.unwrap();
let playback = tokio::time::timeout(Duration::from_secs(3), async {
loop {
if let Some(AppEvent::JamPlayback(snapshot)) = guest_rx.recv().await {
break snapshot;
}
}
})
.await
.unwrap();
assert_eq!(playback.device_id, "dev_host");
assert_eq!(playback.state.volume, 73);
guest
.submit_command(PlaybackCommand::SetState {
state: crate::devices::PlaybackStateWire {
paused: true,
..playback_state
},
seek: false,
})
.unwrap();
guest.poll_once(Arc::clone(&guest_service)).await.unwrap();
let command = tokio::time::timeout(Duration::from_secs(3), async {
loop {
if let Some(AppEvent::JamCommand(command)) = host_rx.recv().await {
break command;
}
}
})
.await
.unwrap();
assert!(matches!(
command,
PlaybackCommand::SetState {
state: crate::devices::PlaybackStateWire { paused: true, .. },
seek: false
}
));
server.abort();
host_drain.abort();
guest_drain.abort();
host_service.shutdown().await.unwrap();
guest_service.shutdown().await.unwrap();
let _ = std::fs::remove_dir_all(host_dir);
let _ = std::fs::remove_dir_all(guest_dir);
}
}
+1
View File
@@ -3,6 +3,7 @@ mod art;
mod config;
mod devices;
mod federation;
mod jam;
mod library;
mod media;
mod player;
+164 -3
View File
@@ -2,6 +2,7 @@
use ratatui::Frame;
use ratatui::layout::{Constraint, Layout, Rect};
use ratatui::style::{Color, Modifier, Style};
use ratatui::text::{Line, Span};
use ratatui::widgets::{Block, Paragraph};
@@ -338,6 +339,20 @@ fn draw_settings_rows(frame: &mut Frame, area: Rect, state: &AppState) {
);
}
fn protocol_label(id: &str) -> &str {
match id {
"federation_net" => "Federation transport",
"ticket" => "Peer ticket",
"rendezvous" => "Rendezvous",
"music_dht" => "Music DHT",
"catalog" => "Catalog",
"audio" => "Audio transfer",
"device_sync" => "Device sync",
"jam" => "Jam",
other => other,
}
}
fn draw_section(frame: &mut Frame, area: Rect, state: &AppState, y: &mut u16, title: &'static str) {
if *y >= area.y + area.height {
return;
@@ -485,11 +500,16 @@ fn draw_status(frame: &mut Frame, area: Rect, state: &AppState) {
return;
}
if area.width >= 60 && area.height >= 15 {
let [top_area, _, bottom_area] = Layout::vertical([
if area.width >= 60 && area.height >= 20 {
let protocols_height =
protocol_card_height(state, area.width.saturating_sub(2), area.height);
let [top_area, _, bottom_area, _, protocols_area, _] = Layout::vertical([
Constraint::Length(7),
Constraint::Length(1),
Constraint::Length(7),
Constraint::Length(1),
Constraint::Length(protocols_height),
Constraint::Min(0),
])
.areas(area);
let [status_area, _, local_area] = Layout::horizontal([
@@ -532,10 +552,17 @@ fn draw_status(frame: &mut Frame, area: Rect, state: &AppState) {
" Connected Devices ",
device_summary_lines(state),
);
draw_summary_card(
frame,
protocols_area,
state,
" Protocol Versions ",
protocol_summary_lines(state, protocols_area.width.saturating_sub(2)),
);
return;
}
if area.height < 31 {
if area.height < 39 {
frame.render_widget(
Paragraph::new(compact_status_lines(state))
.wrap(ratatui::widgets::Wrap { trim: false }),
@@ -553,6 +580,8 @@ fn draw_status(frame: &mut Frame, area: Rect, state: &AppState) {
_,
local_area,
_,
protocols_area,
_,
] = Layout::vertical([
Constraint::Length(7),
Constraint::Length(1),
@@ -561,6 +590,12 @@ fn draw_status(frame: &mut Frame, area: Rect, state: &AppState) {
Constraint::Length(7),
Constraint::Length(1),
Constraint::Length(7),
Constraint::Length(1),
Constraint::Length(protocol_card_height(
state,
area.width.saturating_sub(2),
area.height,
)),
Constraint::Min(0),
])
.areas(area);
@@ -593,6 +628,13 @@ fn draw_status(frame: &mut Frame, area: Rect, state: &AppState) {
" Local Data ",
local_data_summary_lines(state),
);
draw_summary_card(
frame,
protocols_area,
state,
" Protocol Versions ",
protocol_summary_lines(state, protocols_area.width.saturating_sub(2)),
);
}
fn compact_status_lines(state: &AppState) -> Vec<Line<'static>> {
@@ -608,6 +650,9 @@ fn compact_status_lines(state: &AppState) -> Vec<Line<'static>> {
lines.push(Line::default());
lines.push(Line::styled("Connected Devices", theme::header_for(state)));
lines.extend(device_summary_lines(state).into_iter().take(2));
lines.push(Line::default());
lines.push(Line::styled("Protocol Versions", theme::header_for(state)));
lines.extend(protocol_summary_lines(state, 0).into_iter().take(3));
lines
}
@@ -637,6 +682,122 @@ fn summary_line(label: &'static str, value: String) -> Line<'static> {
])
}
fn protocol_summary_line(
label: &str,
value: String,
style: Style,
label_width: usize,
) -> Line<'static> {
Line::from(vec![
Span::styled(
format!("{:<label_width$}", protocol_label(label)),
theme::dim(),
),
Span::styled(value, style),
])
}
fn protocol_summary_lines(state: &AppState, width: u16) -> Vec<Line<'static>> {
let Some(status) = state.federation.status.as_ref() else {
return vec![protocol_summary_line(
"status",
"[UNKNOWN]".to_string(),
theme::dim(),
22,
)];
};
let protocols = &status.protocols;
let newer = !protocols.newer.is_empty();
let badge = if newer {
"[NEWER VERSION SEEN]"
} else if status.running && protocols.observed_peers == 0 {
"[CURRENT · waiting for peers]"
} else {
"[CURRENT]"
};
let badge_style = Style::new()
.fg(if newer { Color::LightRed } else { Color::Green })
.add_modifier(Modifier::BOLD);
let mut lines = vec![protocol_summary_line(
"status",
badge.to_string(),
badge_style,
22,
)];
let mut entries = Vec::new();
for (id, local) in &protocols.local {
let observed = protocols.observed.get(id).copied();
let value = match observed {
Some(remote) if remote > *local => format!("local {local} · network {remote}"),
Some(remote) => format!("{local} · seen {remote}"),
None => local.to_string(),
};
let style = if observed.is_some_and(|remote| remote > *local) {
Style::new()
.fg(Color::LightRed)
.add_modifier(Modifier::BOLD)
} else {
Style::default()
};
entries.push((id.as_str(), value, style));
}
if width >= 58 {
let cell_width = width as usize / 2;
let label_width = 22.min(cell_width.saturating_sub(3));
let value_width = cell_width.saturating_sub(label_width + 2);
for pair in entries.chunks(2) {
let mut spans =
protocol_cell_spans(pair[0].0, &pair[0].1, pair[0].2, label_width, value_width);
if let Some(second) = pair.get(1) {
spans.push(Span::styled(" ", theme::dim()));
spans.extend(protocol_cell_spans(
second.0,
&second.1,
second.2,
label_width,
value_width,
));
}
lines.push(Line::from(spans));
}
} else {
lines.extend(
entries
.into_iter()
.map(|(id, value, style)| protocol_summary_line(id, value, style, 22)),
);
}
if newer {
lines.push(Line::styled(
"A newer protocol was observed; update Furumi for compatibility.",
Style::new().fg(Color::LightRed),
));
}
lines
}
fn protocol_card_height(state: &AppState, width: u16, available: u16) -> u16 {
let content = protocol_summary_lines(state, width).len() as u16;
content.saturating_add(2).min(available)
}
fn protocol_cell_spans(
label: &str,
value: &str,
style: Style,
label_width: usize,
value_width: usize,
) -> Vec<Span<'static>> {
let value = value.chars().take(value_width).collect::<String>();
vec![
Span::styled(
format!("{:<label_width$}", protocol_label(label)),
theme::dim(),
),
Span::styled(format!("{value:<value_width$}"), style),
]
}
fn node_summary_lines(state: &AppState) -> Vec<Line<'static>> {
match &state.federation.status {
None => vec![
+56 -10
View File
@@ -440,7 +440,7 @@ fn draw_connected_devices(frame: &mut Frame, state: &AppState, cursor: usize) {
display_lines.push(DisplayLine::Row(index));
}
let height =
(display_lines.len() as u16 + 9).clamp(11, frame.area().height.saturating_sub(2).max(11));
(display_lines.len() as u16 + 14).clamp(16, frame.area().height.saturating_sub(2).max(16));
let area = centered(frame.area(), 76, height);
let block = Block::bordered()
.title(" Connected devices ")
@@ -450,9 +450,10 @@ fn draw_connected_devices(frame: &mut Frame, state: &AppState, cursor: usize) {
frame.render_widget(Clear, area);
frame.render_widget(block, area);
let [summary_area, this_area, other_area, hint_area] = Layout::vertical([
let [summary_area, this_area, jam_area, other_area, hint_area] = Layout::vertical([
Constraint::Length(1),
Constraint::Length(3),
Constraint::Length(4),
Constraint::Min(1),
Constraint::Length(2),
])
@@ -479,12 +480,11 @@ fn draw_connected_devices(frame: &mut Frame, state: &AppState, cursor: usize) {
width: this_area.width,
height: 1,
};
let action_label =
if state.device_playback.role == crate::app::state::DevicePlaybackRole::Active {
"This device is active"
} else {
"Make this device active"
};
let action_label = if state.device_playback.is_audio_owner() {
"This device is active"
} else {
"Make this device active"
};
let self_name = self_row
.map(|row| row.name.as_str())
.unwrap_or(state.device_playback.self_device_name.as_str());
@@ -501,7 +501,53 @@ fn draw_connected_devices(frame: &mut Frame, state: &AppState, cursor: usize) {
state,
);
render_subtitle(frame, other_area, state, "Other devices");
render_subtitle(frame, jam_area, state, "Jam");
let jam_status = match state.jam.role {
crate::jam::JamRole::None => "inactive · h host · J join".to_string(),
crate::jam::JamRole::Host => format!(
"HOST · {} peer(s) · c copy · h new · l leave",
state.jam.participants.len()
),
crate::jam::JamRole::Participant => format!(
"{} · {} · {} participant(s) · l leave",
if state.jam.connected {
"CONNECTED"
} else {
"RECONNECTING"
},
state.jam.host_name.as_deref().unwrap_or("Jam host"),
state.jam.participants.len()
),
};
frame.render_widget(
Paragraph::new(Line::styled(
format!(" {jam_status}"),
if state.jam.role == crate::jam::JamRole::None {
theme::dim()
} else {
theme::accent_for(state)
},
)),
Rect {
x: jam_area.x,
y: jam_area.y + 1,
width: jam_area.width,
height: 1,
},
);
if let Some(error) = state.jam.last_error.as_deref() {
frame.render_widget(
Paragraph::new(Line::styled(format!(" {error}"), theme::dim())),
Rect {
x: jam_area.x,
y: jam_area.y + 2,
width: jam_area.width,
height: 1,
},
);
}
render_subtitle(frame, other_area, state, "My devices");
let list_area = Rect {
x: other_area.x,
y: other_area.y + 1,
@@ -560,7 +606,7 @@ fn draw_connected_devices(frame: &mut Frame, state: &AppState, cursor: usize) {
);
frame.render_widget(
Paragraph::new(Line::styled(
"enter: activate selected device / control active selected · esc close",
"enter device · h host Jam · J join · c copy · l leave · esc close",
theme::dim(),
))
.alignment(Alignment::Center),
+14 -12
View File
@@ -4,6 +4,7 @@ use crate::app::state::{AppState, DevicePlaybackRole};
pub const ACCENT: Color = Color::Cyan;
pub const CONTROL_ACCENT: Color = Color::Yellow;
pub const JAM_ACCENT: Color = Color::Magenta;
pub const DIM: Color = Color::DarkGray;
pub fn accent() -> Style {
@@ -37,10 +38,10 @@ pub fn selection() -> Style {
}
pub fn selection_for(state: &AppState) -> Style {
if state.device_playback.is_control() {
Style::new().fg(Color::White).bg(Color::Rgb(92, 72, 0))
} else {
selection()
match state.device_playback.role {
DevicePlaybackRole::Jam => Style::new().fg(Color::White).bg(Color::Rgb(80, 24, 96)),
DevicePlaybackRole::Control => Style::new().fg(Color::White).bg(Color::Rgb(92, 72, 0)),
DevicePlaybackRole::Active => selection(),
}
}
@@ -49,10 +50,10 @@ pub fn header_for(state: &AppState) -> Style {
}
pub fn border_for(state: &AppState) -> Style {
if state.device_playback.is_control() {
Style::new().fg(CONTROL_ACCENT)
} else {
dim()
match state.device_playback.role {
DevicePlaybackRole::Jam => Style::new().fg(JAM_ACCENT),
DevicePlaybackRole::Control => Style::new().fg(CONTROL_ACCENT),
DevicePlaybackRole::Active => dim(),
}
}
@@ -64,6 +65,7 @@ pub fn role_pill(role: DevicePlaybackRole) -> Style {
let bg = match role {
DevicePlaybackRole::Active => Color::Green,
DevicePlaybackRole::Control => CONTROL_ACCENT,
DevicePlaybackRole::Jam => JAM_ACCENT,
};
Style::new()
.fg(Color::Black)
@@ -72,9 +74,9 @@ pub fn role_pill(role: DevicePlaybackRole) -> Style {
}
fn accent_color_for(state: &AppState) -> Color {
if state.device_playback.is_control() {
CONTROL_ACCENT
} else {
ACCENT
match state.device_playback.role {
DevicePlaybackRole::Jam => JAM_ACCENT,
DevicePlaybackRole::Control => CONTROL_ACCENT,
DevicePlaybackRole::Active => ACCENT,
}
}
+23
View File
@@ -0,0 +1,23 @@
[package]
name = "block"
version = "0.1.6"
authors = ["Steven Sheldon"]
description = "Rust interface for Apple's C language extension of blocks."
keywords = ["blocks", "osx", "ios", "objective-c"]
readme = "README.md"
repository = "http://github.com/SSheldon/rust-block"
documentation = "http://ssheldon.github.io/rust-objc/block/"
license = "MIT"
exclude = [
".gitignore",
".travis.yml",
"travis_install.sh",
"travis_test.sh",
"tests-ios/**",
]
[dev-dependencies.objc_test_utils]
version = "0.0"
path = "test_utils"
+42
View File
@@ -0,0 +1,42 @@
Rust interface for Apple's C language extension of blocks.
For more information on the specifics of the block implementation, see
Clang's documentation: http://clang.llvm.org/docs/Block-ABI-Apple.html
## Invoking blocks
The `Block` struct is used for invoking blocks from Objective-C. For example,
consider this Objective-C function:
``` objc
int32_t sum(int32_t (^block)(int32_t, int32_t)) {
return block(5, 8);
}
```
We could write it in Rust as the following:
``` rust
unsafe fn sum(block: &Block<(i32, i32), i32>) -> i32 {
block.call((5, 8))
}
```
Note the extra parentheses in the `call` method, since the arguments must be
passed as a tuple.
## Creating blocks
Creating a block to pass to Objective-C can be done with the `ConcreteBlock`
struct. For example, to create a block that adds two `i32`s, we could write:
``` rust
let block = ConcreteBlock::new(|a: i32, b: i32| a + b);
let block = block.copy();
assert!(unsafe { block.call((5, 8)) } == 13);
```
It is important to copy your block to the heap (with the `copy` method) before
passing it to Objective-C; this is because our `ConcreteBlock` is only meant
to be copied once, and we can enforce this in Rust, but if Objective-C code
were to copy it twice we could have a double free.
+399
View File
@@ -0,0 +1,399 @@
/*!
A Rust interface for Objective-C blocks.
For more information on the specifics of the block implementation, see
Clang's documentation: http://clang.llvm.org/docs/Block-ABI-Apple.html
# Invoking blocks
The `Block` struct is used for invoking blocks from Objective-C. For example,
consider this Objective-C function:
``` objc
int32_t sum(int32_t (^block)(int32_t, int32_t)) {
return block(5, 8);
}
```
We could write it in Rust as the following:
```
# use block::Block;
unsafe fn sum(block: &Block<(i32, i32), i32>) -> i32 {
block.call((5, 8))
}
```
Note the extra parentheses in the `call` method, since the arguments must be
passed as a tuple.
# Creating blocks
Creating a block to pass to Objective-C can be done with the `ConcreteBlock`
struct. For example, to create a block that adds two `i32`s, we could write:
```
# use block::ConcreteBlock;
let block = ConcreteBlock::new(|a: i32, b: i32| a + b);
let block = block.copy();
assert!(unsafe { block.call((5, 8)) } == 13);
```
It is important to copy your block to the heap (with the `copy` method) before
passing it to Objective-C; this is because our `ConcreteBlock` is only meant
to be copied once, and we can enforce this in Rust, but if Objective-C code
were to copy it twice we could have a double free.
*/
#[cfg(test)]
mod test_utils;
use std::marker::PhantomData;
use std::mem;
use std::ops::{Deref, DerefMut};
use std::os::raw::{c_int, c_ulong, c_void};
use std::ptr;
#[repr(C)]
struct Class {
_private: [u8; 0],
}
#[cfg_attr(any(target_os = "macos", target_os = "ios"),
link(name = "System", kind = "dylib"))]
#[cfg_attr(not(any(target_os = "macos", target_os = "ios")),
link(name = "BlocksRuntime", kind = "dylib"))]
extern "C" {
static _NSConcreteStackBlock: Class;
fn _Block_copy(block: *const c_void) -> *mut c_void;
fn _Block_release(block: *const c_void);
}
/// Types that may be used as the arguments to an Objective-C block.
pub trait BlockArguments: Sized {
/// Calls the given `Block` with self as the arguments.
///
/// Unsafe because `block` must point to a valid `Block` and this invokes
/// foreign code whose safety the compiler cannot verify.
unsafe fn call_block<R>(self, block: *mut Block<Self, R>) -> R;
}
macro_rules! block_args_impl {
($($a:ident : $t:ident),*) => (
impl<$($t),*> BlockArguments for ($($t,)*) {
unsafe fn call_block<R>(self, block: *mut Block<Self, R>) -> R {
let invoke: unsafe extern "C" fn(*mut Block<Self, R> $(, $t)*) -> R = {
let base = block as *mut BlockBase<Self, R>;
mem::transmute((*base).invoke)
};
let ($($a,)*) = self;
invoke(block $(, $a)*)
}
}
);
}
block_args_impl!();
block_args_impl!(a: A);
block_args_impl!(a: A, b: B);
block_args_impl!(a: A, b: B, c: C);
block_args_impl!(a: A, b: B, c: C, d: D);
block_args_impl!(a: A, b: B, c: C, d: D, e: E);
block_args_impl!(a: A, b: B, c: C, d: D, e: E, f: F);
block_args_impl!(a: A, b: B, c: C, d: D, e: E, f: F, g: G);
block_args_impl!(a: A, b: B, c: C, d: D, e: E, f: F, g: G, h: H);
block_args_impl!(a: A, b: B, c: C, d: D, e: E, f: F, g: G, h: H, i: I);
block_args_impl!(a: A, b: B, c: C, d: D, e: E, f: F, g: G, h: H, i: I, j: J);
block_args_impl!(a: A, b: B, c: C, d: D, e: E, f: F, g: G, h: H, i: I, j: J, k: K);
block_args_impl!(a: A, b: B, c: C, d: D, e: E, f: F, g: G, h: H, i: I, j: J, k: K, l: L);
#[repr(C)]
struct BlockBase<A, R> {
isa: *const Class,
flags: c_int,
_reserved: c_int,
invoke: unsafe extern "C" fn(*mut Block<A, R>, ...) -> R,
}
/// An Objective-C block that takes arguments of `A` when called and
/// returns a value of `R`.
#[repr(C)]
pub struct Block<A, R> {
_base: PhantomData<BlockBase<A, R>>,
}
impl<A: BlockArguments, R> Block<A, R> where A: BlockArguments {
/// Call self with the given arguments.
///
/// Unsafe because this invokes foreign code that the caller must verify
/// doesn't violate any of Rust's safety rules. For example, if this block
/// is shared with multiple references, the caller must ensure that calling
/// it will not cause a data race.
pub unsafe fn call(&self, args: A) -> R {
args.call_block(self as *const _ as *mut _)
}
}
/// A reference-counted Objective-C block.
pub struct RcBlock<A, R> {
ptr: *mut Block<A, R>,
}
impl<A, R> RcBlock<A, R> {
/// Construct an `RcBlock` for the given block without copying it.
/// The caller must ensure the block has a +1 reference count.
///
/// Unsafe because `ptr` must point to a valid `Block` and must have a +1
/// reference count or it will be overreleased when the `RcBlock` is
/// dropped.
pub unsafe fn new(ptr: *mut Block<A, R>) -> Self {
RcBlock { ptr: ptr }
}
/// Constructs an `RcBlock` by copying the given block.
///
/// Unsafe because `ptr` must point to a valid `Block`.
pub unsafe fn copy(ptr: *mut Block<A, R>) -> Self {
let ptr = _Block_copy(ptr as *const c_void) as *mut Block<A, R>;
RcBlock { ptr: ptr }
}
}
impl<A, R> Clone for RcBlock<A, R> {
fn clone(&self) -> RcBlock<A, R> {
unsafe {
RcBlock::copy(self.ptr)
}
}
}
impl<A, R> Deref for RcBlock<A, R> {
type Target = Block<A, R>;
fn deref(&self) -> &Block<A, R> {
unsafe { &*self.ptr }
}
}
impl<A, R> Drop for RcBlock<A, R> {
fn drop(&mut self) {
unsafe {
_Block_release(self.ptr as *const c_void);
}
}
}
/// Types that may be converted into a `ConcreteBlock`.
pub trait IntoConcreteBlock<A>: Sized where A: BlockArguments {
/// The return type of the resulting `ConcreteBlock`.
type Ret;
/// Consumes self to create a `ConcreteBlock`.
fn into_concrete_block(self) -> ConcreteBlock<A, Self::Ret, Self>;
}
macro_rules! concrete_block_impl {
($f:ident) => (
concrete_block_impl!($f,);
);
($f:ident, $($a:ident : $t:ident),*) => (
impl<$($t,)* R, X> IntoConcreteBlock<($($t,)*)> for X
where X: Fn($($t,)*) -> R {
type Ret = R;
fn into_concrete_block(self) -> ConcreteBlock<($($t,)*), R, X> {
unsafe extern "C" fn $f<$($t,)* R, X>(
block_ptr: *mut ConcreteBlock<($($t,)*), R, X>
$(, $a: $t)*) -> R
where X: Fn($($t,)*) -> R {
let block = &*block_ptr;
(block.closure)($($a),*)
}
let f: unsafe extern "C" fn(*mut ConcreteBlock<($($t,)*), R, X> $(, $a: $t)*) -> R = $f;
unsafe {
ConcreteBlock::with_invoke(mem::transmute(f), self)
}
}
}
);
}
concrete_block_impl!(concrete_block_invoke_args0);
concrete_block_impl!(concrete_block_invoke_args1, a: A);
concrete_block_impl!(concrete_block_invoke_args2, a: A, b: B);
concrete_block_impl!(concrete_block_invoke_args3, a: A, b: B, c: C);
concrete_block_impl!(concrete_block_invoke_args4, a: A, b: B, c: C, d: D);
concrete_block_impl!(concrete_block_invoke_args5, a: A, b: B, c: C, d: D, e: E);
concrete_block_impl!(concrete_block_invoke_args6, a: A, b: B, c: C, d: D, e: E, f: F);
concrete_block_impl!(concrete_block_invoke_args7, a: A, b: B, c: C, d: D, e: E, f: F, g: G);
concrete_block_impl!(concrete_block_invoke_args8, a: A, b: B, c: C, d: D, e: E, f: F, g: G, h: H);
concrete_block_impl!(concrete_block_invoke_args9, a: A, b: B, c: C, d: D, e: E, f: F, g: G, h: H, i: I);
concrete_block_impl!(concrete_block_invoke_args10, a: A, b: B, c: C, d: D, e: E, f: F, g: G, h: H, i: I, j: J);
concrete_block_impl!(concrete_block_invoke_args11, a: A, b: B, c: C, d: D, e: E, f: F, g: G, h: H, i: I, j: J, k: K);
concrete_block_impl!(concrete_block_invoke_args12, a: A, b: B, c: C, d: D, e: E, f: F, g: G, h: H, i: I, j: J, k: K, l: L);
/// An Objective-C block whose size is known at compile time and may be
/// constructed on the stack.
#[repr(C)]
pub struct ConcreteBlock<A, R, F> {
base: BlockBase<A, R>,
descriptor: Box<BlockDescriptor<ConcreteBlock<A, R, F>>>,
closure: F,
}
impl<A, R, F> ConcreteBlock<A, R, F>
where A: BlockArguments, F: IntoConcreteBlock<A, Ret=R> {
/// Constructs a `ConcreteBlock` with the given closure.
/// When the block is called, it will return the value that results from
/// calling the closure.
pub fn new(closure: F) -> Self {
closure.into_concrete_block()
}
}
impl<A, R, F> ConcreteBlock<A, R, F> {
/// Constructs a `ConcreteBlock` with the given invoke function and closure.
/// Unsafe because the caller must ensure the invoke function takes the
/// correct arguments.
unsafe fn with_invoke(invoke: unsafe extern "C" fn(*mut Self, ...) -> R,
closure: F) -> Self {
ConcreteBlock {
base: BlockBase {
isa: &_NSConcreteStackBlock,
// 1 << 25 = BLOCK_HAS_COPY_DISPOSE
flags: 1 << 25,
_reserved: 0,
invoke: mem::transmute(invoke),
},
descriptor: Box::new(BlockDescriptor::new()),
closure: closure,
}
}
}
impl<A, R, F> ConcreteBlock<A, R, F> where F: 'static {
/// Copy self onto the heap as an `RcBlock`.
pub fn copy(self) -> RcBlock<A, R> {
unsafe {
let mut block = self;
let copied = RcBlock::copy(&mut *block);
// At this point, our copy helper has been run so the block will
// be moved to the heap and we can forget the original block
// because the heap block will drop in our dispose helper.
mem::forget(block);
copied
}
}
}
impl<A, R, F> Clone for ConcreteBlock<A, R, F> where F: Clone {
fn clone(&self) -> Self {
unsafe {
ConcreteBlock::with_invoke(mem::transmute(self.base.invoke),
self.closure.clone())
}
}
}
impl<A, R, F> Deref for ConcreteBlock<A, R, F> {
type Target = Block<A, R>;
fn deref(&self) -> &Block<A, R> {
unsafe { &*(&self.base as *const _ as *const Block<A, R>) }
}
}
impl<A, R, F> DerefMut for ConcreteBlock<A, R, F> {
fn deref_mut(&mut self) -> &mut Block<A, R> {
unsafe { &mut *(&mut self.base as *mut _ as *mut Block<A, R>) }
}
}
unsafe extern "C" fn block_context_dispose<B>(block: &mut B) {
// Read the block onto the stack and let it drop
ptr::read(block);
}
unsafe extern "C" fn block_context_copy<B>(_dst: &mut B, _src: &B) {
// The runtime memmoves the src block into the dst block, nothing to do
}
#[repr(C)]
struct BlockDescriptor<B> {
_reserved: c_ulong,
block_size: c_ulong,
copy_helper: unsafe extern "C" fn(&mut B, &B),
dispose_helper: unsafe extern "C" fn(&mut B),
}
impl<B> BlockDescriptor<B> {
fn new() -> BlockDescriptor<B> {
BlockDescriptor {
_reserved: 0,
block_size: mem::size_of::<B>() as c_ulong,
copy_helper: block_context_copy::<B>,
dispose_helper: block_context_dispose::<B>,
}
}
}
#[cfg(test)]
mod tests {
use test_utils::*;
use super::{ConcreteBlock, RcBlock};
#[test]
fn test_call_block() {
let block = get_int_block_with(13);
unsafe {
assert!(block.call(()) == 13);
}
}
#[test]
fn test_call_block_args() {
let block = get_add_block_with(13);
unsafe {
assert!(block.call((2,)) == 15);
}
}
#[test]
fn test_create_block() {
let block = ConcreteBlock::new(|| 13);
let result = invoke_int_block(&block);
assert!(result == 13);
}
#[test]
fn test_create_block_args() {
let block = ConcreteBlock::new(|a: i32| a + 5);
let result = invoke_add_block(&block, 6);
assert!(result == 11);
}
#[test]
fn test_concrete_block_copy() {
let s = "Hello!".to_string();
let expected_len = s.len() as i32;
let block = ConcreteBlock::new(move || s.len() as i32);
assert!(invoke_int_block(&block) == expected_len);
let copied = block.copy();
assert!(invoke_int_block(&copied) == expected_len);
}
#[test]
fn test_concrete_block_stack_copy() {
fn make_block() -> RcBlock<(), i32> {
let x = 7;
let block = ConcreteBlock::new(move || x);
block.copy()
}
let block = make_block();
assert!(invoke_int_block(&block) == 7);
}
}
+31
View File
@@ -0,0 +1,31 @@
extern crate objc_test_utils;
use {Block, RcBlock};
pub fn get_int_block_with(i: i32) -> RcBlock<(), i32> {
unsafe {
let ptr = objc_test_utils::get_int_block_with(i);
RcBlock::new(ptr as *mut _)
}
}
pub fn get_add_block_with(i: i32) -> RcBlock<(i32,), i32> {
unsafe {
let ptr = objc_test_utils::get_add_block_with(i);
RcBlock::new(ptr as *mut _)
}
}
pub fn invoke_int_block(block: &Block<(), i32>) -> i32 {
let ptr = block as *const _;
unsafe {
objc_test_utils::invoke_int_block(ptr as *mut _)
}
}
pub fn invoke_add_block(block: &Block<(i32,), i32>, a: i32) -> i32 {
let ptr = block as *const _;
unsafe {
objc_test_utils::invoke_add_block(ptr as *mut _, a)
}
}