Init commit

This commit is contained in:
Ultradesu
2026-08-12 01:30:57 +01:00
commit 423ea59371
67 changed files with 23280 additions and 0 deletions
+481
View File
@@ -0,0 +1,481 @@
use super::{
Actor, ArtistKey, ArtistRef, AudioSource, ConnectedDeviceSnapshot, ConnectedDevicesSnapshot,
ContentId, ControlPlaybackAnchor, DeviceOperationResult, DevicePlaybackRole, DevicePresence,
DeviceTrust, Duration, Instant, InternalEvent, PendingControlState, PendingPairingSnapshot,
PlaybackStatus, ReleaseKey, Track, TrackKey, library_track, normalize_device_name,
portable_playback_placeholder, track_to_playback_track, unix_time_ms, volume_percent,
};
impl Actor {
pub(super) fn create_device_invite(&mut self) {
let Some(service) = self.device_service.clone() else {
self.state.connected_devices.error =
Some("Federation network is still starting".into());
self.publish();
return;
};
self.state.connected_devices.busy = true;
self.state.connected_devices.error = None;
self.publish();
let devices = std::sync::Arc::clone(&self.devices);
let internal = self.internal.clone();
tokio::spawn(async move {
let result = devices
.create_invite(service)
.await
.map(DeviceOperationResult::Invite)
.map_err(|error| format!("{error:#}"));
let _ = internal
.send(InternalEvent::DeviceOperationFinished(result))
.await;
});
}
pub(super) fn connect_device(&mut self, invite: String) {
let Some(service) = self.device_service.clone() else {
self.state.connected_devices.error =
Some("Federation network is still starting".into());
self.publish();
return;
};
self.state.connected_devices.busy = true;
self.state.connected_devices.error = None;
self.publish();
let devices = std::sync::Arc::clone(&self.devices);
let internal = self.internal.clone();
tokio::spawn(async move {
let result = devices
.connect_invite(service, &invite)
.await
.map(DeviceOperationResult::Connected)
.map_err(|error| format!("{error:#}"));
let _ = internal
.send(InternalEvent::DeviceOperationFinished(result))
.await;
});
}
pub(super) fn apply_local_device_name(&mut self, name: &str) {
let name = normalize_device_name(name);
if let Err(error) = self.devices.set_device_name(&name, None) {
self.state.connected_devices.error = Some(format!("device name: {error:#}"));
return;
}
if self.active_device_id == self.state.connected_devices.this_device_id {
self.active_device_name = name;
}
self.refresh_connected_devices();
}
pub(super) fn schedule_device_name_publish(&self, candidate: String) {
let internal = self.internal.clone();
tokio::spawn(async move {
tokio::time::sleep(Duration::from_millis(500)).await;
let _ = internal
.send(InternalEvent::DeviceNamePublishDue(candidate))
.await;
});
}
pub(super) fn publish_device_name(&self, name: String) {
let Some(service) = self.device_service.clone() else {
return;
};
let devices = std::sync::Arc::clone(&self.devices);
let internal = self.internal.clone();
tokio::spawn(async move {
let result = async {
let ticket = service
.ticket()
.await
.map_err(|error| format!("device name: {error:#}"))?;
devices
.set_device_name(&name, Some(&ticket.to_string()))
.map_err(|error| format!("device name: {error:#}"))
}
.await;
let _ = internal
.send(InternalEvent::DeviceNamePublished(result))
.await;
});
}
pub(super) fn refresh_connected_devices(&mut self) {
let status = self.devices.status();
let invite = self.state.connected_devices.invite.take();
let busy = self.state.connected_devices.busy;
let existing_error = self.state.connected_devices.error.take();
self.state.connected_devices = ConnectedDevicesSnapshot {
this_device_id: status.this_device_id,
this_device_name: status.this_device_name,
group_id: status.group_id,
role: self.device_role,
active_device_id: self.active_device_id.clone(),
active_device_name: self.active_device_name.clone(),
devices: status
.devices
.into_iter()
.filter(|device| !device.revoked)
.map(|device| ConnectedDeviceSnapshot {
is_active: device.id == self.active_device_id,
id: device.id,
name: device.name,
client_version: device.client_version,
is_self: device.is_self,
presence: if device.online {
DevicePresence::Online
} else {
DevicePresence::Offline
},
trust: if device.revoked {
DeviceTrust::Revoked
} else {
DeviceTrust::Trusted
},
})
.collect(),
pending_pairings: status
.pending
.into_iter()
.map(|pending| PendingPairingSnapshot {
request_id: pending.request_id,
device_id: pending.device_id,
name: pending.name,
client_version: pending.client_version,
requester_group_id: pending.requester_group_id,
requester_group_active_devices: pending.requester_group_active_devices,
})
.collect(),
invite,
busy,
last_sync: status.last_sync,
error: existing_error.or(status.error),
};
}
pub(super) fn select_playback_device(&mut self, device_id: &str) {
// Selecting the device that already owns playback must not rebuild or
// restart the current audio stream.
if device_id == self.active_device_id {
return;
}
let target = self
.state
.connected_devices
.devices
.iter()
.find(|device| device.id == device_id && device.trust == DeviceTrust::Trusted)
.cloned();
let Some(target) = target else {
return;
};
let wire = self.playback_wire_state();
let previous = self.active_device_id.clone();
let mut urgent_targets = Vec::with_capacity(2);
if target.is_self {
if previous != target.id {
let command = music_dht::device_sync::PlaybackCommand::ActiveChanged {
active_device_id: target.id.clone(),
active_device_name: target.name.clone(),
state: wire.clone(),
};
if let Err(error) = self.devices.record_playback_command(&previous, command) {
self.state.connected_devices.error = Some(format!("device handoff: {error:#}"));
self.publish();
return;
}
urgent_targets.push(previous.clone());
}
self.device_role = DevicePlaybackRole::Active;
self.active_device_id = target.id;
self.active_device_name = target.name;
self.control_anchor = None;
self.pending_control = None;
if self.state.playback.status == PlaybackStatus::Playing {
self.play_current();
}
} else {
let command = music_dht::device_sync::PlaybackCommand::ActiveChanged {
active_device_id: target.id.clone(),
active_device_name: target.name.clone(),
state: wire.clone(),
};
if let Err(error) = self
.devices
.record_playback_command(&target.id, command.clone())
{
self.state.connected_devices.error = Some(format!("device handoff: {error:#}"));
self.publish();
return;
}
urgent_targets.push(target.id.clone());
if previous != target.id && previous != self.state.connected_devices.this_device_id {
if let Err(error) = self.devices.record_playback_command(&previous, command) {
self.state.connected_devices.error =
Some(format!("previous device handoff: {error:#}"));
}
urgent_targets.push(previous);
}
self.audio.stop();
self.device_role = DevicePlaybackRole::Control;
self.active_device_id.clone_from(&target.id);
self.active_device_name = target.name;
self.control_anchor = Some(ControlPlaybackAnchor {
device_id: target.id.clone(),
state: wire.clone(),
observed_at: Instant::now(),
});
self.pending_control = Some(PendingControlState {
device_id: target.id,
state: wire,
seek: true,
sent_at: Instant::now(),
});
}
self.devices.request_sync();
if let Some(service) = self.device_service.clone() {
urgent_targets.sort();
urgent_targets.dedup();
for target_id in urgent_targets {
let devices = std::sync::Arc::clone(&self.devices);
let service = std::sync::Arc::clone(&service);
tokio::spawn(async move {
let _ = devices.sync_target(service, &target_id).await;
});
}
}
self.refresh_connected_devices();
self.publish_device_playback();
self.publish();
}
pub(super) fn playback_wire_state(&self) -> music_dht::device_sync::PlaybackStateWire {
music_dht::device_sync::PlaybackStateWire {
queue: self
.state
.queue
.items()
.iter()
.map(|item| track_to_playback_track(&item.track))
.collect(),
queue_pos: self.state.queue.current_index().unwrap_or(0),
playing: self.state.playback.status != PlaybackStatus::Stopped,
paused: self.state.playback.status == PlaybackStatus::Paused,
idle_since_ms: (self.state.playback.status != PlaybackStatus::Playing)
.then_some(unix_time_ms()),
position_secs: self.state.playback.position_seconds,
volume: volume_percent(self.state.playback.volume),
shuffle: false,
repeat: music_dht::device_sync::PlaybackRepeat::Off,
}
}
pub(super) fn publish_device_playback(&self) {
let snapshot = music_dht::device_sync::PlaybackSnapshot {
device_id: self.state.connected_devices.this_device_id.clone(),
device_name: self.state.connected_devices.this_device_name.clone(),
active: self.device_role == DevicePlaybackRole::Active,
updated_at_ms: unix_time_ms(),
state: self.playback_wire_state(),
};
self.devices.publish_playback(snapshot);
}
pub(super) fn send_control_state(&mut self, seek: bool) {
if self.device_role != DevicePlaybackRole::Control {
return;
}
let state = self.playback_wire_state();
let command = music_dht::device_sync::PlaybackCommand::SetState {
state: state.clone(),
seek,
};
if let Err(error) = self
.devices
.record_playback_command(&self.active_device_id, command)
{
self.state.connected_devices.error = Some(format!("device control: {error:#}"));
} else {
self.pending_control = Some(PendingControlState {
device_id: self.active_device_id.clone(),
state,
seek,
sent_at: Instant::now(),
});
}
}
pub(super) fn apply_device_playback_state(
&mut self,
wire: &music_dht::device_sync::PlaybackStateWire,
start_audio: bool,
seek: bool,
) {
let tracks = wire
.queue
.iter()
.filter_map(|track| self.resolve_playback_track(track))
.collect::<Vec<_>>();
if !tracks.is_empty() {
self.state.queue.replace_context(tracks, wire.queue_pos);
self.resolve_queue_artwork();
}
self.state.playback.volume = f32::from(wire.volume.min(100)) / 100.0;
self.state.playback.position_seconds = wire.position_secs.max(0.0);
self.state.playback.duration_seconds = self
.state
.queue
.current()
.map_or(0.0, |item| item.track.duration_seconds);
self.state.playback.status = if !wire.playing {
PlaybackStatus::Stopped
} else if wire.paused {
PlaybackStatus::Paused
} else {
PlaybackStatus::Playing
};
self.audio.set_volume(self.state.playback.volume);
if !start_audio {
return;
}
if wire.playing {
self.play_current();
if seek || wire.position_secs > 0.0 {
self.audio
.seek(Duration::from_secs_f64(wire.position_secs.max(0.0)));
self.state.playback.position_seconds = wire.position_secs.max(0.0);
}
if wire.paused {
self.audio.pause();
self.state.playback.status = PlaybackStatus::Paused;
}
} else {
self.audio.stop();
}
}
pub(super) fn apply_device_playback_command(
&mut self,
command: music_dht::device_sync::PlaybackCommand,
) {
match command {
music_dht::device_sync::PlaybackCommand::SetState { state, seek } => {
self.device_role = DevicePlaybackRole::Active;
self.active_device_id = self.state.connected_devices.this_device_id.clone();
self.active_device_name = self.state.connected_devices.this_device_name.clone();
self.control_anchor = None;
self.pending_control = None;
self.apply_device_playback_state(&state, true, seek);
}
music_dht::device_sync::PlaybackCommand::ActiveChanged {
active_device_id,
active_device_name,
state,
} => {
let is_self = active_device_id == self.state.connected_devices.this_device_id;
self.active_device_id = active_device_id;
self.active_device_name = active_device_name;
self.device_role = if is_self {
DevicePlaybackRole::Active
} else {
DevicePlaybackRole::Control
};
self.pending_control = None;
self.control_anchor = (!is_self).then(|| ControlPlaybackAnchor {
device_id: self.active_device_id.clone(),
state: state.clone(),
observed_at: Instant::now(),
});
if !is_self {
self.audio.stop();
}
self.apply_device_playback_state(&state, is_self, true);
}
}
self.refresh_connected_devices();
self.publish();
}
pub(super) fn resolve_playback_track(
&self,
wire: &music_dht::device_sync::PlaybackTrack,
) -> Option<Track> {
if let Some(content_id) = wire
.content_id
.as_deref()
.and_then(|id| ContentId::parse(id).ok())
{
let key = TrackKey::remote(content_id.clone());
if let Some(track) = self.track(&key) {
return Some(track.clone());
}
if let Ok(Some(track)) = self.catalog.track_by_content_id(content_id.as_str()) {
return Some(library_track(track, ""));
}
}
let Some(fed) = wire.fed.as_ref() else {
let content_id = wire
.content_id
.as_deref()
.and_then(|id| ContentId::parse(id).ok())?;
return Some(portable_playback_placeholder(wire, content_id));
};
let content_id = ContentId::parse(&fed.content_id).ok()?;
let release_key = ReleaseKey::Federation {
peer_id: fed.owner.clone(),
id: format!(
"name:{}",
music_dht::normalize_name(fed.release_title.as_deref().unwrap_or_default())
),
};
let refs = |names: &[String]| {
names
.iter()
.map(|name| ArtistRef {
key: ArtistKey::Federation {
peer_id: fed.owner.clone(),
id: format!("name:{}", music_dht::normalize_name(name)),
},
name: name.clone(),
})
.collect::<Vec<_>>()
};
Some(Track {
key: TrackKey::federation(
fed.owner.clone(),
fed.item_id.clone(),
Some(content_id.clone()),
),
title: wire.title.clone(),
artist: wire.artist_names.join(", "),
artists: refs(&wire.artist_names),
featured_artists: refs(&wire.featured_artist_names),
release: wire.release_title.clone(),
release_id: release_key,
duration_seconds: wire.duration_seconds,
track_number: wire
.track_number
.and_then(|value| u32::try_from(value).ok()),
disc_number: wire.disc_number.and_then(|value| u32::try_from(value).ok()),
cover_uri: None,
audio_format: wire.audio_format.clone(),
audio_bitrate_kbps: wire
.audio_bitrate
.and_then(|value| u32::try_from(value).ok()),
audio_sample_rate_hz: wire
.audio_sample_rate
.and_then(|value| u32::try_from(value).ok()),
audio_bit_depth: wire
.audio_bit_depth
.and_then(|value| u32::try_from(value).ok()),
file_size_bytes: wire
.file_size_bytes
.and_then(|value| u64::try_from(value).ok()),
liked: false,
audio_source: AudioSource::Federation {
peer_id: fed.owner.clone(),
content_id,
},
})
}
}
+281
View File
@@ -0,0 +1,281 @@
//! Real audio output isolated from the backend actor and UI thread.
use std::fs::File;
use std::io;
use std::path::PathBuf;
use std::sync::Arc;
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::mpsc::{Receiver, RecvTimeoutError, Sender};
use std::time::Duration;
use rodio::{Decoder, DeviceSinkBuilder, Player, stream::MixerDeviceSink};
pub trait TrackReadSeek: std::io::Read + std::io::Seek + Send + Sync {}
impl<T> TrackReadSeek for T where T: std::io::Read + std::io::Seek + Send + Sync {}
pub type TrackReader = Box<dyn TrackReadSeek>;
#[derive(Debug)]
pub enum Event {
Started,
Finished,
Failed(String),
}
enum Command {
Play {
path: PathBuf,
volume: f32,
},
PlayStream {
reader: TrackReader,
mime_type: String,
volume: f32,
},
Pause,
Resume,
Stop,
Seek(Duration),
SetVolume(f32),
}
#[derive(Debug, Default)]
pub struct Shared {
position_ms: AtomicU64,
}
impl Shared {
pub fn position_seconds(&self) -> f64 {
Duration::from_millis(self.position_ms.load(Ordering::Relaxed)).as_secs_f64()
}
}
#[derive(Clone)]
pub struct Controller {
commands: Sender<Command>,
pub shared: Arc<Shared>,
}
impl Controller {
pub fn play(&self, path: PathBuf, volume: f32) {
self.shared.position_ms.store(0, Ordering::Relaxed);
let _ = self.commands.send(Command::Play { path, volume });
}
pub fn play_stream(&self, reader: TrackReader, mime_type: String, volume: f32) {
self.shared.position_ms.store(0, Ordering::Relaxed);
let _ = self.commands.send(Command::PlayStream {
reader,
mime_type,
volume,
});
}
pub fn pause(&self) {
let _ = self.commands.send(Command::Pause);
}
pub fn resume(&self) {
let _ = self.commands.send(Command::Resume);
}
pub fn stop(&self) {
let _ = self.commands.send(Command::Stop);
}
pub fn seek(&self, position: Duration) {
let millis = u64::try_from(position.as_millis()).unwrap_or(u64::MAX);
self.shared.position_ms.store(millis, Ordering::Relaxed);
let _ = self.commands.send(Command::Seek(position));
}
pub fn set_volume(&self, volume: f32) {
let _ = self.commands.send(Command::SetVolume(volume));
}
}
pub fn spawn(on_event: impl Fn(Event) + Send + 'static) -> io::Result<Controller> {
let (commands, receiver) = std::sync::mpsc::channel();
let shared = Arc::new(Shared::default());
let thread_shared = Arc::clone(&shared);
std::thread::Builder::new()
.name("furumi-audio".into())
.spawn(move || run(&receiver, &thread_shared, &on_event))?;
Ok(Controller { commands, shared })
}
struct Output {
_device: MixerDeviceSink,
player: Player,
}
fn run(receiver: &Receiver<Command>, shared: &Arc<Shared>, on_event: &impl Fn(Event)) {
let mut output: Option<Output> = None;
let mut track_loaded = false;
let mut previous_queue_len = 0;
loop {
match receiver.recv_timeout(Duration::from_millis(50)) {
Ok(command) => {
handle(command, shared, &mut output, &mut track_loaded, on_event);
previous_queue_len = output.as_ref().map_or(0, |output| output.player.len());
}
Err(RecvTimeoutError::Timeout) => {
if let Some(output) = &output {
let millis =
u64::try_from(output.player.get_pos().as_millis()).unwrap_or(u64::MAX);
shared.position_ms.store(millis, Ordering::Relaxed);
let queue_len = output.player.len();
if track_loaded && queue_len < previous_queue_len {
track_loaded = false;
on_event(Event::Finished);
}
previous_queue_len = queue_len;
}
}
Err(RecvTimeoutError::Disconnected) => return,
}
}
}
#[allow(clippy::too_many_lines, reason = "exhaustive audio command dispatcher")]
fn handle(
command: Command,
shared: &Arc<Shared>,
output: &mut Option<Output>,
track_loaded: &mut bool,
on_event: &impl Fn(Event),
) {
match command {
Command::Play { path, volume } => {
let output = match ensure_output(output) {
Ok(output) => output,
Err(error) => {
on_event(Event::Failed(format!("cannot open audio output: {error}")));
return;
}
};
let file = match File::open(&path) {
Ok(file) => file,
Err(error) => {
on_event(Event::Failed(format!(
"cannot open {}: {error}",
path.display()
)));
return;
}
};
let byte_len = file.metadata().ok().map(|metadata| metadata.len());
let mut decoder = Decoder::builder()
.with_data(file)
.with_seekable(true)
.with_gapless(true);
if let Some(byte_len) = byte_len {
decoder = decoder.with_byte_len(byte_len);
}
match decoder.build() {
Ok(source) => {
output.player.stop();
output.player.set_volume(amplitude(volume));
output.player.append(source);
output.player.play();
shared.position_ms.store(0, Ordering::Relaxed);
*track_loaded = true;
on_event(Event::Started);
}
Err(error) => on_event(Event::Failed(format!(
"cannot decode {}: {error}",
path.display()
))),
}
}
Command::PlayStream {
reader,
mime_type,
volume,
} => {
let output = match ensure_output(output) {
Ok(output) => output,
Err(error) => {
on_event(Event::Failed(format!("cannot open audio output: {error}")));
return;
}
};
match Decoder::builder()
.with_data(reader)
.with_mime_type(&mime_type)
.with_seekable(false)
.with_gapless(true)
.build()
{
Ok(source) => {
output.player.stop();
output.player.set_volume(amplitude(volume));
output.player.append(source);
output.player.play();
shared.position_ms.store(0, Ordering::Relaxed);
*track_loaded = true;
on_event(Event::Started);
}
Err(error) => on_event(Event::Failed(format!(
"cannot decode federated stream: {error}"
))),
}
}
Command::Pause => {
if let Some(output) = output {
output.player.pause();
}
}
Command::Resume => {
if let Some(output) = output {
output.player.play();
}
}
Command::Stop => {
if let Some(output) = output {
output.player.stop();
}
shared.position_ms.store(0, Ordering::Relaxed);
*track_loaded = false;
}
Command::Seek(position) => {
if let Some(output) = output
&& let Err(error) = output.player.try_seek(position)
{
on_event(Event::Failed(format!("cannot seek: {error}")));
}
}
Command::SetVolume(volume) => {
if let Some(output) = output {
output.player.set_volume(amplitude(volume));
}
}
}
}
fn ensure_output(output: &mut Option<Output>) -> Result<&Output, rodio::stream::DeviceSinkError> {
if output.is_none() {
let device = DeviceSinkBuilder::open_default_sink()?;
let player = Player::connect_new(device.mixer());
*output = Some(Output {
_device: device,
player,
});
}
Ok(output.as_ref().expect("audio output initialized above"))
}
fn amplitude(volume: f32) -> f32 {
volume.clamp(0.0, 1.0).powi(3)
}
#[cfg(test)]
mod tests {
use super::amplitude;
#[test]
fn perceptual_volume_is_clamped_and_cubic() {
assert!(amplitude(-1.0).abs() < f32::EPSILON);
assert!((amplitude(0.5) - 0.125).abs() < f32::EPSILON);
assert!((amplitude(2.0) - 1.0).abs() < f32::EPSILON);
}
}
File diff suppressed because it is too large Load Diff
+875
View File
@@ -0,0 +1,875 @@
use super::*;
impl DeviceSync {
pub(super) fn ensure_identity(&self) -> Result<Identity> {
let conn = lock(&self.conn);
if let (Some(device_id), Some(group_id), Some(name)) = (
get_meta(&conn, "device_id")?,
get_meta(&conn, "group_id")?,
get_meta(&conn, "device_name")?,
) {
return Ok(Identity {
device_id,
group_id,
name,
});
}
let seed = random_hex(32);
let digest = hash_secret(&seed);
let device_id = format!("dev_{}", &digest[..24]);
let group_id = format!("grp_{}", &hash_secret(&device_id)[..24]);
let name = format!("Furumi on {}", std::env::consts::OS);
set_meta(&conn, "device_id", &device_id)?;
set_meta(&conn, "device_secret", &seed)?;
set_meta(&conn, "group_id", &group_id)?;
set_meta(&conn, "device_name", &name)?;
set_meta(&conn, "local_seq", "0")?;
set_meta(&conn, "last_hlc_ms", "0")?;
conn.execute(
"INSERT INTO sync_devices
(device_id, name, client_version, protocol_version, trusted_at_ms, last_seen_ms)
VALUES (?1, ?2, ?3, ?4, ?5, ?5)
ON CONFLICT(device_id) DO UPDATE SET name = excluded.name,
client_version = excluded.client_version,
protocol_version = excluded.protocol_version,
trusted_at_ms = COALESCE(sync_devices.trusted_at_ms, excluded.trusted_at_ms)",
params![
device_id,
name,
CLIENT_VERSION,
DEVICE_SYNC_PROTOCOL_VERSION,
now_ms()
],
)?;
Ok(Identity {
device_id,
group_id,
name,
})
}
pub(super) fn set_group_id(&self, group_id: &str) -> Result<()> {
let conn = lock(&self.conn);
set_meta(&conn, "group_id", group_id)
}
pub(super) fn set_meta(&self, key: &str, value: &str) -> Result<()> {
set_meta(&lock(&self.conn), key, value)
}
pub(super) fn own_profile(&self, ticket: &str) -> Result<DeviceProfileWire> {
let identity = self.ensure_identity()?;
Ok(DeviceProfileWire {
device_id: identity.device_id,
name: identity.name,
client_version: CLIENT_VERSION.into(),
protocol_version: DEVICE_SYNC_PROTOCOL_VERSION,
endpoint_id: ticket_endpoint_id(ticket).unwrap_or_default(),
endpoint_ticket: ticket.into(),
revoked: false,
revoke_cutoff_seq: None,
updated_at_ms: now_ms(),
})
}
pub(super) fn active_device_count(&self) -> Result<usize> {
let count = lock(&self.conn).query_row(
"SELECT COUNT(*) FROM sync_devices
WHERE trusted_at_ms IS NOT NULL AND revoked_at_ms IS NULL",
[],
|row| row.get::<_, i64>(0),
)?;
Ok(usize::try_from(count.max(0)).unwrap_or(usize::MAX))
}
pub(super) fn status_inner(&self) -> Result<Status> {
let identity = self.ensure_identity()?;
let conn = lock(&self.conn);
let now = now_ms();
let devices = {
let mut stmt = conn.prepare(
"SELECT device_id, name, client_version, last_seen_ms
FROM sync_devices
WHERE trusted_at_ms IS NOT NULL AND revoked_at_ms IS NULL
ORDER BY name COLLATE NOCASE, device_id",
)?;
stmt.query_map([], |row| {
let id: String = row.get(0)?;
let last_seen: Option<i64> = row.get(3)?;
Ok(DeviceRow {
is_self: id == identity.device_id,
online: id == identity.device_id
|| last_seen.is_some_and(|seen| now.saturating_sub(seen) <= ONLINE_TTL_MS),
id,
name: row.get(1)?,
client_version: row.get(2)?,
revoked: false,
})
})?
.collect::<rusqlite::Result<Vec<_>>>()?
};
let pending = {
let mut stmt = conn.prepare(
"SELECT request_id, device_id, name, client_version,
requester_group_id, requester_group_active_devices
FROM sync_pending_pairing WHERE status = 'pending'
ORDER BY created_at_ms",
)?;
stmt.query_map([], |row| {
Ok(PendingPairing {
request_id: row.get(0)?,
device_id: row.get(1)?,
name: row.get(2)?,
client_version: row.get(3)?,
requester_group_id: row.get(4)?,
requester_group_active_devices: usize::try_from(row.get::<_, i64>(5)?.max(0))
.unwrap_or(usize::MAX),
})
})?
.collect::<rusqlite::Result<Vec<_>>>()?
};
Ok(Status {
this_device_id: identity.device_id,
this_device_name: identity.name,
group_id: identity.group_id,
devices,
pending,
last_sync: get_meta(&conn, "last_sync")?,
error: get_meta(&conn, "last_error")?,
})
}
pub(super) fn device_profiles(&self) -> Result<Vec<DeviceProfileWire>> {
let conn = lock(&self.conn);
let mut stmt = conn.prepare(
"SELECT device_id, name, client_version, protocol_version, endpoint_id,
endpoint_ticket, revoked_at_ms IS NOT NULL, revoke_cutoff_seq,
MAX(COALESCE(last_seen_ms, 0), COALESCE(trusted_at_ms, 0),
COALESCE(revoked_at_ms, 0))
FROM sync_devices WHERE trusted_at_ms IS NOT NULL",
)?;
Ok(stmt
.query_map([], |row| {
Ok(DeviceProfileWire {
device_id: row.get(0)?,
name: row.get(1)?,
client_version: row.get(2)?,
protocol_version: row.get::<_, u16>(3)?,
endpoint_id: row.get(4)?,
endpoint_ticket: row.get(5)?,
revoked: row.get::<_, i64>(6)? != 0,
revoke_cutoff_seq: row.get(7)?,
updated_at_ms: row.get(8)?,
})
})?
.collect::<rusqlite::Result<Vec<_>>>()?)
}
pub(super) fn apply_device_profiles(&self, profiles: &[DeviceProfileWire]) -> Result<()> {
for profile in profiles {
self.apply_device_profile(profile, false)?;
}
Ok(())
}
pub(super) fn apply_device_profile(
&self,
profile: &DeviceProfileWire,
trusted: bool,
) -> Result<()> {
if profile.device_id == self.ensure_identity()?.device_id {
return Ok(());
}
let now = now_ms();
lock(&self.conn).execute(
"INSERT INTO sync_devices
(device_id, name, client_version, protocol_version, endpoint_id,
endpoint_ticket, trusted_at_ms, last_seen_ms, revoked_at_ms,
revoke_cutoff_seq)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10)
ON CONFLICT(device_id) DO UPDATE SET
name = CASE WHEN excluded.last_seen_ms >= COALESCE(sync_devices.last_seen_ms, 0)
THEN excluded.name ELSE sync_devices.name END,
client_version = excluded.client_version,
protocol_version = excluded.protocol_version,
endpoint_id = CASE WHEN excluded.endpoint_id != '' THEN excluded.endpoint_id
ELSE sync_devices.endpoint_id END,
endpoint_ticket = CASE WHEN excluded.endpoint_ticket != '' THEN excluded.endpoint_ticket
ELSE sync_devices.endpoint_ticket END,
trusted_at_ms = COALESCE(sync_devices.trusted_at_ms, excluded.trusted_at_ms),
last_seen_ms = MAX(COALESCE(sync_devices.last_seen_ms, 0), excluded.last_seen_ms),
revoked_at_ms = CASE WHEN excluded.revoked_at_ms IS NOT NULL
THEN excluded.revoked_at_ms ELSE sync_devices.revoked_at_ms END,
revoke_cutoff_seq = COALESCE(excluded.revoke_cutoff_seq,
sync_devices.revoke_cutoff_seq)",
params![
profile.device_id,
profile.name,
profile.client_version,
profile.protocol_version,
profile.endpoint_id,
profile.endpoint_ticket,
trusted.then_some(now),
profile.updated_at_ms.max(now),
profile.revoked.then_some(profile.updated_at_ms.max(now)),
profile.revoke_cutoff_seq,
],
)?;
Ok(())
}
pub(super) fn mark_seen(&self, device_id: &str, endpoint_id: &str) -> Result<()> {
lock(&self.conn).execute(
"UPDATE sync_devices SET last_seen_ms = ?2,
endpoint_id = CASE WHEN ?3 != '' THEN ?3 ELSE endpoint_id END
WHERE device_id = ?1",
params![device_id, now_ms(), endpoint_id],
)?;
Ok(())
}
pub(super) fn record_op(&self, payload: SyncOpPayload) -> Result<()> {
let op = {
let conn = lock(&self.conn);
let identity = Self::ensure_identity_with_conn(&conn)?;
let seq = get_meta(&conn, "local_seq")?
.and_then(|value| value.parse::<i64>().ok())
.unwrap_or(0)
.saturating_add(1);
let previous_hlc = get_meta(&conn, "last_hlc_ms")?
.and_then(|value| value.parse::<i64>().ok())
.unwrap_or(0);
let hlc_ms = now_ms().max(previous_hlc.saturating_add(1));
set_meta(&conn, "local_seq", &seq.to_string())?;
set_meta(&conn, "last_hlc_ms", &hlc_ms.to_string())?;
SyncOpWire {
op_id: format!("{}:{seq}", identity.device_id),
origin_device_id: identity.device_id,
seq,
hlc_ms,
payload,
}
};
self.store_and_apply_op(&op)?;
self.request_sync();
self.notify();
Ok(())
}
pub(super) fn ensure_identity_with_conn(conn: &Connection) -> Result<Identity> {
Ok(Identity {
device_id: get_meta(conn, "device_id")?.context("missing device id")?,
group_id: get_meta(conn, "group_id")?.context("missing device group")?,
name: get_meta(conn, "device_name")?.context("missing device name")?,
})
}
pub(super) fn store_and_apply_op(&self, op: &SyncOpWire) -> Result<()> {
let inserted = lock(&self.conn).execute(
"INSERT OR IGNORE INTO sync_ops
(op_id, origin_device_id, seq, kind, payload_json, hlc_ms,
received_at_ms, tombstone)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8)",
params![
op.op_id,
op.origin_device_id,
op.seq,
payload_kind(&op.payload),
serde_json::to_string(&op.payload)?,
op.hlc_ms,
now_ms(),
i64::from(op.payload.is_tombstone()),
],
)?;
if inserted == 0 {
return Ok(());
}
lock(&self.conn).execute(
"INSERT INTO sync_vectors (device_id, max_seq) VALUES (?1, ?2)
ON CONFLICT(device_id) DO UPDATE SET max_seq = MAX(max_seq, excluded.max_seq)",
params![op.origin_device_id, op.seq],
)?;
self.apply_op(op)
}
pub(super) fn apply_ops(&self, ops: Vec<SyncOpWire>) -> Result<()> {
for op in ops {
if self.should_accept_op(&op)? {
self.store_and_apply_op(&op)?;
}
}
self.notify();
Ok(())
}
pub(super) fn should_accept_op(&self, op: &SyncOpWire) -> Result<bool> {
if op.origin_device_id == self.ensure_identity()?.device_id {
return Ok(true);
}
let conn = lock(&self.conn);
let row = conn
.query_row(
"SELECT trusted_at_ms IS NOT NULL, revoked_at_ms IS NOT NULL,
COALESCE(revoke_cutoff_seq, 0)
FROM sync_devices WHERE device_id = ?1",
[&op.origin_device_id],
|row| {
Ok((
row.get::<_, i64>(0)? != 0,
row.get::<_, i64>(1)? != 0,
row.get::<_, i64>(2)?,
))
},
)
.optional()?;
Ok(row.is_some_and(|(trusted, revoked, cutoff)| trusted && (!revoked || op.seq <= cutoff)))
}
pub(super) fn vector(&self) -> Result<BTreeMap<String, i64>> {
let conn = lock(&self.conn);
let mut stmt = conn.prepare("SELECT device_id, max_seq FROM sync_vectors")?;
Ok(stmt
.query_map([], |row| Ok((row.get(0)?, row.get(1)?)))?
.collect::<rusqlite::Result<BTreeMap<_, _>>>()?)
}
pub(super) fn ops_for_peer(&self, peer: &str) -> Result<Vec<SyncOpWire>> {
let conn = lock(&self.conn);
let mut stmt = conn.prepare(
"SELECT payload_json, op_id, origin_device_id, seq, hlc_ms
FROM sync_ops o
WHERE seq > COALESCE((SELECT max_seq FROM sync_peer_acks a
WHERE a.peer_device_id = ?1
AND a.origin_device_id = o.origin_device_id), 0)
AND (kind != 'playback_command' OR hlc_ms >= ?2)
ORDER BY received_at_ms, origin_device_id, seq LIMIT ?3",
)?;
Ok(stmt
.query_map(
params![
peer,
now_ms().saturating_sub(PLAYBACK_COMMAND_TTL_MS),
i64::try_from(MAX_OPS_PER_BATCH).unwrap_or(i64::MAX)
],
|row| {
let payload: String = row.get(0)?;
Ok(SyncOpWire {
op_id: row.get(1)?,
origin_device_id: row.get(2)?,
seq: row.get(3)?,
hlc_ms: row.get(4)?,
payload: serde_json::from_str(&payload).map_err(|error| {
rusqlite::Error::FromSqlConversionFailure(
0,
rusqlite::types::Type::Text,
Box::new(error),
)
})?,
})
},
)?
.collect::<rusqlite::Result<Vec<_>>>()?)
}
pub(super) fn note_peer_vector(
&self,
peer: &str,
vector: &BTreeMap<String, i64>,
) -> Result<()> {
let conn = lock(&self.conn);
for (origin, seq) in vector {
conn.execute(
"INSERT INTO sync_peer_acks
(peer_device_id, origin_device_id, max_seq, updated_at_ms)
VALUES (?1, ?2, ?3, ?4)
ON CONFLICT(peer_device_id, origin_device_id) DO UPDATE SET
max_seq = MAX(max_seq, excluded.max_seq),
updated_at_ms = excluded.updated_at_ms",
params![peer, origin, seq, now_ms()],
)?;
}
Ok(())
}
}
impl DeviceSync {
pub(super) fn apply_op(&self, op: &SyncOpWire) -> Result<()> {
match &op.payload {
SyncOpPayload::TrackLikeSet {
content_id,
liked,
fed,
} => self.apply_like(content_id, *liked, fed.as_ref(), op.hlc_ms, &op.op_id)?,
SyncOpPayload::PlaylistCreated { playlist_id, title }
| SyncOpPayload::PlaylistRenamed { playlist_id, title } => {
self.apply_playlist(playlist_id, title, false, op.hlc_ms, &op.op_id)?;
}
SyncOpPayload::PlaylistDeleted { playlist_id } => {
self.apply_playlist(playlist_id, "", true, op.hlc_ms, &op.op_id)?;
}
SyncOpPayload::PlaylistTrackAdded {
playlist_id,
content_id,
position,
fed,
} => self.apply_playlist_item(
playlist_id,
content_id,
true,
*position,
fed.as_ref(),
op.hlc_ms,
&op.op_id,
)?,
SyncOpPayload::PlaylistTrackRemoved {
playlist_id,
content_id,
} => self.apply_playlist_item(
playlist_id,
content_id,
false,
0,
None,
op.hlc_ms,
&op.op_id,
)?,
SyncOpPayload::DeviceProfileSet {
name,
client_version,
endpoint_ticket,
endpoint_id,
} => self.apply_device_profile(
&DeviceProfileWire {
device_id: op.origin_device_id.clone(),
name: name.clone(),
client_version: client_version.clone(),
protocol_version: DEVICE_SYNC_PROTOCOL_VERSION,
endpoint_id: endpoint_id.clone(),
endpoint_ticket: endpoint_ticket.clone(),
revoked: false,
revoke_cutoff_seq: None,
updated_at_ms: op.hlc_ms,
},
false,
)?,
SyncOpPayload::DeviceTrusted { target_device_id } => {
lock(&self.conn).execute(
"UPDATE sync_devices SET trusted_at_ms = MAX(COALESCE(trusted_at_ms, 0), ?2),
revoked_at_ms = CASE WHEN COALESCE(revoked_at_ms, 0) <= ?2
THEN NULL ELSE revoked_at_ms END
WHERE device_id = ?1",
params![target_device_id, op.hlc_ms],
)?;
}
SyncOpPayload::DeviceRevoked {
target_device_id,
target_max_seq_seen,
} => {
lock(&self.conn).execute(
"UPDATE sync_devices SET revoked_at_ms = ?2, revoked_by = ?3,
revoke_cutoff_seq = ?4 WHERE device_id = ?1",
params![
target_device_id,
op.hlc_ms,
op.origin_device_id,
target_max_seq_seen
],
)?;
}
SyncOpPayload::PlaybackCommand {
target_device_id,
command,
} => self.apply_playback_command(target_device_id, command, &op.op_id)?,
SyncOpPayload::ListenRecorded { event } => {
self.library
.apply_listen_event(event, &op.origin_device_id)?;
}
}
Ok(())
}
pub(super) fn apply_like(
&self,
content_id: &str,
liked: bool,
fed: Option<&SyncedFedTrack>,
hlc_ms: i64,
op_id: &str,
) -> Result<()> {
let Some(content_id) = music_dht::normalize_content_id(content_id) else {
return Ok(());
};
if !self.lww_wins("sync_state_likes", "content_id", &content_id, hlc_ms, op_id)? {
return Ok(());
}
lock(&self.conn).execute(
"INSERT INTO sync_state_likes (content_id, liked, hlc_ms, op_id)
VALUES (?1, ?2, ?3, ?4)
ON CONFLICT(content_id) DO UPDATE SET liked = excluded.liked,
hlc_ms = excluded.hlc_ms, op_id = excluded.op_id",
params![content_id, i64::from(liked), hlc_ms, op_id],
)?;
if let Some(track_id) = self.library.track_id_by_content_id(&content_id)? {
self.library.set_synced_like(track_id, liked, hlc_ms)?;
if liked {
self.library.remove_fed_like_by_content_id(&content_id)?;
}
} else if liked {
if let Some(fed) = fed {
self.library
.upsert_synced_fed_like(&to_library_fed(fed), hlc_ms)?;
}
} else {
self.library.remove_fed_like_by_content_id(&content_id)?;
}
self.notify_library();
Ok(())
}
pub(super) fn apply_playlist(
&self,
playlist_id: &str,
title: &str,
deleted: bool,
hlc_ms: i64,
op_id: &str,
) -> Result<()> {
if !self.lww_wins(
"sync_state_playlists",
"playlist_id",
playlist_id,
hlc_ms,
op_id,
)? {
return Ok(());
}
lock(&self.conn).execute(
"INSERT INTO sync_state_playlists
(playlist_id, title, deleted, hlc_ms, op_id)
VALUES (?1, ?2, ?3, ?4, ?5)
ON CONFLICT(playlist_id) DO UPDATE SET title = excluded.title,
deleted = excluded.deleted, hlc_ms = excluded.hlc_ms,
op_id = excluded.op_id",
params![playlist_id, title, i64::from(deleted), hlc_ms, op_id],
)?;
if deleted {
self.library.delete_playlist_by_sync_id(playlist_id)?;
} else if !title.trim().is_empty() {
self.library.upsert_synced_playlist(playlist_id, title)?;
}
self.notify_library();
Ok(())
}
#[allow(clippy::too_many_arguments)]
pub(super) fn apply_playlist_item(
&self,
playlist_id: &str,
content_id: &str,
present: bool,
position: i64,
fed: Option<&SyncedFedTrack>,
hlc_ms: i64,
op_id: &str,
) -> Result<()> {
let Some(content_id) = music_dht::normalize_content_id(content_id) else {
return Ok(());
};
let current = lock(&self.conn)
.query_row(
"SELECT hlc_ms, op_id FROM sync_state_playlist_items
WHERE playlist_id = ?1 AND content_id = ?2",
params![playlist_id, content_id],
|row| Ok((row.get::<_, i64>(0)?, row.get::<_, String>(1)?)),
)
.optional()?;
if current.is_some_and(|(current_hlc, current_op)| {
(current_hlc, current_op.as_str()) >= (hlc_ms, op_id)
}) {
return Ok(());
}
lock(&self.conn).execute(
"INSERT INTO sync_state_playlist_items
(playlist_id, content_id, present, position, hlc_ms, op_id)
VALUES (?1, ?2, ?3, ?4, ?5, ?6)
ON CONFLICT(playlist_id, content_id) DO UPDATE SET
present = excluded.present, position = excluded.position,
hlc_ms = excluded.hlc_ms, op_id = excluded.op_id",
params![
playlist_id,
content_id,
i64::from(present),
position,
hlc_ms,
op_id
],
)?;
if present {
self.library
.add_content_id_to_synced_playlist(playlist_id, &content_id)?;
if let Some(fed) = fed {
self.library.upsert_fed_playlist_track(
playlist_id,
&to_library_fed(fed),
position,
)?;
} else if let Some(fed) = self.library.fed_like_by_content_id(&content_id)? {
self.library
.upsert_fed_playlist_track(playlist_id, &fed, position)?;
}
} else {
self.library
.remove_content_id_from_synced_playlist(playlist_id, &content_id)?;
}
self.notify_library();
Ok(())
}
pub(super) fn lww_wins(
&self,
table: &str,
key_name: &str,
key: &str,
hlc_ms: i64,
op_id: &str,
) -> Result<bool> {
let sql = format!("SELECT hlc_ms, op_id FROM {table} WHERE {key_name} = ?1");
let current = lock(&self.conn)
.query_row(&sql, [key], |row| {
Ok((row.get::<_, i64>(0)?, row.get::<_, String>(1)?))
})
.optional()?;
Ok(current.is_none_or(|(current_hlc, current_op)| {
(hlc_ms, op_id) > (current_hlc, current_op.as_str())
}))
}
pub(super) fn apply_playback_command(
&self,
target_device_id: &str,
command: &PlaybackCommand,
op_id: &str,
) -> Result<()> {
if target_device_id != self.ensure_identity()?.device_id {
return Ok(());
}
let inserted = lock(&self.conn).execute(
"INSERT OR IGNORE INTO sync_playback_applied (op_id, applied_at_ms)
VALUES (?1, ?2)",
params![op_id, now_ms()],
)?;
if inserted > 0 {
let _ = self
.events
.try_send(InternalEvent::DevicePlaybackCommand(command.clone()));
}
Ok(())
}
pub(super) fn apply_playback_snapshot(&self, snapshot: PlaybackSnapshot) {
if self
.ensure_identity()
.is_ok_and(|identity| identity.device_id == snapshot.device_id)
{
return;
}
let changed = {
let mut playback = lock(&self.playback);
let changed = playback
.remote
.get(&snapshot.device_id)
.is_none_or(|current| snapshot.updated_at_ms > current.updated_at_ms);
if changed {
playback
.remote
.insert(snapshot.device_id.clone(), snapshot.clone());
}
changed
};
if changed {
let _ = self
.events
.try_send(InternalEvent::DevicePlaybackSnapshot(snapshot));
}
}
#[allow(
clippy::too_many_lines,
reason = "the snapshot is one transactional projection of all synchronized entities"
)]
pub(super) fn snapshot(&self) -> Result<SyncSnapshot> {
let conn = lock(&self.conn);
let mut snapshot = SyncSnapshot::default();
{
let mut stmt =
conn.prepare("SELECT content_id, liked, hlc_ms, op_id FROM sync_state_likes")?;
let rows = stmt.query_map([], |row| {
Ok((
row.get::<_, String>(0)?,
row.get::<_, i64>(1)? != 0,
row.get::<_, i64>(2)?,
row.get::<_, String>(3)?,
))
})?;
for row in rows {
let (content_id, liked, hlc_ms, op_id) = row?;
if liked {
snapshot.likes.push(SnapshotLike {
content_id,
hlc_ms,
op_id,
fed: None,
});
} else {
snapshot.unlikes.push(SnapshotLikeTombstone {
content_id,
hlc_ms,
op_id,
});
}
}
}
let mut playlists = BTreeMap::<String, SnapshotPlaylist>::new();
{
let mut stmt = conn.prepare(
"SELECT playlist_id, title, deleted, hlc_ms, op_id
FROM sync_state_playlists",
)?;
let rows = stmt.query_map([], |row| {
Ok((
row.get::<_, String>(0)?,
row.get::<_, String>(1)?,
row.get::<_, i64>(2)? != 0,
row.get::<_, i64>(3)?,
row.get::<_, String>(4)?,
))
})?;
for row in rows {
let (playlist_id, title, deleted, hlc_ms, op_id) = row?;
if deleted {
snapshot.deleted_playlists.push(SnapshotPlaylistTombstone {
playlist_id,
hlc_ms,
op_id,
});
} else {
playlists.insert(
playlist_id.clone(),
SnapshotPlaylist {
playlist_id,
title,
hlc_ms,
op_id,
items: Vec::new(),
},
);
}
}
}
{
let mut stmt = conn.prepare(
"SELECT playlist_id, content_id, present, position, hlc_ms, op_id
FROM sync_state_playlist_items",
)?;
let rows = stmt.query_map([], |row| {
Ok((
row.get::<_, String>(0)?,
row.get::<_, String>(1)?,
row.get::<_, i64>(2)? != 0,
row.get::<_, i64>(3)?,
row.get::<_, i64>(4)?,
row.get::<_, String>(5)?,
))
})?;
for row in rows {
let (playlist_id, content_id, present, position, hlc_ms, op_id) = row?;
if present {
if let Some(playlist) = playlists.get_mut(&playlist_id) {
playlist.items.push(SnapshotPlaylistItem {
content_id,
position,
hlc_ms,
op_id,
fed: None,
});
}
} else {
snapshot
.removed_playlist_items
.push(SnapshotPlaylistItemTombstone {
playlist_id,
content_id,
hlc_ms,
op_id,
});
}
}
}
snapshot.playlists = playlists.into_values().collect();
Ok(snapshot)
}
pub(super) fn apply_snapshot(&self, snapshot: SyncSnapshot) -> Result<()> {
for like in snapshot.likes {
self.apply_like(
&like.content_id,
true,
like.fed.as_ref(),
like.hlc_ms,
&like.op_id,
)?;
}
for like in snapshot.unlikes {
self.apply_like(&like.content_id, false, None, like.hlc_ms, &like.op_id)?;
}
for playlist in snapshot.playlists {
self.apply_playlist(
&playlist.playlist_id,
&playlist.title,
false,
playlist.hlc_ms,
&playlist.op_id,
)?;
for item in playlist.items {
self.apply_playlist_item(
&playlist.playlist_id,
&item.content_id,
true,
item.position,
item.fed.as_ref(),
item.hlc_ms,
&item.op_id,
)?;
}
}
for playlist in snapshot.deleted_playlists {
self.apply_playlist(
&playlist.playlist_id,
"",
true,
playlist.hlc_ms,
&playlist.op_id,
)?;
}
for item in snapshot.removed_playlist_items {
self.apply_playlist_item(
&item.playlist_id,
&item.content_id,
false,
0,
None,
item.hlc_ms,
&item.op_id,
)?;
}
self.notify();
Ok(())
}
pub(super) fn notify(&self) {
let _ = self.events.try_send(InternalEvent::DevicesChanged);
}
pub(super) fn notify_library(&self) {
let _ = self.events.try_send(InternalEvent::DeviceLibraryChanged);
}
}
+858
View File
@@ -0,0 +1,858 @@
//! Federation lifecycle, DHT search and catalog artwork fetching.
use std::collections::{HashMap, HashSet};
use std::hash::{DefaultHasher, Hash, Hasher};
use std::path::PathBuf;
use std::sync::Arc;
use anyhow::{Context, Result};
use furumi_backend_api::{FederationDebugSnapshot, SearchResults, SearchStats};
use furumi_domain::{
Artist, ArtistKey, ArtistRef, Artwork, AudioSource, CatalogSource, ContentId, Release,
ReleaseKey, Track, TrackKey,
};
use music_dht::catalog::{
CATALOG_ALPN, CatalogArtist, CatalogImageHeader, CatalogRequest, CatalogResponse,
};
use music_dht::{
EndpointId, ItemKind, ItemSpec, LibraryItem, MusicDhtConfig, MusicDhtService, NetworkId,
RendezvousConfig,
};
use serde::{Deserialize, Serialize};
use tokio::io::{AsyncRead, AsyncReadExt as _};
pub const AUDIO_ALPN: &[u8] = b"furumi-fd/audio/1";
pub const AUDIO_PROTOCOL_VERSION: u16 = 1;
const STREAM_BUFFER: u64 = 2 * 1024 * 1024;
#[derive(Serialize)]
struct AudioRequest {
item_id: String,
offset: u64,
want_cover: bool,
metadata_only: bool,
}
#[derive(Deserialize)]
struct AudioHeader {
ok: bool,
#[serde(default)]
error: Option<String>,
#[serde(default)]
mime_type: String,
#[serde(default)]
total_size: u64,
#[serde(default)]
cover_size: u64,
#[serde(default)]
artist_image_size: u64,
#[serde(default)]
metadata: Option<TrackMetadata>,
}
#[derive(Clone, Default, Deserialize)]
pub struct TrackMetadata {
#[serde(default)]
pub title: String,
#[serde(default)]
pub artists: Vec<String>,
#[serde(default)]
pub featured_artists: Vec<String>,
#[serde(default)]
pub album_artists: Vec<String>,
#[serde(default)]
pub release_title: String,
#[serde(default)]
pub release_type: Option<String>,
pub year: Option<i32>,
pub track_number: Option<i32>,
pub disc_number: Option<i32>,
pub duration_seconds: Option<f64>,
pub audio_format: Option<String>,
pub audio_bitrate: Option<i32>,
pub audio_sample_rate: Option<i32>,
pub audio_bit_depth: Option<i32>,
}
pub enum StreamEvent {
Ready(crate::streaming::GrowingFileReader, String),
Complete(PathBuf, Box<Option<TrackMetadata>>),
Failed(String),
}
const MAX_IMAGE_BYTES: u64 = 16 * 1024 * 1024;
const IMAGE_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(4);
pub struct Client {
service: Arc<MusicDhtService>,
media_dir: PathBuf,
}
impl Client {
pub async fn start(data_dir: PathBuf, media_dir: PathBuf, network: &str) -> Result<Arc<Self>> {
tokio::fs::create_dir_all(&data_dir).await?;
tokio::fs::create_dir_all(&media_dir).await?;
let config = MusicDhtConfig::builder()
.data_dir(data_dir)
.network_id(NetworkId::from_name(network))
.rendezvous(RendezvousConfig::default())
.stream_protocol(CATALOG_ALPN)
.stream_protocol(AUDIO_ALPN)
.stream_protocol(music_dht::device_sync::SYNC_ALPN_V1)
.stream_protocol(music_dht::device_sync::SYNC_ALPN_V2)
.build()
.context("invalid federation configuration")?;
let (service, mut events) = MusicDhtService::start(config)
.await
.context("starting federation node")?;
tokio::spawn(async move { while events.recv().await.is_some() {} });
Ok(Arc::new(Self {
service: Arc::new(service),
media_dir,
}))
}
pub fn service(&self) -> Arc<MusicDhtService> {
Arc::clone(&self.service)
}
pub async fn debug_snapshot(&self) -> FederationDebugSnapshot {
let stored_dht_records = self.service.dht_record_count().await.ok();
let published_items = self
.service
.list_local_items()
.await
.map_or(0, |items| items.len());
FederationDebugSnapshot {
running: true,
endpoint_id: self.service.endpoint_id().to_string(),
dht_node_id: self.service.node_id().to_string(),
connected_peers: self.service.connected_peers().len(),
known_contacts: self.service.known_peers().len(),
stored_dht_records,
published_items,
error: None,
}
}
/// Resolves album artwork for a queue item independently of the screen
/// from which the track was enqueued.
pub async fn artwork_for_track(&self, track: &Track) -> Option<PathBuf> {
let peer_id = match &track.audio_source {
AudioSource::Federation { peer_id, .. } if !peer_id.is_empty() => peer_id.as_str(),
_ => track.key.federation_id()?.0,
};
let owner = peer_id.parse::<EndpointId>().ok()?;
let artist = track
.artists
.first()
.map(|artist| artist.name.as_str())
.or_else(|| (!track.artist.is_empty()).then_some(track.artist.as_str()))?;
if track.release.is_empty() {
return None;
}
self.cached_image(
owner,
artist,
Some(&track.release),
&format!("release-{peer_id}-{artist}-{}", track.release),
)
.await
}
pub async fn search(&self, query: &str) -> Result<(SearchResults, SearchStats)> {
let outcome = self.service.search_network(query).await?;
let own = self.service.endpoint_id();
let mut results = convert_items(&outcome.network_results, own);
self.fetch_artwork(&mut results).await;
let stats = SearchStats {
tracks: results.tracks.len(),
artists: results.artists.len(),
peers_queried: outcome.queried_nodes,
duration_ms: u64::try_from(outcome.duration.as_millis()).unwrap_or(u64::MAX),
};
Ok((results, stats))
}
/// Resolves a portable queue entry when another connected device only
/// knows its stable audio content id.
pub async fn track_by_content_id(&self, content_id: &str) -> Result<Track> {
let outcome = self.service.search_content_id(content_id).await?;
let own = self.service.endpoint_id();
let items = outcome
.local_results
.into_iter()
.chain(outcome.network_results)
.collect::<Vec<_>>();
convert_items(&items, own)
.tracks
.into_iter()
.next()
.context("no federation peer currently publishes this track")
}
pub async fn publish(&self, specs: Vec<ItemSpec>) -> Result<()> {
self.service.sync_library(specs).await?;
Ok(())
}
pub async fn artist_card(&self, name: &str) -> Result<(SearchResults, SearchStats)> {
let started = std::time::Instant::now();
let outcome = self.service.search_network(name).await?;
let own = self.service.endpoint_id();
let normalized = music_dht::normalize_name(name);
let owners = outcome
.network_results
.iter()
.filter(|item| {
item.owner != own
&& ((item.kind == ItemKind::Artist && item.normalized_name == normalized)
|| item
.artist_names
.iter()
.chain(item.featured_artist_names.iter())
.any(|artist| music_dht::normalize_name(artist) == normalized))
})
.map(|item| item.owner)
.collect::<HashSet<_>>();
let mut catalogs = Vec::new();
for owner in owners {
if let Ok(Ok(catalog)) = tokio::time::timeout(
std::time::Duration::from_secs(5),
fetch_catalog(&self.service, owner, name),
)
.await
{
catalogs.push((owner.to_string(), catalog));
}
}
let mut result = card_results(name, catalogs);
self.fetch_artwork(&mut result).await;
let stats = SearchStats {
tracks: result.tracks.len(),
artists: result.artists.len(),
peers_queried: outcome.queried_nodes,
duration_ms: u64::try_from(started.elapsed().as_millis()).unwrap_or(u64::MAX),
};
Ok((result, stats))
}
pub async fn stream_track(
&self,
peer: &str,
item_id: &str,
directory: &std::path::Path,
stem: &str,
events: tokio::sync::mpsc::Sender<StreamEvent>,
) -> Result<()> {
use tokio::io::AsyncWriteExt as _;
let owner: EndpointId = peer.parse().context("invalid peer id")?;
tokio::fs::create_dir_all(directory).await?;
let mut stream = self.service.open_stream(owner, AUDIO_ALPN).await?;
let mut request = serde_json::to_vec(&AudioRequest {
item_id: item_id.into(),
offset: 0,
want_cover: false,
metadata_only: false,
})?;
request.push(b'\n');
stream.send.write_all(&request).await?;
stream.send.finish()?;
let header: AudioHeader = serde_json::from_slice(&read_line(&mut stream.recv).await?)?;
anyhow::ensure!(
header.ok,
"{}",
header.error.unwrap_or_else(|| "peer refused audio".into())
);
anyhow::ensure!(
header.cover_size == 0 && header.artist_image_size == 0,
"unexpected image data in audio stream"
);
let extension = match header.mime_type.as_str() {
"audio/mpeg" => "mp3",
"audio/flac" | "audio/x-flac" => "flac",
"audio/ogg" => "ogg",
"audio/opus" => "opus",
"audio/wav" | "audio/x-wav" => "wav",
"audio/mp4" | "audio/x-m4a" => "m4a",
"audio/aac" => "aac",
_ => "bin",
};
let final_path = directory.join(format!("{stem}.{extension}"));
let part_path = directory.join(format!(".{stem}.{extension}.part"));
let mut file = tokio::fs::File::create(&part_path).await?;
let (reader, writer) = crate::streaming::growing_file(&part_path)?;
let mut reader = Some(reader);
let mut started = false;
let mut received = 0u64;
let threshold = header.total_size.clamp(1, STREAM_BUFFER);
let mut chunk = vec![0; 64 * 1024];
while let Some(n) = stream.recv.read(&mut chunk).await? {
file.write_all(&chunk[..n]).await?;
received += n as u64;
writer.add_available(n as u64);
if !started
&& received >= threshold
&& let Some(reader) = reader.take()
{
let _ = events
.send(StreamEvent::Ready(reader, header.mime_type.clone()))
.await;
started = true;
}
}
if !started
&& received > 0
&& let Some(reader) = reader.take()
{
let _ = events
.send(StreamEvent::Ready(reader, header.mime_type.clone()))
.await;
}
writer.finish();
file.flush().await?;
drop(file);
anyhow::ensure!(
header.total_size == 0 || received == header.total_size,
"incomplete audio download: {received}/{}",
header.total_size
);
tokio::fs::rename(&part_path, &final_path).await?;
let _ = events
.send(StreamEvent::Complete(final_path, Box::new(header.metadata)))
.await;
Ok(())
}
async fn fetch_artwork(&self, results: &mut SearchResults) {
for artist in &mut results.artists {
let CatalogSource::Federation { peer_id } = &artist.source else {
continue;
};
let Ok(owner) = peer_id.parse::<EndpointId>() else {
continue;
};
if let Some(path) = self
.cached_image(
owner,
&artist.name,
None,
&format!("artist-{peer_id}-{}", artist.name),
)
.await
{
artist.artwork.uri = Some(path.to_string_lossy().into_owned());
}
}
for release in &mut results.releases {
let CatalogSource::Federation { peer_id } = &release.source else {
continue;
};
let Some(artist) = release.artists.first() else {
continue;
};
let Ok(owner) = peer_id.parse::<EndpointId>() else {
continue;
};
if let Some(path) = self
.cached_image(
owner,
&artist.name,
Some(&release.title),
&format!("release-{peer_id}-{}-{}", artist.name, release.title),
)
.await
{
let uri = path.to_string_lossy().into_owned();
release.artwork.uri = Some(uri.clone());
for track in &mut release.tracks {
if track.cover_uri.is_none() {
track.cover_uri = Some(uri.clone());
}
}
for track in &mut results.tracks {
if track.release == release.title
&& track
.artists
.first()
.is_some_and(|candidate| candidate.name == artist.name)
{
track.cover_uri = Some(uri.clone());
}
}
}
}
for track in &mut results.tracks {
if track.cover_uri.is_some() || track.release.is_empty() {
continue;
}
let AudioSource::Federation { peer_id, .. } = &track.audio_source else {
continue;
};
let Some(artist) = track.artists.first() else {
continue;
};
let Ok(owner) = peer_id.parse::<EndpointId>() else {
continue;
};
if let Some(path) = self
.cached_image(
owner,
&artist.name,
Some(&track.release),
&format!("release-{peer_id}-{}-{}", artist.name, track.release),
)
.await
{
track.cover_uri = Some(path.to_string_lossy().into_owned());
}
}
}
async fn cached_image(
&self,
owner: EndpointId,
artist: &str,
release: Option<&str>,
cache_key: &str,
) -> Option<PathBuf> {
let mut hasher = DefaultHasher::new();
cache_key.hash(&mut hasher);
let base = self.media_dir.join(format!("{:016x}", hasher.finish()));
for extension in ["jpg", "png", "webp", "gif", "bmp"] {
let path = base.with_extension(extension);
if path.is_file() {
return Some(path);
}
}
let fetched = tokio::time::timeout(
IMAGE_TIMEOUT,
fetch_image(&self.service, owner, artist, release),
)
.await
.ok()?
.ok()?;
let (bytes, extension) = fetched?;
let path = base.with_extension(extension);
tokio::fs::write(&path, bytes).await.ok()?;
Some(path)
}
}
async fn fetch_catalog(
service: &MusicDhtService,
owner: EndpointId,
artist: &str,
) -> Result<CatalogArtist> {
let mut stream = service.open_stream(owner, CATALOG_ALPN).await?;
let mut request = serde_json::to_vec(&CatalogRequest {
artist: artist.into(),
..CatalogRequest::default()
})?;
request.push(b'\n');
stream.send.write_all(&request).await?;
stream.send.finish()?;
let mut payload = Vec::new();
stream
.recv
.take(4 * 1024 * 1024 + 1)
.read_to_end(&mut payload)
.await?;
anyhow::ensure!(
payload.len() <= 4 * 1024 * 1024,
"catalog response is too large"
);
let response: CatalogResponse = serde_json::from_slice(&payload)?;
anyhow::ensure!(
response.ok,
"{}",
response.error.unwrap_or_else(|| "catalog rejected".into())
);
response.artist.context("empty artist catalog")
}
#[allow(
clippy::too_many_lines,
reason = "wire catalog conversion is clearest as one pass"
)]
fn card_results(name: &str, catalogs: Vec<(String, CatalogArtist)>) -> SearchResults {
let mut results = SearchResults::default();
if let Some((peer, _)) = catalogs.first() {
results.artists.push(Artist {
key: ArtistKey::Federation {
peer_id: peer.clone(),
id: music_dht::normalize_name(name),
},
source: CatalogSource::Federation {
peer_id: peer.clone(),
},
name: name.into(),
artwork: Artwork::default(),
release_count: 0,
track_count: 0,
});
}
let mut release_slots = HashMap::<String, usize>::new();
for (peer, catalog) in catalogs {
let artist_key = ArtistKey::Federation {
peer_id: peer.clone(),
id: music_dht::normalize_name(name),
};
for remote in catalog.releases {
let normalized = music_dht::normalize_name(&remote.title);
let slot = *release_slots.entry(normalized.clone()).or_insert_with(|| {
results.releases.push(Release {
key: ReleaseKey::Federation {
peer_id: peer.clone(),
id: format!("name:{normalized}"),
},
source: CatalogSource::Federation {
peer_id: peer.clone(),
},
title: remote.title.clone(),
artists: vec![ArtistRef {
key: artist_key.clone(),
name: name.into(),
}],
featured_artists: Vec::new(),
release_type: remote.release_type.clone(),
year: remote.year,
artwork: Artwork::default(),
tracks: Vec::new(),
});
results.releases.len() - 1
});
let release = &mut results.releases[slot];
for item in remote.tracks {
let duplicate = release.tracks.iter().any(|track| {
music_dht::normalize_name(&track.title)
== music_dht::normalize_name(&item.title)
});
if duplicate || item.item_id.is_empty() {
continue;
}
let content_id = item
.content_id
.as_deref()
.and_then(|id| ContentId::parse(id).ok());
let artists = artist_refs(&peer, &item.artists);
let featured = artist_refs(&peer, &item.featured_artists);
let track = Track {
key: TrackKey::federation(
peer.clone(),
item.item_id.clone(),
content_id.clone(),
),
title: item.title,
artist: artist_line(&item.artists, &item.featured_artists),
artists,
featured_artists: featured,
release: release.title.clone(),
release_id: release.key.clone(),
duration_seconds: item.duration_seconds.unwrap_or_default(),
track_number: item
.track_number
.and_then(|value| u32::try_from(value).ok()),
disc_number: item.disc_number.and_then(|value| u32::try_from(value).ok()),
cover_uri: release.artwork.uri.clone(),
audio_format: None,
audio_bitrate_kbps: None,
audio_sample_rate_hz: None,
audio_bit_depth: None,
file_size_bytes: None,
liked: false,
audio_source: AudioSource::Federation {
peer_id: peer.clone(),
content_id: content_id.unwrap_or_else(|| {
ContentId::parse(format!("b3:{}", "0".repeat(64)))
.expect("valid fallback")
}),
},
};
release.tracks.push(track.clone());
results.tracks.push(track);
}
}
for appearance in catalog.appears_on {
if appearance.track.item_id.is_empty() {
continue;
}
let content_id = appearance
.track
.content_id
.as_deref()
.and_then(|id| ContentId::parse(id).ok());
let release_key = ReleaseKey::Federation {
peer_id: peer.clone(),
id: format!(
"name:{}",
music_dht::normalize_name(&appearance.release_title)
),
};
results.tracks.push(Track {
key: TrackKey::federation(
peer.clone(),
appearance.track.item_id.clone(),
content_id.clone(),
),
title: appearance.track.title,
artist: artist_line(
&appearance.track.artists,
&appearance.track.featured_artists,
),
artists: artist_refs(&peer, &appearance.track.artists),
featured_artists: artist_refs(&peer, &appearance.track.featured_artists),
release: appearance.release_title,
release_id: release_key,
duration_seconds: appearance.track.duration_seconds.unwrap_or_default(),
track_number: appearance
.track
.track_number
.and_then(|value| u32::try_from(value).ok()),
disc_number: appearance
.track
.disc_number
.and_then(|value| u32::try_from(value).ok()),
cover_uri: None,
audio_format: None,
audio_bitrate_kbps: None,
audio_sample_rate_hz: None,
audio_bit_depth: None,
file_size_bytes: None,
liked: false,
audio_source: AudioSource::Federation {
peer_id: peer.clone(),
content_id: content_id.unwrap_or_else(|| {
ContentId::parse(format!("b3:{}", "0".repeat(64))).expect("valid fallback")
}),
},
});
}
}
for release in &mut results.releases {
populate_release_contributors(release);
}
if let Some(artist) = results.artists.first_mut() {
artist.release_count = results.releases.len();
artist.track_count = results.tracks.len();
}
results
}
fn populate_release_contributors(release: &mut Release) {
let mut known = release
.artists
.iter()
.chain(release.featured_artists.iter())
.map(|artist| music_dht::normalize_name(&artist.name))
.collect::<HashSet<_>>();
for artist in release
.tracks
.iter()
.flat_map(|track| track.artists.iter().chain(track.featured_artists.iter()))
{
if known.insert(music_dht::normalize_name(&artist.name)) {
release.featured_artists.push(artist.clone());
}
}
}
fn convert_items(items: &[LibraryItem], own: EndpointId) -> SearchResults {
let mut results = SearchResults::default();
let mut artists = HashMap::<(String, String), Artist>::new();
let mut releases = HashSet::<(String, String)>::new();
let mut tracks = HashSet::<(String, String)>::new();
for item in items.iter().filter(|item| item.owner != own) {
let peer = item.owner.to_string();
let source = CatalogSource::Federation {
peer_id: peer.clone(),
};
let refs = artist_refs(&peer, &item.artist_names);
let featured = artist_refs(&peer, &item.featured_artist_names);
for name in item
.artist_names
.iter()
.chain(item.featured_artist_names.iter())
{
note_artist(&mut artists, &peer, name);
}
match item.kind {
ItemKind::Artist => {
note_artist(&mut artists, &peer, &item.name);
}
ItemKind::Release => {
if releases.insert((peer.clone(), item.id.to_string())) {
results.releases.push(Release {
key: ReleaseKey::Federation {
peer_id: peer.clone(),
id: item.id.to_string(),
},
source,
title: item.name.clone(),
artists: refs,
featured_artists: Vec::new(),
release_type: item
.release_type
.clone()
.unwrap_or_else(|| "release".into()),
year: item.year,
artwork: Artwork::default(),
tracks: Vec::new(),
});
}
}
ItemKind::Track => {
if !tracks.insert((peer.clone(), item.id.to_string())) {
continue;
}
let Some(content_id) = item
.content_id
.as_deref()
.and_then(|value| ContentId::parse(value).ok())
else {
continue;
};
let release = item.release_title.clone().unwrap_or_default();
results.tracks.push(Track {
key: TrackKey::federation(
peer.clone(),
item.id.to_string(),
Some(content_id.clone()),
),
title: item.name.clone(),
artist: artist_line(&item.artist_names, &item.featured_artist_names),
artists: refs.clone(),
featured_artists: featured,
release: release.clone(),
release_id: ReleaseKey::Federation {
peer_id: peer.clone(),
id: format!("name:{}", music_dht::normalize_name(&release)),
},
duration_seconds: item.duration_seconds.unwrap_or_default(),
track_number: item
.track_number
.and_then(|value| u32::try_from(value).ok()),
disc_number: item.disc_number.and_then(|value| u32::try_from(value).ok()),
cover_uri: None,
audio_format: None,
audio_bitrate_kbps: None,
audio_sample_rate_hz: None,
audio_bit_depth: None,
file_size_bytes: None,
liked: false,
audio_source: AudioSource::Federation {
peer_id: peer,
content_id,
},
});
}
}
}
results.artists = artists.into_values().collect();
results
.artists
.sort_by(|left, right| left.name.cmp(&right.name));
results
}
fn note_artist(artists: &mut HashMap<(String, String), Artist>, peer: &str, name: &str) {
let normalized = music_dht::normalize_name(name);
if normalized.is_empty() {
return;
}
artists
.entry((peer.to_owned(), normalized.clone()))
.or_insert_with(|| Artist {
key: ArtistKey::Federation {
peer_id: peer.to_owned(),
id: normalized,
},
source: CatalogSource::Federation {
peer_id: peer.to_owned(),
},
name: name.to_owned(),
artwork: Artwork::default(),
release_count: 0,
track_count: 0,
});
}
fn artist_refs(peer: &str, names: &[String]) -> Vec<ArtistRef> {
names
.iter()
.map(|name| ArtistRef {
key: ArtistKey::Federation {
peer_id: peer.to_owned(),
id: music_dht::normalize_name(name),
},
name: name.clone(),
})
.collect()
}
fn artist_line(main: &[String], featured: &[String]) -> String {
let mut line = main.join(", ");
if !featured.is_empty() {
if !line.is_empty() {
line.push_str(" feat. ");
}
line.push_str(&featured.join(", "));
}
line
}
async fn fetch_image(
service: &MusicDhtService,
owner: EndpointId,
artist: &str,
release: Option<&str>,
) -> Result<Option<(Vec<u8>, &'static str)>> {
let mut stream = service.open_stream(owner, CATALOG_ALPN).await?;
let request = CatalogRequest {
artist: artist.to_owned(),
want: Some(
if release.is_some() {
"release_cover"
} else {
"artist_image"
}
.into(),
),
release: release.map(str::to_owned),
cursor: None,
limit: None,
};
let mut line = serde_json::to_vec(&request)?;
line.push(b'\n');
stream.send.write_all(&line).await?;
stream.send.finish()?;
let header: CatalogImageHeader = serde_json::from_slice(&read_line(&mut stream.recv).await?)?;
if !header.ok || header.size == 0 {
return Ok(None);
}
anyhow::ensure!(
header.size <= MAX_IMAGE_BYTES,
"federated image exceeds size limit"
);
let image_size =
usize::try_from(header.size).context("image is too large for this platform")?;
let mut bytes = vec![0; image_size];
stream.recv.read_exact(&mut bytes).await?;
let extension = match header.mime_type.as_str() {
"image/png" => "png",
"image/webp" => "webp",
"image/gif" => "gif",
"image/bmp" => "bmp",
_ => "jpg",
};
Ok(Some((bytes, extension)))
}
async fn read_line(reader: &mut (impl AsyncRead + Unpin)) -> Result<Vec<u8>> {
let mut line = Vec::new();
loop {
let byte = reader.read_u8().await?;
if byte == b'\n' {
return Ok(line);
}
anyhow::ensure!(line.len() < 64 * 1024, "catalog header is too large");
line.push(byte);
}
}
File diff suppressed because it is too large Load Diff
+153
View File
@@ -0,0 +1,153 @@
use std::fs;
use std::path::Path;
use furumi_backend_api::SettingsSnapshot;
use rusqlite::{Connection, OptionalExtension, params};
const MIGRATIONS: &[(i64, &str)] = &[
(
1,
r"
CREATE TABLE app_settings (
singleton_id INTEGER PRIMARY KEY CHECK (singleton_id = 1),
network_id TEXT NOT NULL,
library_path TEXT NOT NULL,
federation_enabled INTEGER NOT NULL CHECK (federation_enabled IN (0, 1)),
language TEXT NOT NULL
);
INSERT INTO app_settings (
singleton_id, network_id, library_path, federation_enabled, language
) VALUES (1, 'furumi', '~/Music/Furumi', 1, 'English');
",
),
(
2,
"ALTER TABLE app_settings ADD COLUMN save_federated_on_listen INTEGER NOT NULL DEFAULT 1 CHECK (save_federated_on_listen IN (0, 1));",
),
(
3,
"ALTER TABLE app_settings ADD COLUMN device_name TEXT NOT NULL DEFAULT '';",
),
];
pub struct SettingsStore {
connection: Connection,
}
impl SettingsStore {
pub fn open(path: &Path) -> rusqlite::Result<Self> {
if let Some(parent) = path.parent() {
fs::create_dir_all(parent)
.map_err(|error| rusqlite::Error::ToSqlConversionFailure(Box::new(error)))?;
}
let connection = Connection::open(path)?;
let mut store = Self { connection };
store.migrate()?;
Ok(store)
}
#[cfg(test)]
fn in_memory() -> rusqlite::Result<Self> {
let connection = Connection::open_in_memory()?;
let mut store = Self { connection };
store.migrate()?;
Ok(store)
}
pub fn load(&self) -> rusqlite::Result<SettingsSnapshot> {
self.connection.query_row(
"SELECT network_id, library_path, federation_enabled, language,
save_federated_on_listen, device_name
FROM app_settings WHERE singleton_id = 1",
[],
|row| {
Ok(SettingsSnapshot {
network_id: row.get(0)?,
library_path: row.get(1)?,
federation_enabled: row.get::<_, i64>(2)? != 0,
language: row.get(3)?,
save_federated_on_listen: row.get::<_, i64>(4)? != 0,
device_name: row.get(5)?,
})
},
)
}
pub fn save(&self, settings: &SettingsSnapshot) -> rusqlite::Result<()> {
self.connection.execute(
"UPDATE app_settings
SET network_id = ?1,
library_path = ?2,
federation_enabled = ?3,
language = ?4,
save_federated_on_listen = ?5,
device_name = ?6
WHERE singleton_id = 1",
params![
settings.network_id,
settings.library_path,
i64::from(settings.federation_enabled),
settings.language,
i64::from(settings.save_federated_on_listen),
settings.device_name,
],
)?;
Ok(())
}
fn migrate(&mut self) -> rusqlite::Result<()> {
self.connection.execute_batch(
"PRAGMA foreign_keys = ON;
CREATE TABLE IF NOT EXISTS schema_migrations (
version INTEGER PRIMARY KEY,
applied_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP
);",
)?;
for &(version, sql) in MIGRATIONS {
let applied = self
.connection
.query_row(
"SELECT 1 FROM schema_migrations WHERE version = ?1",
[version],
|_| Ok(()),
)
.optional()?
.is_some();
if applied {
continue;
}
let transaction = self.connection.transaction()?;
transaction.execute_batch(sql)?;
transaction.execute(
"INSERT INTO schema_migrations (version) VALUES (?1)",
[version],
)?;
transaction.commit()?;
}
Ok(())
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn migration_creates_defaults_and_settings_round_trip() {
let store = SettingsStore::in_memory().unwrap();
let mut settings = store.load().unwrap();
assert_eq!(settings.network_id, "furumi");
assert!(settings.federation_enabled);
assert!(settings.save_federated_on_listen);
assert!(settings.device_name.is_empty());
settings.network_id = "friends".into();
settings.device_name = "Studio Mac".into();
settings.library_path = "/music/library".into();
settings.federation_enabled = false;
store.save(&settings).unwrap();
assert_eq!(store.load().unwrap(), settings);
}
}
+97
View File
@@ -0,0 +1,97 @@
use std::fs::File;
use std::io::{self, Read, Seek, SeekFrom};
use std::path::Path;
use std::sync::{Arc, Condvar, Mutex};
#[derive(Default)]
struct State {
available: u64,
complete: bool,
}
#[derive(Default)]
struct Shared {
state: Mutex<State>,
changed: Condvar,
}
pub struct GrowingFileReader {
file: File,
pos: u64,
shared: Arc<Shared>,
}
pub struct GrowingFileWriter {
shared: Arc<Shared>,
}
pub fn growing_file(path: &Path) -> io::Result<(GrowingFileReader, GrowingFileWriter)> {
let shared = Arc::new(Shared::default());
Ok((
GrowingFileReader {
file: File::open(path)?,
pos: 0,
shared: Arc::clone(&shared),
},
GrowingFileWriter { shared },
))
}
impl GrowingFileWriter {
pub fn add_available(&self, bytes: u64) {
let mut state = self
.shared
.state
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
state.available = state.available.saturating_add(bytes);
self.shared.changed.notify_all();
}
pub fn finish(&self) {
let mut state = self
.shared
.state
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
state.complete = true;
self.shared.changed.notify_all();
}
}
impl Drop for GrowingFileWriter {
fn drop(&mut self) {
self.finish();
}
}
impl Read for GrowingFileReader {
fn read(&mut self, buf: &mut [u8]) -> io::Result<usize> {
loop {
let mut state = self
.shared
.state
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
while self.pos >= state.available && !state.complete {
state = self
.shared
.changed
.wait(state)
.unwrap_or_else(std::sync::PoisonError::into_inner);
}
if self.pos >= state.available {
return Ok(0);
}
let n = usize::try_from((state.available - self.pos).min(buf.len() as u64))
.unwrap_or(buf.len());
drop(state);
let read = self.file.read(&mut buf[..n])?;
self.pos += read as u64;
if read > 0 {
return Ok(read);
}
}
}
}
impl Seek for GrowingFileReader {
fn seek(&mut self, _: SeekFrom) -> io::Result<u64> {
Err(io::Error::new(
io::ErrorKind::Unsupported,
"federated stream is not seekable while downloading",
))
}
}
+955
View File
@@ -0,0 +1,955 @@
use super::{
Artist, ArtistId, ArtistKey, ArtistRef, Artwork, AudioSource, BuildInfoSnapshot,
CONTROL_POSITION_ACK_TOLERANCE_SECONDS, CatalogSource, ContentId, DevicePlaybackRole, Duration,
HashMap, HashSet, InternalEvent, LibrarySnapshot, LocalTrackId, PlaybackStatus,
PlaylistSnapshot, Release, ReleaseId, ReleaseKey, SearchResults, SettingsSnapshot,
SettingsStore, Track, TrackKey, VersionEntrySnapshot, federation, mpsc, std_mpsc, thread,
};
pub(super) fn expand_tilde(value: &str) -> std::path::PathBuf {
if value == "~" {
directories::UserDirs::new().map_or_else(|| value.into(), |dirs| dirs.home_dir().into())
} else if let Some(rest) = value.strip_prefix("~/") {
directories::UserDirs::new().map_or_else(|| value.into(), |dirs| dirs.home_dir().join(rest))
} else {
value.into()
}
}
pub(super) fn normalize_device_name(value: &str) -> String {
let value = value.trim();
if value.is_empty() {
"furumi".into()
} else {
value.into()
}
}
pub(super) fn selected_track_position(tracks: &[Track], selected: &TrackKey) -> usize {
tracks
.iter()
.position(|track| track.key.matches(selected))
.unwrap_or(0)
}
pub(super) fn runtime_build_info() -> BuildInfoSnapshot {
use music_dht::capabilities::{
CATALOG_ID, CapabilityManifest, DEVICE_SYNC_ID, FEDERATION_NET_ID, MUSIC_DHT_ID,
RENDEZVOUS_ID, TICKET_ID,
};
let manifest = CapabilityManifest::frid("furumi-desktop", env!("CARGO_PKG_VERSION"));
let protocol = |name: &str, id: &str| VersionEntrySnapshot {
name: name.into(),
version: manifest
.protocols
.get(id)
.map_or_else(|| "unknown".into(), u16::to_string),
};
BuildInfoSnapshot {
software: vec![
VersionEntrySnapshot {
name: "furumi-desktop".into(),
version: env!("CARGO_PKG_VERSION").into(),
},
VersionEntrySnapshot {
name: "furumi-library".into(),
version: env!("FURUMI_LIBRARY_VERSION").into(),
},
VersionEntrySnapshot {
name: "music-dht".into(),
version: env!("FURUMI_MUSIC_DHT_VERSION").into(),
},
VersionEntrySnapshot {
name: "federation-net".into(),
version: env!("FURUMI_FEDERATION_NET_VERSION").into(),
},
],
protocols: vec![
protocol("Federation transport", FEDERATION_NET_ID),
protocol("Peer ticket", TICKET_ID),
protocol("Rendezvous", RENDEZVOUS_ID),
protocol("Music DHT", MUSIC_DHT_ID),
protocol("Catalog", CATALOG_ID),
VersionEntrySnapshot {
name: "Audio transfer".into(),
version: federation::AUDIO_PROTOCOL_VERSION.to_string(),
},
VersionEntrySnapshot {
name: "Connected devices".into(),
version: format!(
"{} · accepts v1",
manifest
.protocols
.get(DEVICE_SYNC_ID)
.copied()
.unwrap_or_default()
),
},
],
}
}
pub(super) fn find_catalog_track<'a>(
library: &'a LibrarySnapshot,
search: &'a SearchResults,
key: &TrackKey,
) -> Option<&'a Track> {
library
.featured_releases
.iter()
.flat_map(|release| release.tracks.iter())
.chain(
library
.playlists
.iter()
.flat_map(|playlist| playlist.tracks.iter()),
)
.chain(
search
.releases
.iter()
.flat_map(|release| release.tracks.iter()),
)
.chain(search.tracks.iter())
.find(|track| track.key.matches(key))
}
pub(super) fn portable_playback_placeholder(
wire: &music_dht::device_sync::PlaybackTrack,
content_id: ContentId,
) -> Track {
let peer_id = "content".to_owned();
let refs = |names: &[String]| {
names
.iter()
.map(|name| ArtistRef {
key: ArtistKey::Federation {
peer_id: peer_id.clone(),
id: format!("name:{}", music_dht::normalize_name(name)),
},
name: name.clone(),
})
.collect::<Vec<_>>()
};
Track {
key: TrackKey::remote(content_id.clone()),
title: wire.title.clone(),
artist: playback_artist_line(&wire.artist_names, &wire.featured_artist_names),
artists: refs(&wire.artist_names),
featured_artists: refs(&wire.featured_artist_names),
release: wire.release_title.clone(),
release_id: ReleaseKey::Federation {
peer_id: peer_id.clone(),
id: format!("name:{}", music_dht::normalize_name(&wire.release_title)),
},
duration_seconds: wire.duration_seconds,
track_number: wire
.track_number
.and_then(|value| u32::try_from(value).ok()),
disc_number: wire.disc_number.and_then(|value| u32::try_from(value).ok()),
cover_uri: None,
audio_format: wire.audio_format.clone(),
audio_bitrate_kbps: wire
.audio_bitrate
.and_then(|value| u32::try_from(value).ok()),
audio_sample_rate_hz: wire
.audio_sample_rate
.and_then(|value| u32::try_from(value).ok()),
audio_bit_depth: wire
.audio_bit_depth
.and_then(|value| u32::try_from(value).ok()),
file_size_bytes: wire
.file_size_bytes
.and_then(|value| u64::try_from(value).ok()),
liked: false,
audio_source: AudioSource::Federation {
peer_id: String::new(),
content_id,
},
}
}
pub(super) fn playback_artist_line(main: &[String], featured: &[String]) -> String {
match (main.is_empty(), featured.is_empty()) {
(false, false) => format!("{} feat. {}", main.join(", "), featured.join(", ")),
(false, true) => main.join(", "),
(true, false) => format!("feat. {}", featured.join(", ")),
(true, true) => String::new(),
}
}
pub(super) fn extrapolated_control_position(
state: &music_dht::device_sync::PlaybackStateWire,
elapsed: Duration,
) -> f64 {
let elapsed = if state.playing && !state.paused {
elapsed.as_secs_f64()
} else {
0.0
};
(state.position_secs + elapsed).max(0.0)
}
pub(super) fn remote_snapshot_has_authority(
role: DevicePlaybackRole,
active_device_id: &str,
local_status: PlaybackStatus,
local_queue_empty: bool,
snapshot: &music_dht::device_sync::PlaybackSnapshot,
) -> bool {
if !snapshot.active {
return false;
}
match role {
DevicePlaybackRole::Control => snapshot.device_id == active_device_id,
DevicePlaybackRole::Active => local_status == PlaybackStatus::Stopped && local_queue_empty,
}
}
pub(super) fn playback_state_acknowledges_command(
expected: &music_dht::device_sync::PlaybackStateWire,
actual: &music_dht::device_sync::PlaybackStateWire,
seek: bool,
elapsed: Duration,
) -> bool {
let same_queue = expected.queue.len() == actual.queue.len()
&& expected
.queue
.iter()
.zip(&actual.queue)
.all(|(expected, actual)| playback_tracks_equivalent(expected, actual));
let same_transport = expected.queue_pos == actual.queue_pos
&& expected.playing == actual.playing
&& expected.paused == actual.paused
&& expected.volume == actual.volume
&& expected.shuffle == actual.shuffle
&& expected.repeat == actual.repeat;
if !same_queue || !same_transport {
return false;
}
if !seek {
return true;
}
let expected_position = extrapolated_control_position(expected, elapsed);
(expected_position - actual.position_secs).abs() <= CONTROL_POSITION_ACK_TOLERANCE_SECONDS
}
pub(super) fn playback_tracks_equivalent(
left: &music_dht::device_sync::PlaybackTrack,
right: &music_dht::device_sync::PlaybackTrack,
) -> bool {
let content_id = |track: &music_dht::device_sync::PlaybackTrack| {
track
.content_id
.as_deref()
.and_then(music_dht::normalize_content_id)
.or_else(|| {
track
.fed
.as_ref()
.and_then(|fed| music_dht::normalize_content_id(&fed.content_id))
})
};
match (content_id(left), content_id(right)) {
(Some(left), Some(right)) => left == right,
_ => {
music_dht::normalize_name(&left.title) == music_dht::normalize_name(&right.title)
&& music_dht::normalize_name(&left.release_title)
== music_dht::normalize_name(&right.release_title)
&& left.track_number == right.track_number
&& left.disc_number == right.disc_number
&& normalized_artist_names(&left.artist_names)
== normalized_artist_names(&right.artist_names)
}
}
}
pub(super) fn normalized_artist_names(names: &[String]) -> Vec<String> {
names
.iter()
.map(|name| music_dht::normalize_name(name))
.collect()
}
pub(super) fn unix_time_ms() -> i64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map_or(0, |duration| {
i64::try_from(duration.as_millis()).unwrap_or(i64::MAX)
})
}
#[allow(
clippy::cast_possible_truncation,
clippy::cast_sign_loss,
reason = "the playback volume is clamped to the wire protocol's 0..=100 range"
)]
pub(super) fn volume_percent(volume: f32) -> u8 {
(volume.clamp(0.0, 1.0) * 100.0).round() as u8
}
pub(super) fn track_to_synced_fed(track: &Track) -> Option<music_dht::device_sync::SyncedFedTrack> {
let (owner, item_id) = track.key.federation_id()?;
let content_id = track.key.content_id()?.as_str().to_owned();
Some(music_dht::device_sync::SyncedFedTrack {
item_id: item_id.to_owned(),
owner: owner.to_owned(),
title: track.title.clone(),
artist_names: track
.artists
.iter()
.map(|artist| artist.name.clone())
.collect(),
featured_artist_names: track
.featured_artists
.iter()
.map(|artist| artist.name.clone())
.collect(),
year: None,
duration_seconds: (track.duration_seconds.is_finite() && track.duration_seconds > 0.0)
.then(|| {
i64::try_from(Duration::from_secs_f64(track.duration_seconds).as_secs())
.unwrap_or(i64::MAX)
}),
content_id,
release_title: (!track.release.is_empty()).then(|| track.release.clone()),
track_number: track
.track_number
.and_then(|value| i32::try_from(value).ok()),
disc_number: track
.disc_number
.and_then(|value| i32::try_from(value).ok()),
})
}
pub(super) fn track_to_library_fed(track: &Track) -> Option<furumi_library::FederatedTrack> {
let synced = track_to_synced_fed(track)?;
Some(furumi_library::FederatedTrack {
item_id: synced.item_id,
owner: synced.owner,
own: false,
title: synced.title,
artist_names: synced.artist_names,
featured_artist_names: synced.featured_artist_names,
year: synced.year,
duration_seconds: synced.duration_seconds,
content_id: Some(synced.content_id),
release_title: synced.release_title,
track_number: synced.track_number,
disc_number: synced.disc_number,
})
}
pub(super) fn track_to_playback_track(track: &Track) -> music_dht::device_sync::PlaybackTrack {
let id = track.key.local_id().map_or(-1, LocalTrackId::get);
let release_id = match track.release_id {
ReleaseKey::Local(id) => id.get(),
ReleaseKey::Federation { .. } => -1,
};
music_dht::device_sync::PlaybackTrack {
id,
title: track.title.clone(),
track_number: track
.track_number
.and_then(|value| i32::try_from(value).ok()),
disc_number: track
.disc_number
.and_then(|value| i32::try_from(value).ok()),
duration_seconds: track.duration_seconds,
artist_names: track
.artists
.iter()
.map(|artist| artist.name.clone())
.collect(),
featured_artist_names: track
.featured_artists
.iter()
.map(|artist| artist.name.clone())
.collect(),
release_id,
release_title: track.release.clone(),
release_year: None,
file_path: String::new(),
content_id: track.key.content_id().map(|id| id.as_str().to_owned()),
audio_format: track.audio_format.clone(),
audio_bitrate: track
.audio_bitrate_kbps
.and_then(|value| i32::try_from(value).ok()),
audio_sample_rate: track
.audio_sample_rate_hz
.and_then(|value| i32::try_from(value).ok()),
audio_bit_depth: track
.audio_bit_depth
.and_then(|value| i32::try_from(value).ok()),
file_size_bytes: track
.file_size_bytes
.and_then(|value| i64::try_from(value).ok()),
play_count: 0,
fed: track_to_synced_fed(track),
}
}
pub(super) fn sanitize_filename(value: &str) -> String {
let clean: String = value
.chars()
.map(|c| {
if c.is_alphanumeric() || matches!(c, '-' | '_') {
c
} else {
'_'
}
})
.collect();
if clean.is_empty() {
"track".into()
} else {
clean
}
}
pub(super) fn federation_specs(
library: &furumi_library::Library,
) -> anyhow::Result<Vec<music_dht::ItemSpec>> {
let export = library.federation_export()?;
let mut specs = Vec::new();
for (id, name) in export.artists {
specs.push(music_dht::ItemSpec {
local_key: format!("artist:{id}"),
kind: music_dht::ItemKind::Artist,
name,
artist_names: Vec::new(),
featured_artist_names: Vec::new(),
year: None,
release_type: None,
release_title: None,
track_number: None,
disc_number: None,
duration_seconds: None,
content_id: None,
});
}
for release in export.releases {
specs.push(music_dht::ItemSpec {
local_key: format!("release:{}", release.id),
kind: music_dht::ItemKind::Release,
name: release.title,
artist_names: release.artist_names,
featured_artist_names: Vec::new(),
year: release.year,
release_type: Some(release.release_type),
release_title: None,
track_number: None,
disc_number: None,
duration_seconds: None,
content_id: None,
});
}
for track in export.tracks {
specs.push(music_dht::ItemSpec {
local_key: format!("track:{}", track.id),
kind: music_dht::ItemKind::Track,
name: track.title,
artist_names: track.artist_names,
featured_artist_names: track.featured_artist_names,
year: track.year,
release_type: Some(track.release_type),
release_title: Some(track.release_title),
track_number: track.track_number,
disc_number: track.disc_number,
duration_seconds: (track.duration_seconds > 0.0).then_some(track.duration_seconds),
content_id: track.content_id,
});
}
Ok(specs)
}
pub(super) fn apply_federated_metadata(
import: &mut furumi_library::import::TrackImport,
metadata: Option<federation::TrackMetadata>,
) {
let Some(metadata) = metadata else {
return;
};
if !metadata.title.trim().is_empty() {
import.title = metadata.title;
}
if !metadata.artists.is_empty() {
import.artists = metadata.artists;
}
if !metadata.featured_artists.is_empty() {
import.featured_artists = metadata.featured_artists;
}
if !metadata.album_artists.is_empty() {
import.album_artists = metadata.album_artists;
}
if !metadata.release_title.trim().is_empty() {
import.release_title = metadata.release_title;
}
import.release_type = metadata.release_type.or(import.release_type.take());
import.year = metadata.year.or(import.year);
import.track_number = metadata.track_number.or(import.track_number);
import.disc_number = metadata.disc_number.or(import.disc_number);
import.duration_seconds = metadata.duration_seconds.unwrap_or(import.duration_seconds);
import.audio_format = metadata.audio_format.or(import.audio_format.take());
import.audio_bitrate = metadata.audio_bitrate.or(import.audio_bitrate);
import.audio_sample_rate = metadata.audio_sample_rate.or(import.audio_sample_rate);
import.audio_bit_depth = metadata.audio_bit_depth.or(import.audio_bit_depth);
}
pub(super) fn spawn_settings_worker(
store: SettingsStore,
receiver: std_mpsc::Receiver<SettingsSnapshot>,
events: mpsc::Sender<InternalEvent>,
) {
let report = events.clone();
if let Err(error) = thread::Builder::new()
.name("furumi-settings-storage".into())
.spawn(move || {
while let Ok(mut settings) = receiver.recv() {
for newer in receiver.try_iter() {
settings = newer;
}
let result = store.save(&settings).map_err(|error| error.to_string());
if events
.blocking_send(InternalEvent::SettingsPersisted(result))
.is_err()
{
break;
}
}
})
{
let _ = report.blocking_send(InternalEvent::SettingsPersisted(Err(format!(
"settings worker: {error}"
))));
}
}
pub(super) fn local_search_results(
catalog: &furumi_library::Library,
query: &str,
) -> anyhow::Result<SearchResults> {
let found = catalog.search(query, 50)?;
let artists = found
.artists
.into_iter()
.map(|artist| Artist {
key: ArtistKey::local(ArtistId::new(artist.id)),
source: CatalogSource::Local,
name: artist.name,
artwork: Artwork {
uri: artist.image_path,
},
release_count: usize::try_from(artist.release_count.max(0)).unwrap_or(usize::MAX),
track_count: usize::try_from(artist.track_count.max(0)).unwrap_or(usize::MAX),
})
.collect();
let mut releases = Vec::with_capacity(found.releases.len());
for card in found.releases {
let detail = catalog.release(card.id)?;
releases.push(library_release(detail));
}
let tracks = found
.tracks
.into_iter()
.map(|track| library_track(track, ""))
.collect();
Ok(SearchResults {
artists,
releases,
tracks,
})
}
pub(super) fn merge_search_results(target: &mut SearchResults, incoming: SearchResults) {
for artist in incoming.artists {
if let Some(existing) = target
.artists
.iter_mut()
.find(|item| item.key == artist.key)
{
existing.release_count = existing.release_count.max(artist.release_count);
existing.track_count = existing.track_count.max(artist.track_count);
if existing.artwork.uri.is_none() && artist.artwork.uri.is_some() {
existing.artwork = artist.artwork;
}
} else {
target.artists.push(artist);
}
}
for mut release in incoming.releases {
populate_release_contributors(&mut release);
sort_release_tracks(&mut release.tracks);
if let Some(existing) = target
.releases
.iter_mut()
.find(|item| item.key == release.key)
{
merge_release_preserving_local(existing, release);
} else {
target.releases.push(release);
}
}
for track in incoming.tracks {
if let Some(existing) = target
.tracks
.iter_mut()
.find(|item| item.same_catalog_track(&track))
{
if matches!(existing.audio_source, AudioSource::Federation { .. })
&& matches!(track.audio_source, AudioSource::LocalFile(_))
{
*existing = track;
}
} else {
target.tracks.push(track);
}
}
}
pub(super) fn merge_release_preserving_local(target: &mut Release, incoming: Release) {
if matches!(target.source, CatalogSource::Federation { .. })
&& matches!(incoming.source, CatalogSource::Local)
{
let previous = std::mem::replace(target, incoming);
merge_release_preserving_local(target, previous);
return;
}
if target.artwork.uri.is_none() && incoming.artwork.uri.is_some() {
target.artwork = incoming.artwork;
}
merge_artist_refs(&mut target.artists, incoming.artists);
merge_artist_refs(&mut target.featured_artists, incoming.featured_artists);
for track in incoming.tracks {
if let Some(current) = target
.tracks
.iter_mut()
.find(|item| tracks_match_within_release(item, &track))
{
if matches!(current.audio_source, AudioSource::Federation { .. })
&& matches!(track.audio_source, AudioSource::LocalFile(_))
{
*current = track;
}
} else {
target.tracks.push(track);
}
}
sort_release_tracks(&mut target.tracks);
populate_release_contributors(target);
}
fn sort_release_tracks(tracks: &mut [Track]) {
tracks.sort_by(|left, right| {
left.disc_number
.unwrap_or(1)
.cmp(&right.disc_number.unwrap_or(1))
.then_with(|| {
left.track_number
.unwrap_or(u32::MAX)
.cmp(&right.track_number.unwrap_or(u32::MAX))
})
.then_with(|| left.title.cmp(&right.title))
});
}
/// Match provider records within an already identified release. Track/disc
/// position is more reliable here than optional remote descriptive metadata.
pub(super) fn tracks_match_within_release(left: &Track, right: &Track) -> bool {
if left.key.matches(&right.key) {
return true;
}
if let (Some(left_number), Some(right_number)) = (left.track_number, right.track_number)
&& left_number == right_number
&& left.disc_number.unwrap_or(1) == right.disc_number.unwrap_or(1)
{
return true;
}
music_dht::normalize_name(&left.title) == music_dht::normalize_name(&right.title)
}
pub(super) fn merge_artist_refs(target: &mut Vec<ArtistRef>, incoming: Vec<ArtistRef>) {
let mut known = target
.iter()
.map(|artist| music_dht::normalize_name(&artist.name))
.collect::<HashSet<_>>();
target.extend(
incoming
.into_iter()
.filter(|artist| known.insert(music_dht::normalize_name(&artist.name))),
);
}
pub(super) fn populate_release_contributors(release: &mut Release) {
let mut known = release
.artists
.iter()
.chain(release.featured_artists.iter())
.map(|artist| music_dht::normalize_name(&artist.name))
.collect::<HashSet<_>>();
for artist in release
.tracks
.iter()
.flat_map(|track| track.artists.iter().chain(track.featured_artists.iter()))
{
if known.insert(music_dht::normalize_name(&artist.name)) {
release.featured_artists.push(artist.clone());
}
}
}
pub(super) fn library_release(detail: furumi_library::ReleaseDetail) -> Release {
let (artists, featured_artists) = release_artist_roles(&detail);
let fallback = artists
.iter()
.map(|artist| artist.name.as_str())
.collect::<Vec<_>>()
.join(", ");
Release {
key: ReleaseKey::local(ReleaseId::new(detail.id)),
source: CatalogSource::Local,
title: detail.title,
artists,
featured_artists,
release_type: detail.release_type,
year: detail.year,
artwork: Artwork {
uri: detail.cover_path,
},
tracks: detail
.tracks
.into_iter()
.map(|track| library_track(track, &fallback))
.collect(),
}
}
pub(super) fn library_snapshot(
catalog: &furumi_library::Library,
) -> anyhow::Result<LibrarySnapshot> {
let liked_ids = catalog
.liked_content_ids()?
.into_iter()
.chain(catalog.fed_like_ids()?)
.collect::<HashSet<_>>();
let artist_cards = catalog.artists(
1,
i64::MAX,
furumi_library::LibraryFilters {
source_mode: furumi_library::LibrarySourceMode::Local,
..furumi_library::LibraryFilters::default()
},
)?;
let artists = artist_cards
.items
.into_iter()
.map(|artist| Artist {
key: ArtistKey::local(ArtistId::new(artist.id)),
source: CatalogSource::Local,
name: artist.name,
artwork: Artwork {
uri: artist.image_path,
},
release_count: usize::try_from(artist.release_count.max(0)).unwrap_or(usize::MAX),
track_count: usize::try_from(artist.track_count.max(0)).unwrap_or(usize::MAX),
})
.collect();
let mut releases = Vec::new();
for card in catalog.releases()? {
let detail = catalog.release(card.id)?;
let (artist_refs, featured_artist_refs) = release_artist_roles(&detail);
let artist_line = artist_refs
.iter()
.map(|artist| artist.name.as_str())
.collect::<Vec<_>>()
.join(", ");
let mut tracks = detail
.tracks
.into_iter()
.map(|track| library_track(track, &artist_line))
.collect::<Vec<_>>();
for track in &mut tracks {
track.liked = track_is_liked(track, &liked_ids);
}
releases.push(Release {
key: ReleaseKey::local(ReleaseId::new(detail.id)),
source: CatalogSource::Local,
title: detail.title,
artists: artist_refs,
featured_artists: featured_artist_refs,
release_type: detail.release_type,
year: detail.year,
artwork: Artwork {
uri: detail.cover_path,
},
tracks,
});
}
let recently_played = releases
.iter()
.flat_map(|release| release.tracks.iter().cloned())
.take(12)
.collect();
let mut playlists = Vec::new();
for card in catalog.playlists()? {
let detail = catalog.playlist(card.id)?;
let mut tracks = detail
.tracks
.into_iter()
.map(|track| library_track(track, ""))
.collect::<Vec<_>>();
for track in &mut tracks {
track.liked = track_is_liked(track, &liked_ids);
}
playlists.push(PlaylistSnapshot {
id: card.id,
title: detail.title,
is_likes: card.kind == "likes",
tracks,
});
}
Ok(LibrarySnapshot {
artists,
featured_releases: releases,
recently_played,
playlists,
})
}
pub(super) fn track_is_liked(track: &Track, liked_ids: &HashSet<String>) -> bool {
track
.key
.content_id()
.is_some_and(|id| liked_ids.contains(id.as_str()))
|| track
.key
.federation_id()
.is_some_and(|(_, item)| liked_ids.contains(item))
}
pub(super) fn release_artist_roles(
detail: &furumi_library::ReleaseDetail,
) -> (Vec<ArtistRef>, Vec<ArtistRef>) {
let mut main_counts = HashMap::<i64, usize>::new();
let mut featured = HashMap::<i64, String>::new();
for track in &detail.tracks {
for artist in &track.artists {
*main_counts.entry(artist.id).or_default() += 1;
}
for artist in &track.featured_artists {
featured
.entry(artist.id)
.or_insert_with(|| artist.name.clone());
}
}
let mut main_artists = Vec::new();
let mut featured_artists = Vec::new();
let mut known = HashSet::new();
for artist in &detail.artists {
known.insert(artist.id);
let reference = ArtistRef {
key: ArtistKey::local(ArtistId::new(artist.id)),
name: artist.name.clone(),
};
if main_counts.contains_key(&artist.id) || !featured.contains_key(&artist.id) {
main_artists.push(reference);
} else {
featured_artists.push(reference);
}
}
for track in &detail.tracks {
for artist in &track.artists {
if known.insert(artist.id) {
main_artists.push(ArtistRef {
key: ArtistKey::local(ArtistId::new(artist.id)),
name: artist.name.clone(),
});
}
}
}
for (id, name) in featured {
if known.insert(id) {
featured_artists.push(ArtistRef {
key: ArtistKey::local(ArtistId::new(id)),
name,
});
}
}
main_artists.sort_by(|left, right| {
let count = |artist: &ArtistRef| match artist.key {
ArtistKey::Local(id) => main_counts.get(&id.get()).copied().unwrap_or_default(),
ArtistKey::Federation { .. } => 0,
};
count(right).cmp(&count(left))
});
(main_artists, featured_artists)
}
pub(super) fn library_track(track: furumi_library::TrackItem, fallback_artist: &str) -> Track {
let content_id = track
.content_id
.as_deref()
.and_then(|content_id| ContentId::parse(content_id).ok());
let local_id = LocalTrackId::new(track.id);
let key = content_id.map_or_else(
|| TrackKey::local(local_id),
|content_id| TrackKey::new(local_id, content_id),
);
let artists = track
.artists
.iter()
.map(|artist| ArtistRef {
key: ArtistKey::local(ArtistId::new(artist.id)),
name: artist.name.clone(),
})
.collect::<Vec<_>>();
let featured_artists = track
.featured_artists
.iter()
.map(|artist| ArtistRef {
key: ArtistKey::local(ArtistId::new(artist.id)),
name: artist.name.clone(),
})
.collect::<Vec<_>>();
let artist = {
let value = track.artist_line();
if value.is_empty() {
fallback_artist.to_string()
} else {
value
}
};
Track {
key,
title: track.title,
artist,
artists,
featured_artists,
release: track.release_title,
release_id: ReleaseKey::local(ReleaseId::new(track.release_id)),
duration_seconds: track.duration_seconds,
track_number: track
.track_number
.and_then(|value| u32::try_from(value).ok()),
disc_number: track
.disc_number
.and_then(|value| u32::try_from(value).ok()),
cover_uri: track.cover_path,
audio_format: track.audio_format,
audio_bitrate_kbps: track
.audio_bitrate
.and_then(|value| u32::try_from(value).ok()),
audio_sample_rate_hz: track
.audio_sample_rate
.and_then(|value| u32::try_from(value).ok()),
audio_bit_depth: track
.audio_bit_depth
.and_then(|value| u32::try_from(value).ok()),
file_size_bytes: track
.file_size_bytes
.and_then(|value| u64::try_from(value).ok()),
liked: false,
audio_source: AudioSource::LocalFile(track.file_path.into()),
}
}
+462
View File
@@ -0,0 +1,462 @@
use super::*;
fn merge_test_track(key: TrackKey, audio_source: AudioSource) -> Track {
Track {
key,
title: "Track".into(),
artist: "Artist".into(),
artists: vec![ArtistRef {
key: ArtistKey::local(ArtistId::new(1)),
name: "Artist".into(),
}],
featured_artists: Vec::new(),
release: "Album".into(),
release_id: ReleaseKey::local(ReleaseId::new(1)),
duration_seconds: 180.0,
track_number: Some(1),
disc_number: Some(1),
cover_uri: None,
audio_format: Some("flac".into()),
audio_bitrate_kbps: None,
audio_sample_rate_hz: None,
audio_bit_depth: None,
file_size_bytes: None,
liked: false,
audio_source,
}
}
#[test]
fn merging_search_results_deduplicates_source_keys() {
let artist = Artist {
key: ArtistKey::Federation {
peer_id: "peer".into(),
id: "artist".into(),
},
source: CatalogSource::Federation {
peer_id: "peer".into(),
},
name: "Artist".into(),
artwork: Artwork::default(),
release_count: 0,
track_count: 0,
};
let mut target = SearchResults {
artists: vec![artist.clone()],
..SearchResults::default()
};
merge_search_results(
&mut target,
SearchResults {
artists: vec![artist],
..SearchResults::default()
},
);
assert_eq!(target.artists.len(), 1);
}
#[test]
fn newly_received_release_tracks_are_sorted_by_disc_and_track_number() {
let mut fifth = merge_test_track(
TrackKey::federation("peer".into(), "five".into(), None),
AudioSource::Federation {
peer_id: "peer".into(),
content_id: ContentId::parse(format!("b3:{}", "5".repeat(64))).unwrap(),
},
);
fifth.title = "Los".into();
fifth.track_number = Some(5);
let mut first = merge_test_track(
TrackKey::federation("peer".into(), "one".into(), None),
AudioSource::Federation {
peer_id: "peer".into(),
content_id: ContentId::parse(format!("b3:{}", "1".repeat(64))).unwrap(),
},
);
first.title = "Reise, Reise".into();
let mut target = SearchResults::default();
merge_search_results(
&mut target,
SearchResults {
releases: vec![Release {
key: ReleaseKey::Federation {
peer_id: "peer".into(),
id: "reise-reise".into(),
},
source: CatalogSource::Federation {
peer_id: "peer".into(),
},
title: "Reise, Reise".into(),
artists: first.artists.clone(),
featured_artists: Vec::new(),
release_type: "album".into(),
year: Some(2004),
artwork: Artwork::default(),
tracks: vec![fifth, first],
}],
..SearchResults::default()
},
);
assert_eq!(
target.releases[0]
.tracks
.iter()
.map(|track| track.track_number)
.collect::<Vec<_>>(),
vec![Some(1), Some(5)]
);
}
#[test]
fn selected_track_key_survives_context_filtering_and_reordering() {
let selected = TrackKey::local(LocalTrackId::new(2));
let resolved = vec![
merge_test_track(
TrackKey::local(LocalTrackId::new(3)),
AudioSource::LocalFile("three.flac".into()),
),
merge_test_track(selected.clone(), AudioSource::LocalFile("two.flac".into())),
];
assert_eq!(selected_track_position(&resolved, &selected), 1);
}
#[test]
fn queue_artwork_requests_are_grouped_per_peer_and_release() {
let first_content = ContentId::parse(format!("b3:{}", "1".repeat(64))).unwrap();
let second_content = ContentId::parse(format!("b3:{}", "2".repeat(64))).unwrap();
let mut first = merge_test_track(
TrackKey::federation("peer".into(), "one".into(), Some(first_content.clone())),
AudioSource::Federation {
peer_id: "peer".into(),
content_id: first_content,
},
);
let mut second = merge_test_track(
TrackKey::federation("peer".into(), "two".into(), Some(second_content.clone())),
AudioSource::Federation {
peer_id: "peer".into(),
content_id: second_content,
},
);
first.title = "One".into();
second.title = "Two".into();
assert_eq!(
Actor::queue_artwork_request_id(&first),
Actor::queue_artwork_request_id(&second)
);
second.release = "Another album".into();
assert_ne!(
Actor::queue_artwork_request_id(&first),
Actor::queue_artwork_request_id(&second)
);
}
#[test]
fn federated_release_enrichment_never_replaces_a_local_track() {
let release_key = ReleaseKey::local(ReleaseId::new(1));
let artist = ArtistRef {
key: ArtistKey::local(ArtistId::new(1)),
name: "Artist".into(),
};
let local_track = merge_test_track(
TrackKey::local(LocalTrackId::new(1)),
AudioSource::LocalFile("track.flac".into()),
);
let mut target = SearchResults {
releases: vec![Release {
key: release_key.clone(),
source: CatalogSource::Local,
title: "Album".into(),
artists: vec![artist.clone()],
featured_artists: Vec::new(),
release_type: "album".into(),
year: Some(2026),
artwork: Artwork::default(),
tracks: vec![local_track],
}],
..SearchResults::default()
};
let content_id = ContentId::parse(format!("b3:{}", "a".repeat(64))).unwrap();
let remote_track = merge_test_track(
TrackKey::federation(
"peer".into(),
"remote-item".into(),
Some(content_id.clone()),
),
AudioSource::Federation {
peer_id: "peer".into(),
content_id,
},
);
merge_search_results(
&mut target,
SearchResults {
releases: vec![Release {
key: release_key,
source: CatalogSource::Federation {
peer_id: "peer".into(),
},
title: "Album".into(),
artists: vec![artist],
featured_artists: Vec::new(),
release_type: "album".into(),
year: Some(2026),
artwork: Artwork::default(),
tracks: vec![remote_track],
}],
..SearchResults::default()
},
);
let tracks = &target.releases[0].tracks;
assert_eq!(tracks.len(), 1);
assert!(matches!(tracks[0].audio_source, AudioSource::LocalFile(_)));
}
#[test]
fn incomplete_federated_metadata_cannot_shadow_local_release_metadata() {
let release_key = ReleaseKey::local(ReleaseId::new(1));
let artist = ArtistRef {
key: ArtistKey::local(ArtistId::new(1)),
name: "Artist".into(),
};
let mut local_track = merge_test_track(
TrackKey::local(LocalTrackId::new(1)),
AudioSource::LocalFile("track.flac".into()),
);
local_track.cover_uri = Some("local-cover.png".into());
local_track.audio_bitrate_kbps = Some(1_411);
let mut target = Release {
key: release_key.clone(),
source: CatalogSource::Local,
title: "Complete local album title".into(),
artists: vec![artist.clone()],
featured_artists: Vec::new(),
release_type: "album".into(),
year: Some(2026),
artwork: Artwork {
uri: Some("local-cover.png".into()),
},
tracks: vec![local_track],
};
let content_id = ContentId::parse(format!("b3:{}", "b".repeat(64))).unwrap();
let mut incomplete_remote = merge_test_track(
TrackKey::federation("peer".into(), "item".into(), Some(content_id.clone())),
AudioSource::Federation {
peer_id: "peer".into(),
content_id,
},
);
incomplete_remote.title = String::new();
incomplete_remote.artist = String::new();
incomplete_remote.artists.clear();
incomplete_remote.audio_format = None;
merge_release_preserving_local(
&mut target,
Release {
key: release_key,
source: CatalogSource::Federation {
peer_id: "peer".into(),
},
title: String::new(),
artists: vec![artist],
featured_artists: Vec::new(),
release_type: String::new(),
year: None,
artwork: Artwork::default(),
tracks: vec![incomplete_remote],
},
);
assert!(matches!(target.source, CatalogSource::Local));
assert_eq!(target.title, "Complete local album title");
assert_eq!(target.tracks.len(), 1);
assert_eq!(target.tracks[0].title, "Track");
assert_eq!(target.tracks[0].audio_format.as_deref(), Some("flac"));
assert_eq!(target.tracks[0].audio_bitrate_kbps, Some(1_411));
assert!(matches!(
target.tracks[0].audio_source,
AudioSource::LocalFile(_)
));
}
#[test]
fn queue_resolver_keeps_tracks_nested_in_a_federated_release() {
let content_id = ContentId::parse(format!("b3:{}", "c".repeat(64))).unwrap();
let key = TrackKey::federation(
"peer".into(),
"remote-track".into(),
Some(content_id.clone()),
);
let track = merge_test_track(
key.clone(),
AudioSource::Federation {
peer_id: "peer".into(),
content_id,
},
);
let search = SearchResults {
releases: vec![Release {
key: ReleaseKey::Federation {
peer_id: "peer".into(),
id: "album".into(),
},
source: CatalogSource::Federation {
peer_id: "peer".into(),
},
title: "Album".into(),
artists: track.artists.clone(),
featured_artists: Vec::new(),
release_type: "album".into(),
year: Some(2026),
artwork: Artwork::default(),
tracks: vec![track],
}],
..SearchResults::default()
};
let library = LibrarySnapshot::default();
let resolved = find_catalog_track(&library, &search, &key);
assert!(resolved.is_some());
assert!(matches!(
resolved.unwrap().audio_source,
AudioSource::Federation { .. }
));
}
#[test]
fn portable_device_queue_keeps_unresolved_content_as_a_federated_placeholder() {
let content_id = ContentId::parse(format!("b3:{}", "d".repeat(64))).unwrap();
let wire = music_dht::device_sync::PlaybackTrack {
id: 12,
title: "Portable track".into(),
track_number: Some(2),
disc_number: Some(1),
duration_seconds: 123.0,
artist_names: vec!["Artist".into()],
featured_artist_names: vec!["Guest".into()],
release_id: 4,
release_title: "Release".into(),
release_year: Some(2026),
file_path: String::new(),
content_id: Some(content_id.as_str().into()),
audio_format: Some("flac".into()),
audio_bitrate: Some(1_411),
audio_sample_rate: Some(44_100),
audio_bit_depth: Some(16),
file_size_bytes: Some(42),
play_count: 0,
fed: None,
};
let placeholder = portable_playback_placeholder(&wire, content_id.clone());
assert_eq!(placeholder.key.content_id(), Some(&content_id));
assert_eq!(placeholder.artist, "Artist feat. Guest");
assert!(matches!(
placeholder.audio_source,
AudioSource::Federation { ref peer_id, .. } if peer_id.is_empty()
));
}
fn playback_test_state(content_byte: char) -> music_dht::device_sync::PlaybackStateWire {
music_dht::device_sync::PlaybackStateWire {
queue: vec![music_dht::device_sync::PlaybackTrack {
id: 1,
title: format!("Track {content_byte}"),
track_number: Some(1),
disc_number: Some(1),
duration_seconds: 180.0,
artist_names: vec!["Artist".into()],
featured_artist_names: Vec::new(),
release_id: 1,
release_title: "Album".into(),
release_year: Some(2026),
file_path: String::new(),
content_id: Some(format!("b3:{}", content_byte.to_string().repeat(64))),
audio_format: Some("flac".into()),
audio_bitrate: None,
audio_sample_rate: None,
audio_bit_depth: None,
file_size_bytes: None,
play_count: 0,
fed: None,
}],
queue_pos: 0,
playing: true,
paused: false,
idle_since_ms: None,
position_secs: 40.0,
volume: 72,
shuffle: false,
repeat: music_dht::device_sync::PlaybackRepeat::Off,
}
}
#[test]
fn control_device_ignores_snapshots_from_the_previous_active_device() {
let snapshot = music_dht::device_sync::PlaybackSnapshot {
device_id: "old-device".into(),
device_name: "Old".into(),
active: true,
updated_at_ms: 1,
state: playback_test_state('a'),
};
assert!(!remote_snapshot_has_authority(
DevicePlaybackRole::Control,
"current-device",
PlaybackStatus::Playing,
false,
&snapshot,
));
assert!(remote_snapshot_has_authority(
DevicePlaybackRole::Control,
"old-device",
PlaybackStatus::Playing,
false,
&snapshot,
));
}
#[test]
fn stale_remote_state_cannot_ack_a_new_control_command() {
let expected = playback_test_state('a');
let stale_track = playback_test_state('b');
assert!(!playback_state_acknowledges_command(
&expected,
&stale_track,
true,
Duration::from_secs(2),
));
let mut acknowledged = expected.clone();
acknowledged.position_secs = 42.0;
assert!(playback_state_acknowledges_command(
&expected,
&acknowledged,
true,
Duration::from_secs(2),
));
acknowledged.position_secs = 5.0;
assert!(!playback_state_acknowledges_command(
&expected,
&acknowledged,
true,
Duration::from_secs(2),
));
assert!(playback_state_acknowledges_command(
&expected,
&acknowledged,
false,
Duration::from_secs(2),
));
}