Reworked federation. Added settings page. Fixed connected devices.
Build and Publish / Build and Publish Docker Image (push) Successful in 5m21s

This commit is contained in:
Ultradesu
2026-07-27 16:44:16 +01:00
parent fd766cda24
commit 624e75839d
16 changed files with 4587 additions and 169 deletions
+536
View File
@@ -0,0 +1,536 @@
//! Receiving side of music federation.
//!
//! User-facing identity is content-addressed. An `(owner, item_id)` pair is
//! only a source locator and several locators may resolve the same track.
use std::collections::HashMap;
use anyhow::{Context, Result};
use music_dht::{ItemKind, LibraryItem, normalize_content_id};
use serde::Serialize;
use serde_json::{Value, json};
use sqlx::Row as _;
use tokio::io::AsyncReadExt;
use super::{Federation, now_iso};
const MAX_CATALOG_BYTES: u64 = 4 * 1024 * 1024;
#[derive(Debug, Clone, Serialize)]
pub struct TrackKeyDto {
pub content_id: String,
}
#[derive(Debug, Clone, Serialize)]
pub struct ArtistKeyDto {
pub normalized_name: String,
}
#[derive(Debug, Clone, Serialize)]
pub struct ArtistRefDto {
pub key: ArtistKeyDto,
pub name: String,
pub local_id: Option<i64>,
}
#[derive(Debug, Clone, Serialize)]
pub struct ReleaseKeyDto {
pub normalized_title: String,
pub primary_artists: Vec<String>,
pub release_type: Option<String>,
pub year: Option<i32>,
}
#[derive(Debug, Clone, Serialize)]
pub struct ReleaseRefDto {
pub key: ReleaseKeyDto,
pub local_id: Option<i64>,
pub title: String,
}
#[derive(Debug, Clone, Serialize)]
pub struct FederationSourceDto {
pub owner: String,
pub item_id: String,
}
#[derive(Debug, Clone, Serialize)]
pub struct LocalAvailabilityDto {
pub track_id: i64,
pub stream_url: String,
}
#[derive(Debug, Clone, Serialize)]
pub struct TrackMetadataDto {
pub title: String,
pub artists: Vec<ArtistRefDto>,
pub featured_artists: Vec<ArtistRefDto>,
pub release: Option<ReleaseRefDto>,
pub year: Option<i32>,
pub duration_seconds: Option<f64>,
pub track_number: Option<i32>,
pub disc_number: Option<i32>,
pub cover_url: Option<String>,
}
#[derive(Debug, Clone, Serialize)]
pub struct TrackAvailabilityDto {
pub state: &'static str,
pub local: Option<LocalAvailabilityDto>,
pub federation: Vec<FederationSourceDto>,
}
#[derive(Debug, Clone, Serialize)]
pub struct TrackDto {
pub key: TrackKeyDto,
pub metadata: TrackMetadataDto,
pub availability: TrackAvailabilityDto,
}
#[derive(Debug, Clone, Serialize)]
pub struct SearchEvent {
pub search_id: String,
pub sequence: u64,
pub kind: &'static str,
pub peer: Option<String>,
pub entity_key: Value,
pub entity: Value,
}
impl Federation {
pub fn stream_artist_catalogs(
self: &std::sync::Arc<Self>,
name: String,
) -> tokio::sync::mpsc::UnboundedReceiver<Result<(String, music_dht::catalog::CatalogArtist)>>
{
let (sender, receiver) = tokio::sync::mpsc::unbounded_channel();
let federation = std::sync::Arc::clone(self);
tokio::spawn(async move {
let result = async {
let service = federation.service().await?;
let normalized = music_dht::normalize_name(&name);
let outcome = service
.search_network(&name)
.await
.map_err(|err| anyhow::anyhow!("federated artist search failed: {err}"))?;
let owners: std::collections::HashSet<_> = outcome
.network_results
.iter()
.filter(|item| {
(item.kind == ItemKind::Artist && item.normalized_name == normalized)
|| item
.artist_names
.iter()
.chain(&item.featured_artist_names)
.any(|artist| music_dht::normalize_name(artist) == normalized)
})
.map(|item| item.owner)
.collect();
for owner in owners {
let service = std::sync::Arc::clone(&service);
let sender = sender.clone();
let name = name.clone();
tokio::spawn(async move {
let result = tokio::time::timeout(
std::time::Duration::from_secs(8),
fetch_artist_catalog(&service, owner, &name),
)
.await
.map_err(|_| anyhow::anyhow!("catalog request timed out"))
.and_then(|result| result)
.map(|artist| (owner.to_string(), artist));
let _ = sender.send(result);
});
}
Ok::<(), anyhow::Error>(())
}
.await;
if let Err(err) = result {
let _ = sender.send(Err(err));
}
});
receiver
}
/// Performs one bounded DHT search and returns entity upserts. The HTTP
/// layer streams each upsert independently; catalog fan-out can append
/// events to the same contract without changing the browser model.
pub async fn search_events(&self, search_id: &str, query: &str) -> Result<Vec<SearchEvent>> {
let query = query.trim();
anyhow::ensure!(!query.is_empty(), "search query is empty");
anyhow::ensure!(query.chars().count() <= 200, "search query is too long");
let started = std::time::Instant::now();
let service = self.service().await?;
tracing::info!(
search_id,
query,
connected_peers = service.connected_peers().len(),
known_contacts = service.known_peers().len(),
"federated search started"
);
let own = service.endpoint_id();
let result = tokio::time::timeout(
std::time::Duration::from_secs(20),
service.search_network(query),
)
.await
.map_err(|_| anyhow::anyhow!("federated search timed out after 20 seconds"))?
.map_err(|err| anyhow::anyhow!("federated search failed: {err}"))?;
tracing::info!(
search_id,
query,
local_results = result.local_results.len(),
network_results = result.network_results.len(),
queried_nodes = result.queried_nodes,
elapsed_ms = started.elapsed().as_millis() as u64,
"federated DHT search finished"
);
let all_items: Vec<LibraryItem> = result
.local_results
.into_iter()
.chain(result.network_results)
.collect();
let pool = self.pool().await?;
let mut by_content: HashMap<String, TrackDto> = HashMap::new();
for item in all_items.iter().filter(|item| item.kind == ItemKind::Track) {
let Some(content_id) = item.content_id.as_deref().and_then(normalize_content_id) else {
// A globally usable track reference must be verifiable.
continue;
};
let local = local_availability(&pool, &content_id).await?;
let source = FederationSourceDto {
owner: item.owner.to_string(),
item_id: hex(item.id.as_bytes()),
};
let entry = by_content.entry(content_id.clone()).or_insert_with(|| {
track_from_item(content_id.clone(), item, local, item.owner == own)
});
if !entry.availability.federation.iter().any(|candidate| {
candidate.owner == source.owner && candidate.item_id == source.item_id
}) {
entry.availability.federation.push(source);
}
if entry.availability.local.is_some() {
entry.availability.state = "local";
}
}
let mut tracks: Vec<_> = by_content.into_values().collect();
tracks.sort_by(|left, right| {
left.metadata
.title
.to_lowercase()
.cmp(&right.metadata.title.to_lowercase())
});
let mut events = Vec::with_capacity(all_items.len());
for (index, track) in tracks.into_iter().enumerate() {
persist_track_ref(&pool, &track).await?;
let peer = track
.availability
.federation
.first()
.map(|source| source.owner.clone());
events.push(SearchEvent {
search_id: search_id.to_owned(),
sequence: index as u64 + 1,
kind: "federation.track",
peer,
entity_key: serde_json::to_value(&track.key)?,
entity: serde_json::to_value(track)?,
});
}
let mut artist_peers: HashMap<String, (String, Vec<String>)> = HashMap::new();
let mut releases: HashMap<String, Value> = HashMap::new();
for item in &all_items {
match item.kind {
ItemKind::Artist => {
let key = music_dht::normalize_name(&item.name);
let entry = artist_peers
.entry(key)
.or_insert_with(|| (item.name.clone(), Vec::new()));
let owner = item.owner.to_string();
if !entry.1.contains(&owner) {
entry.1.push(owner);
}
}
ItemKind::Release => {
let artist_keys: Vec<String> = item
.artist_names
.iter()
.map(|name| music_dht::normalize_name(name))
.collect();
let normalized_title = music_dht::normalize_name(&item.name);
let cover_url = all_items
.iter()
.find(|track| {
track.kind == ItemKind::Track
&& track.release_title.as_deref().is_some_and(|title| {
music_dht::normalize_name(title) == normalized_title
})
&& track.year == item.year
})
.map(|track| {
format!(
"/api/player/federation/tracks/artwork?owner={}&item_id={}",
track.owner,
hex(track.id.as_bytes())
)
});
let key = format!(
"{}|{}|{}",
normalized_title,
artist_keys.join(","),
item.year.map_or_else(String::new, |year| year.to_string())
);
releases.entry(key.clone()).or_insert_with(|| {
json!({
"key": {
"normalized_title": music_dht::normalize_name(&item.name),
"primary_artists": artist_keys,
"release_type": null,
"year": item.year,
},
"title": item.name,
"artists": item.artist_names,
"year": item.year,
"cover_url": cover_url,
"sources": [{
"owner": item.owner.to_string(),
"item_id": hex(item.id.as_bytes()),
}],
})
});
}
ItemKind::Track => {}
}
}
for (key, (name, peers)) in artist_peers {
let sequence = events.len() as u64 + 1;
events.push(SearchEvent {
search_id: search_id.to_owned(),
sequence,
kind: "federation.artist",
peer: peers.first().cloned(),
entity_key: json!({ "normalized_name": key }),
entity: json!({
"key": { "normalized_name": key },
"name": name,
"image_url": null,
"peers": peers,
}),
});
}
for (key, release) in releases {
let sequence = events.len() as u64 + 1;
events.push(SearchEvent {
search_id: search_id.to_owned(),
sequence,
kind: "federation.release",
peer: None,
entity_key: json!({ "composite": key }),
entity: release,
});
}
tracing::info!(
search_id,
query,
events = events.len(),
elapsed_ms = started.elapsed().as_millis() as u64,
"federated search response ready"
);
Ok(events)
}
}
async fn fetch_artist_catalog(
service: &music_dht::MusicDhtService,
owner: music_dht::EndpointId,
artist: &str,
) -> Result<music_dht::catalog::CatalogArtist> {
let mut stream = service
.open_stream(owner, super::CATALOG_ALPN)
.await
.map_err(|err| anyhow::anyhow!("cannot reach catalog peer: {err}"))?;
let mut request = serde_json::to_vec(&music_dht::catalog::CatalogRequest {
artist: artist.to_owned(),
want: Some("catalog".to_owned()),
..Default::default()
})?;
request.push(b'\n');
stream.send.write_all(&request).await?;
stream.send.finish()?;
let mut payload = Vec::new();
stream
.recv
.take(MAX_CATALOG_BYTES + 1)
.read_to_end(&mut payload)
.await?;
anyhow::ensure!(
payload.len() as u64 <= MAX_CATALOG_BYTES,
"catalog response is too large"
);
let response: music_dht::catalog::CatalogResponse =
serde_json::from_slice(&payload).context("invalid catalog response")?;
anyhow::ensure!(
response.ok,
"peer refused catalog: {}",
response.error.unwrap_or_else(|| "unknown error".to_owned())
);
response.artist.context("peer returned no artist catalog")
}
fn track_from_item(
content_id: String,
item: &LibraryItem,
local: Option<LocalAvailabilityDto>,
own: bool,
) -> TrackDto {
let owner = item.owner.to_string();
let item_id = hex(item.id.as_bytes());
let artists = artist_refs(&item.artist_names);
let featured_artists = artist_refs(&item.featured_artist_names);
let release = item.release_title.as_ref().map(|title| ReleaseRefDto {
key: ReleaseKeyDto {
normalized_title: music_dht::normalize_name(title),
primary_artists: item
.artist_names
.iter()
.map(|artist| music_dht::normalize_name(artist))
.collect(),
release_type: None,
year: item.year,
},
local_id: None,
title: title.clone(),
});
let state = if local.is_some() || own {
"local"
} else {
"federated"
};
TrackDto {
key: TrackKeyDto { content_id },
metadata: TrackMetadataDto {
title: item.name.clone(),
artists,
featured_artists,
release,
year: item.year,
duration_seconds: item.duration_seconds,
track_number: item.track_number,
disc_number: item.disc_number,
cover_url: Some(format!(
"/api/player/federation/tracks/artwork?owner={owner}&item_id={item_id}"
)),
},
availability: TrackAvailabilityDto {
state,
local,
federation: vec![FederationSourceDto { owner, item_id }],
},
}
}
fn artist_refs(names: &[String]) -> Vec<ArtistRefDto> {
names
.iter()
.map(|name| ArtistRefDto {
key: ArtistKeyDto {
normalized_name: music_dht::normalize_name(name),
},
name: name.clone(),
local_id: None,
})
.collect()
}
async fn local_availability(
pool: &sqlx::PgPool,
content_id: &str,
) -> Result<Option<LocalAvailabilityDto>> {
let row = sqlx::query(
"SELECT t.id
FROM furumusic__federation_content_id_cache c
JOIN furumusic__track t ON t.audio_file_id = c.media_file_id
WHERE c.content_id = $1 AND t.is_hidden = false
LIMIT 1",
)
.bind(content_id)
.fetch_optional(pool)
.await?;
Ok(row.map(|row| {
let track_id: i64 = row.get(0);
LocalAvailabilityDto {
track_id,
stream_url: format!("/api/player/stream/{track_id}"),
}
}))
}
async fn persist_track_ref(pool: &sqlx::PgPool, track: &TrackDto) -> Result<()> {
let metadata = serde_json::to_value(&track.metadata)?;
let local_id = track
.availability
.local
.as_ref()
.map(|local| local.track_id);
let row = sqlx::query(
"INSERT INTO furumusic__track_ref
(content_id, local_track_id, title, release_title, year,
duration_seconds, metadata_json, metadata_authority, created_at, updated_at)
VALUES ($1, $2, $3, $4, $5, $6, $7, 'federation', $8, $8)
ON CONFLICT (content_id) DO UPDATE SET
local_track_id = COALESCE(furumusic__track_ref.local_track_id, EXCLUDED.local_track_id),
title = EXCLUDED.title,
release_title = EXCLUDED.release_title,
year = EXCLUDED.year,
duration_seconds = EXCLUDED.duration_seconds,
metadata_json = EXCLUDED.metadata_json,
updated_at = EXCLUDED.updated_at
RETURNING id",
)
.bind(&track.key.content_id)
.bind(local_id)
.bind(&track.metadata.title)
.bind(
track
.metadata
.release
.as_ref()
.map(|release| &release.title),
)
.bind(track.metadata.year)
.bind(track.metadata.duration_seconds)
.bind(metadata)
.bind(now_iso())
.fetch_one(pool)
.await
.context("persisting content-addressed track reference failed")?;
let track_ref_id: i64 = row.get(0);
for source in &track.availability.federation {
sqlx::query(
"INSERT INTO furumusic__federation_track_source
(track_ref_id, owner_peer_id, item_id, last_seen_ms, metadata_json)
VALUES ($1, $2, $3, $4, $5)
ON CONFLICT (owner_peer_id, item_id) DO UPDATE SET
track_ref_id = EXCLUDED.track_ref_id,
last_seen_ms = EXCLUDED.last_seen_ms,
metadata_json = EXCLUDED.metadata_json",
)
.bind(track_ref_id)
.bind(&source.owner)
.bind(&source.item_id)
.bind(chrono::Utc::now().timestamp_millis())
.bind(json!({ "track": track.metadata }))
.execute(pool)
.await?;
}
Ok(())
}
fn hex(bytes: &[u8]) -> String {
bytes.iter().map(|byte| format!("{byte:02x}")).collect()
}
+168
View File
@@ -771,6 +771,149 @@ pub async fn record_track_like(
Ok(())
}
/// Records a like for a content-addressed track that need not be local yet.
/// The ordinary local likes table remains a projection for materialized
/// tracks; the durable user intent lives under `content_id`.
pub async fn record_content_like(
pool: &sqlx::PgPool,
user_id: i64,
content_id: &str,
liked: bool,
fed: Option<serde_json::Value>,
) -> Result<()> {
let content_id =
music_dht::normalize_content_id(content_id).context("invalid track content id")?;
let fed = fed
.map(serde_json::from_value::<SyncedFedTrack>)
.transpose()
.context("invalid federated track metadata")?;
record_local_op(
pool,
user_id,
SyncOpPayload::TrackLikeSet {
content_id,
liked,
fed: liked.then_some(fed).flatten(),
},
)
.await
}
/// Adds a content-addressed track to a user's playlist without requiring a
/// local numeric track row. Materialization later fills `local_track_id`
/// without replacing this playlist item.
pub async fn record_content_playlist_add(
pool: &sqlx::PgPool,
user_id: i64,
playlist_id: i64,
content_id: &str,
position: i64,
fed: Option<serde_json::Value>,
) -> Result<()> {
let content_id =
music_dht::normalize_content_id(content_id).context("invalid track content id")?;
let fed = fed
.map(serde_json::from_value::<SyncedFedTrack>)
.transpose()
.context("invalid federated track metadata")?;
let title = playlist_title(pool, playlist_id)
.await?
.context("playlist not found")?;
let sync_id = ensure_local_playlist_sync_id(pool, user_id, playlist_id, &title).await?;
record_local_op(
pool,
user_id,
SyncOpPayload::PlaylistTrackAdded {
playlist_id: sync_id,
content_id,
position,
fed,
},
)
.await
}
/// Reconciles content-addressed user state after a track becomes local.
/// History is intentionally not touched: only actual web playback writes it.
pub async fn materialize_content_state(
pool: &sqlx::PgPool,
content_id: &str,
track_id: i64,
) -> Result<()> {
let content_id =
music_dht::normalize_content_id(content_id).context("invalid track content id")?;
let likes = sqlx::query(
"SELECT user_id, liked, hlc_ms
FROM furumusic__fed_state_like WHERE content_id = $1",
)
.bind(&content_id)
.fetch_all(pool)
.await?;
for row in likes {
let user_id: i64 = row.get("user_id");
let liked: bool = row.get("liked");
if liked {
sqlx::query(
"INSERT INTO furumusic__user_liked_track
(user_id, track_id, created_at)
VALUES ($1, $2, $3)
ON CONFLICT (user_id, track_id) DO NOTHING",
)
.bind(user_id)
.bind(track_id)
.bind(iso_from_ms(row.get("hlc_ms")))
.execute(pool)
.await?;
}
}
let playlist_items = sqlx::query(
"SELECT i.user_id, i.playlist_id, i.position
FROM furumusic__fed_state_playlist_item i
WHERE i.content_id = $1 AND i.present = true",
)
.bind(&content_id)
.fetch_all(pool)
.await?;
for row in playlist_items {
let user_id: i64 = row.get("user_id");
let sync_id: String = row.get("playlist_id");
let playlist_id = ensure_playlist_for_item(pool, user_id, &sync_id).await?;
sqlx::query(
"INSERT INTO furumusic__playlist_track
(playlist_id, track_id, position, added_at, added_by_user_id)
SELECT $1, $2, $3, $4, $5
WHERE NOT EXISTS (
SELECT 1 FROM furumusic__playlist_track
WHERE playlist_id = $1 AND track_id = $2
)",
)
.bind(playlist_id)
.bind(track_id)
.bind(row.get::<i64, _>("position") as i32)
.bind(now_iso())
.bind(user_id)
.execute(pool)
.await?;
}
sqlx::query(
"UPDATE furumusic__fed_state_like
SET local_track_id = $2 WHERE content_id = $1",
)
.bind(&content_id)
.bind(track_id)
.execute(pool)
.await?;
sqlx::query(
"UPDATE furumusic__fed_state_playlist_item
SET local_track_id = $2 WHERE content_id = $1",
)
.bind(&content_id)
.bind(track_id)
.execute(pool)
.await?;
Ok(())
}
pub async fn record_playlist_created(
pool: &sqlx::PgPool,
user_id: i64,
@@ -2664,6 +2807,31 @@ pub async fn record_web_active_transfer(
Ok(())
}
pub async fn record_web_active_takeover(
pool: &sqlx::PgPool,
user_id: i64,
previous_device_id: &str,
state: serde_json::Value,
) -> Result<()> {
ensure_web_playback_target(pool, user_id, previous_device_id).await?;
let identity = ensure_identity(pool, user_id, "").await?;
let wire = playback_state_from_browser_json(pool, state).await?;
record_local_op(
pool,
user_id,
SyncOpPayload::PlaybackCommand {
target_device_id: previous_device_id.to_string(),
command: PlaybackCommand::ActiveChanged {
active_device_id: identity.device_id,
active_device_name: identity.name,
state: wire,
},
},
)
.await?;
Ok(())
}
async fn ensure_web_playback_target(
pool: &sqlx::PgPool,
user_id: i64,
+39 -2
View File
@@ -9,10 +9,12 @@
//! or download from other peers.
//!
//! Settings are the regular admin config entries (`federation_enabled`,
//! `federation_network_id`) and apply on the fly — saving the settings
//! `federation_network_id`, `federation_save_on_listen`) and apply on the fly — saving the settings
//! starts, stops or re-joins the node without a server restart.
pub mod client;
pub mod devices;
mod receive;
mod serve;
mod storage;
@@ -51,6 +53,13 @@ struct ContentHashJob {
file_path: String,
}
#[derive(Clone)]
struct CachedArtwork {
bytes: Vec<u8>,
mime: String,
fetched_at: std::time::Instant,
}
#[derive(Debug, Clone)]
struct TransportSample {
at: String,
@@ -207,8 +216,12 @@ pub struct Federation {
data_dir: PathBuf,
database_url: std::sync::Mutex<String>,
storage_dir: std::sync::Mutex<String>,
save_on_listen: std::sync::atomic::AtomicBool,
content_cache: std::sync::Mutex<HashMap<i64, (String, String)>>,
content_pending: std::sync::Mutex<HashSet<i64>>,
prepared_cache: std::sync::Mutex<HashMap<String, (PathBuf, String)>>,
artwork_cache: std::sync::Mutex<HashMap<String, CachedArtwork>>,
download_locks: std::sync::Mutex<HashMap<String, Arc<tokio::sync::Mutex<()>>>>,
pool: tokio::sync::OnceCell<PgPool>,
running: tokio::sync::Mutex<Option<Running>>,
last_sync: std::sync::Mutex<Option<String>>,
@@ -234,8 +247,12 @@ pub fn handle() -> Arc<Federation> {
data_dir: PathBuf::from(crate::media_paths::resolve_config_path("federation")),
database_url: std::sync::Mutex::new(String::new()),
storage_dir: std::sync::Mutex::new(String::new()),
save_on_listen: std::sync::atomic::AtomicBool::new(false),
content_cache: std::sync::Mutex::new(Default::default()),
content_pending: std::sync::Mutex::new(Default::default()),
prepared_cache: std::sync::Mutex::new(Default::default()),
artwork_cache: std::sync::Mutex::new(Default::default()),
download_locks: std::sync::Mutex::new(Default::default()),
pool: tokio::sync::OnceCell::new(),
running: tokio::sync::Mutex::new(None),
last_sync: std::sync::Mutex::new(None),
@@ -285,7 +302,8 @@ impl Federation {
let mut effective = config.clone();
let rows = sqlx::query(
"SELECT key, value FROM furumusic__config_entry
WHERE key IN ('federation_enabled', 'federation_network_id', 'agent_storage_dir')",
WHERE key IN ('federation_enabled', 'federation_network_id',
'federation_save_on_listen', 'agent_storage_dir')",
)
.fetch_all(&pool)
.await
@@ -304,6 +322,11 @@ impl Federation {
}
}
"federation_network_id" => effective.federation_network_id = value,
"federation_save_on_listen" => {
if let Ok(parsed) = value.parse() {
effective.federation_save_on_listen = parsed;
}
}
"agent_storage_dir" => {
effective.agent_storage_dir = crate::media_paths::resolve_config_path(&value);
}
@@ -318,6 +341,10 @@ impl Federation {
pub async fn apply(self: &Arc<Self>, config: &AppConfig) {
*lock(&self.database_url) = config.database_url.clone();
*lock(&self.storage_dir) = config.agent_storage_dir.clone();
self.save_on_listen.store(
config.federation_save_on_listen,
std::sync::atomic::Ordering::Relaxed,
);
let network = config.federation_network_id.trim().to_string();
if config.federation_enabled && !network.is_empty() {
if let Err(err) = self.start(network, config.agent_storage_dir.clone()).await {
@@ -905,6 +932,16 @@ impl Federation {
)
.await
}
pub async fn fed_device_web_active_takeover(
&self,
user_id: i64,
previous_device_id: &str,
state: serde_json::Value,
) -> Result<()> {
let pool = self.pool().await?;
devices::record_web_active_takeover(&pool, user_id, previous_device_id, state).await
}
}
async fn persist_content_id(
+914
View File
@@ -0,0 +1,914 @@
//! Verified federated audio download and trusted materialization.
//!
//! This module never writes inbox, processing-task, or review tables. Peer
//! metadata is the authority for this import path.
use std::path::{Path, PathBuf};
use std::str::FromStr;
use anyhow::{Context, Result};
use music_dht::EndpointId;
use serde::{Deserialize, Serialize};
use sha2::{Digest as _, Sha256};
use tokio::io::{AsyncRead, AsyncReadExt, AsyncWrite, AsyncWriteExt};
use super::{Federation, now_iso};
const MAX_LINE: usize = 4096;
const MAX_AUDIO_BYTES: u64 = 4 * 1024 * 1024 * 1024;
const MAX_IMAGE_BYTES: u64 = 16 * 1024 * 1024;
const ARTWORK_CACHE_TTL: std::time::Duration = std::time::Duration::from_secs(60 * 60);
#[derive(Debug, Clone, Deserialize, Serialize)]
struct TrackMetadata {
title: String,
#[serde(default)]
artists: Vec<String>,
#[serde(default)]
featured_artists: Vec<String>,
#[serde(default)]
album_artists: Vec<String>,
#[serde(default)]
release_title: String,
release_type: Option<String>,
year: Option<i32>,
track_number: Option<i32>,
disc_number: Option<i32>,
}
#[derive(Serialize)]
struct AudioRequest<'a> {
item_id: &'a str,
offset: u64,
want_cover: bool,
}
#[derive(Deserialize)]
struct AudioHeader {
ok: bool,
#[serde(default)]
error: Option<String>,
#[serde(default)]
mime_type: String,
#[serde(default)]
total_size: u64,
#[serde(default)]
metadata: Option<TrackMetadata>,
#[serde(default)]
cover_size: u64,
#[serde(default)]
cover_mime: String,
#[serde(default)]
artist_image_size: u64,
#[serde(default)]
artist_image_mime: String,
}
pub struct PreparedTrack {
pub local_track_id: Option<i64>,
pub stream_url: String,
}
#[derive(Debug, Clone, Serialize)]
pub struct DownloadProgress {
pub phase: &'static str,
pub received: u64,
pub total: u64,
}
struct Downloaded {
path: PathBuf,
mime: String,
metadata: TrackMetadata,
cover: Option<(Vec<u8>, String)>,
artist_image: Option<(Vec<u8>, String)>,
}
impl Federation {
pub async fn discover_catalog_artwork(
&self,
artist: &str,
release: Option<&str>,
) -> Result<Option<(Vec<u8>, String)>> {
let service = self.service().await?;
let normalized = music_dht::normalize_name(artist);
let outcome = service
.search_network(artist)
.await
.map_err(|err| anyhow::anyhow!("artwork peer discovery failed: {err}"))?;
let mut owners = Vec::new();
for item in outcome.local_results.iter().chain(&outcome.network_results) {
let matches = (item.kind == music_dht::ItemKind::Artist
&& item.normalized_name == normalized)
|| item
.artist_names
.iter()
.chain(&item.featured_artist_names)
.any(|name| music_dht::normalize_name(name) == normalized);
let owner = item.owner.to_string();
if matches && !owners.contains(&owner) {
owners.push(owner);
}
}
for owner in owners {
match tokio::time::timeout(
std::time::Duration::from_secs(5),
self.catalog_artwork(&owner, artist, release),
)
.await
{
Ok(Ok(Some(artwork))) => return Ok(Some(artwork)),
Ok(Ok(None)) => {}
Ok(Err(err)) => {
tracing::debug!(%owner, %artist, ?release, "catalog artwork peer failed: {err:#}");
}
Err(_) => {
tracing::debug!(%owner, %artist, ?release, "catalog artwork peer timed out");
}
}
}
Ok(None)
}
pub async fn catalog_artwork(
&self,
owner: &str,
artist: &str,
release: Option<&str>,
) -> Result<Option<(Vec<u8>, String)>> {
let cache_key = format!("catalog:{owner}:{artist}:{}", release.unwrap_or_default());
if let Some(cached) = cached_artwork(self, &cache_key) {
return Ok(Some(cached));
}
let owner = EndpointId::from_str(owner).context("invalid federation owner")?;
let service = self.service().await?;
let mut stream = service
.open_stream(owner, super::CATALOG_ALPN)
.await
.map_err(|err| anyhow::anyhow!("cannot reach catalog peer: {err}"))?;
let mut request = serde_json::to_vec(&music_dht::catalog::CatalogRequest {
artist: artist.to_owned(),
want: Some(if release.is_some() {
"release_cover".to_owned()
} else {
"artist_image".to_owned()
}),
release: release.map(str::to_owned),
..Default::default()
})?;
request.push(b'\n');
stream.send.write_all(&request).await?;
stream.send.finish()?;
let header: music_dht::catalog::CatalogImageHeader =
serde_json::from_slice(&read_line(&mut stream.recv).await?)
.context("invalid catalog artwork response")?;
if !header.ok || header.size == 0 {
return Ok(None);
}
anyhow::ensure!(
header.size <= MAX_IMAGE_BYTES,
"catalog artwork is too large"
);
anyhow::ensure!(
header.mime_type.starts_with("image/"),
"invalid catalog artwork mime type"
);
let mut bytes = vec![0; header.size as usize];
stream.recv.read_exact(&mut bytes).await?;
let artwork = (bytes, header.mime_type);
cache_artwork(self, cache_key, &artwork);
Ok(Some(artwork))
}
pub async fn track_artwork(
&self,
owner: &str,
item_id: &str,
) -> Result<Option<(Vec<u8>, String)>> {
let cache_key = format!("{owner}:{item_id}");
if let Some(cached) = cached_artwork(self, &cache_key) {
return Ok(Some(cached));
}
let owner = EndpointId::from_str(owner).context("invalid federation owner")?;
anyhow::ensure!(item_id.len() == 64, "invalid federation item id");
let service = self.service().await?;
let mut stream = service
.open_stream(owner, super::AUDIO_ALPN)
.await
.map_err(|err| anyhow::anyhow!("cannot reach track owner: {err}"))?;
write_line(
&mut stream.send,
&AudioRequest {
item_id,
offset: 0,
want_cover: true,
},
)
.await?;
stream.send.finish()?;
let header: AudioHeader = serde_json::from_slice(&read_line(&mut stream.recv).await?)
.context("invalid artwork response")?;
anyhow::ensure!(
header.ok,
"peer refused artwork: {}",
header.error.unwrap_or_else(|| "unknown error".to_string())
);
let artwork = read_segment(
&mut stream.recv,
header.cover_size,
&header.cover_mime,
"cover",
)
.await?;
if let Some((bytes, mime)) = artwork {
let artwork = (bytes, mime);
cache_artwork(self, cache_key, &artwork);
Ok(Some(artwork))
} else {
Ok(None)
}
}
pub async fn prepare_content_with_progress<F>(
self: &std::sync::Arc<Self>,
content_id: &str,
owner: &str,
item_id: &str,
mut progress: F,
) -> Result<PreparedTrack>
where
F: FnMut(DownloadProgress) + Send,
{
progress(DownloadProgress {
phase: "checking",
received: 0,
total: 0,
});
let content_id =
music_dht::normalize_content_id(content_id).context("invalid content id")?;
let pool = self.pool().await?;
let token = content_id.trim_start_matches("b3:").to_owned();
let download_lock = {
let mut locks = super::lock(&self.download_locks);
std::sync::Arc::clone(
locks
.entry(content_id.clone())
.or_insert_with(|| std::sync::Arc::new(tokio::sync::Mutex::new(()))),
)
};
let _download_guard = download_lock.lock().await;
if let Some(track_id) = local_track_id(&pool, &content_id).await? {
progress(DownloadProgress {
phase: "ready",
received: 1,
total: 1,
});
return Ok(PreparedTrack {
local_track_id: Some(track_id),
stream_url: format!("/api/player/stream/{track_id}"),
});
}
let save = self
.save_on_listen
.load(std::sync::atomic::Ordering::Relaxed);
if !save
&& let Some((_path, _mime)) = super::lock(&self.prepared_cache).get(&token).cloned()
{
return Ok(PreparedTrack {
local_track_id: None,
stream_url: format!("/api/player/federation/cache/{token}"),
});
}
let owner = EndpointId::from_str(owner).context("invalid federation owner")?;
anyhow::ensure!(item_id.len() == 64, "invalid federation item id");
let service = self.service().await?;
let storage_root = if save {
PathBuf::from(super::lock(&self.storage_dir).clone())
} else {
PathBuf::from(crate::media_paths::resolve_config_path("federation-cache"))
};
anyhow::ensure!(
!storage_root.as_os_str().is_empty(),
"media storage directory is not configured"
);
let dir = storage_root.join("federation");
tokio::fs::create_dir_all(&dir).await?;
let downloaded =
download(&service, owner, item_id, &content_id, &dir, &mut progress).await?;
if !save {
super::lock(&self.prepared_cache)
.insert(token.clone(), (downloaded.path, downloaded.mime));
return Ok(PreparedTrack {
local_track_id: None,
stream_url: format!("/api/player/federation/cache/{token}"),
});
}
progress(DownloadProgress {
phase: "saving",
received: 1,
total: 1,
});
let track_id = materialize(&pool, &storage_root, &content_id, downloaded).await?;
// The normal periodic sync will publish it; this immediate sync keeps
// save-on-listen useful to the federation without waiting a minute.
if let Err(err) = self.sync_now().await {
tracing::warn!(track_id, "post-import federation publish failed: {err:#}");
}
Ok(PreparedTrack {
local_track_id: Some(track_id),
stream_url: format!("/api/player/stream/{track_id}"),
})
}
pub fn prepared_cache_file(&self, token: &str) -> Option<(PathBuf, String)> {
if token.len() != 64 || !token.bytes().all(|byte| byte.is_ascii_hexdigit()) {
return None;
}
super::lock(&self.prepared_cache).get(token).cloned()
}
}
fn cached_artwork(federation: &Federation, key: &str) -> Option<(Vec<u8>, String)> {
let mut cache = super::lock(&federation.artwork_cache);
let cached = cache.get(key).cloned()?;
if cached.fetched_at.elapsed() > ARTWORK_CACHE_TTL {
cache.remove(key);
return None;
}
Some((cached.bytes, cached.mime))
}
fn cache_artwork(federation: &Federation, key: String, artwork: &(Vec<u8>, String)) {
let mut cache = super::lock(&federation.artwork_cache);
if cache.len() >= 512 {
cache.retain(|_, value| value.fetched_at.elapsed() <= ARTWORK_CACHE_TTL);
if cache.len() >= 512
&& let Some(oldest) = cache
.iter()
.min_by_key(|(_, value)| value.fetched_at)
.map(|(key, _)| key.clone())
{
cache.remove(&oldest);
}
}
cache.insert(
key,
super::CachedArtwork {
bytes: artwork.0.clone(),
mime: artwork.1.clone(),
fetched_at: std::time::Instant::now(),
},
);
}
async fn download(
service: &music_dht::MusicDhtService,
owner: EndpointId,
item_id: &str,
content_id: &str,
dir: &Path,
progress: &mut (impl FnMut(DownloadProgress) + Send),
) -> Result<Downloaded> {
progress(DownloadProgress {
phase: "connecting",
received: 0,
total: 0,
});
let mut stream = service
.open_stream(owner, super::AUDIO_ALPN)
.await
.map_err(|err| anyhow::anyhow!("cannot reach track owner: {err}"))?;
write_line(
&mut stream.send,
&AudioRequest {
item_id,
offset: 0,
want_cover: true,
},
)
.await?;
stream.send.finish()?;
let header: AudioHeader = serde_json::from_slice(&read_line(&mut stream.recv).await?)
.context("invalid audio response")?;
anyhow::ensure!(
header.ok,
"peer refused audio: {}",
header.error.unwrap_or_else(|| "unknown error".to_string())
);
anyhow::ensure!(
header.total_size > 0 && header.total_size <= MAX_AUDIO_BYTES,
"invalid federated audio size"
);
progress(DownloadProgress {
phase: "downloading",
received: 0,
total: header.total_size,
});
let metadata = header.metadata.context("peer returned no track metadata")?;
let cover = read_segment(
&mut stream.recv,
header.cover_size,
&header.cover_mime,
"cover",
)
.await?;
let artist_image = read_segment(
&mut stream.recv,
header.artist_image_size,
&header.artist_image_mime,
"artist image",
)
.await?;
let stem = content_id.trim_start_matches("b3:");
let extension = audio_extension(&header.mime_type);
let final_path = dir.join(format!("{stem}.{extension}"));
let part_path = dir.join(format!(".{stem}.{extension}.part"));
let mut file = tokio::fs::File::create(&part_path).await?;
let mut hasher = blake3::Hasher::new();
let mut received = 0u64;
let mut buf = vec![0u8; 64 * 1024];
while received < header.total_size {
let remaining = (header.total_size - received).min(buf.len() as u64) as usize;
let count = stream.recv.read(&mut buf[..remaining]).await?.unwrap_or(0);
anyhow::ensure!(count > 0, "audio stream ended early");
file.write_all(&buf[..count]).await?;
hasher.update(&buf[..count]);
received += count as u64;
progress(DownloadProgress {
phase: "downloading",
received,
total: header.total_size,
});
}
file.flush().await?;
drop(file);
let actual = format!("b3:{}", hasher.finalize().to_hex());
progress(DownloadProgress {
phase: "verifying",
received,
total: header.total_size,
});
if actual != content_id {
let _ = tokio::fs::remove_file(&part_path).await;
anyhow::bail!("downloaded audio content id mismatch");
}
tokio::fs::rename(&part_path, &final_path).await?;
Ok(Downloaded {
path: final_path,
mime: header.mime_type,
metadata,
cover,
artist_image,
})
}
async fn materialize(
pool: &sqlx::PgPool,
storage_root: &Path,
content_id: &str,
downloaded: Downloaded,
) -> Result<i64> {
if let Some(track_id) = local_track_id(pool, content_id).await? {
return Ok(track_id);
}
let bytes = tokio::fs::read(&downloaded.path).await?;
let sha256 = format!("{:x}", Sha256::digest(&bytes));
let relative = downloaded
.path
.strip_prefix(storage_root)
.unwrap_or(&downloaded.path)
.to_string_lossy()
.into_owned();
let mut tx = pool.begin().await?;
let existing_media: Option<i64> = sqlx::query_scalar(
"SELECT id FROM furumusic__media_file
WHERE file_type = 'audio' AND sha256_hash = $1 LIMIT 1",
)
.bind(&sha256)
.fetch_optional(&mut *tx)
.await?;
let media_id = match existing_media {
Some(id) => id,
None => {
sqlx::query_scalar(
"INSERT INTO furumusic__media_file
(file_type, file_path, original_filename, mime_type,
file_size_bytes, sha256_hash, audio_format,
uploaded_by_user_id, uploader_name, created_at)
VALUES ('audio', $1, $2, $3, $4, $5, $6, NULL, 'Federation', $7)
RETURNING id",
)
.bind(&relative)
.bind(
downloaded
.path
.file_name()
.and_then(|name| name.to_str())
.unwrap_or("federated-audio"),
)
.bind(&downloaded.mime)
.bind(bytes.len() as i64)
.bind(&sha256)
.bind(audio_extension(&downloaded.mime))
.bind(now_iso())
.fetch_one(&mut *tx)
.await?
}
};
sqlx::query(
"INSERT INTO furumusic__federation_content_id_cache
(media_file_id, sha256_hash, content_id, updated_at)
VALUES ($1, $2, $3, $4)
ON CONFLICT (media_file_id) DO UPDATE SET
sha256_hash = EXCLUDED.sha256_hash,
content_id = EXCLUDED.content_id,
updated_at = EXCLUDED.updated_at",
)
.bind(media_id)
.bind(&sha256)
.bind(content_id)
.bind(now_iso())
.execute(&mut *tx)
.await?;
if let Some(track_id) = sqlx::query_scalar::<_, i64>(
"SELECT id FROM furumusic__track WHERE audio_file_id = $1 LIMIT 1",
)
.bind(media_id)
.fetch_optional(&mut *tx)
.await?
{
tx.commit().await?;
return Ok(track_id);
}
let release_title = nonempty(&downloaded.metadata.release_title).unwrap_or("Unknown release");
let release_sort = music_dht::normalize_name(release_title);
let release_id: i64 = if let Some(id) = sqlx::query_scalar(
"SELECT id FROM furumusic__release
WHERE title_sort = $1 AND year IS NOT DISTINCT FROM $2
ORDER BY id LIMIT 1",
)
.bind(&release_sort)
.bind(downloaded.metadata.year)
.fetch_optional(&mut *tx)
.await?
{
id
} else {
sqlx::query_scalar(
"INSERT INTO furumusic__release
(title, title_sort, release_type, year, is_hidden, model_name,
created_at, updated_at)
VALUES ($1, $2, $3, $4, false, NULL, $5, $5)
RETURNING id",
)
.bind(release_title)
.bind(&release_sort)
.bind(
downloaded
.metadata
.release_type
.as_deref()
.unwrap_or("album"),
)
.bind(downloaded.metadata.year)
.bind(now_iso())
.fetch_one(&mut *tx)
.await?
};
let duration: f64 = sqlx::query_scalar(
"SELECT COALESCE(duration_seconds, 0)
FROM furumusic__track_ref WHERE content_id = $1",
)
.bind(content_id)
.fetch_optional(&mut *tx)
.await?
.unwrap_or(0.0);
let track_id: i64 = sqlx::query_scalar(
"INSERT INTO furumusic__track
(title, title_sort, release_id, track_number, disc_number,
duration_seconds, audio_file_id, year, is_hidden, model_name,
created_at, updated_at)
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, false, NULL, $9, $9)
RETURNING id",
)
.bind(&downloaded.metadata.title)
.bind(music_dht::normalize_name(&downloaded.metadata.title))
.bind(release_id)
.bind(downloaded.metadata.track_number)
.bind(downloaded.metadata.disc_number)
.bind(duration)
.bind(media_id)
.bind(downloaded.metadata.year)
.bind(now_iso())
.fetch_one(&mut *tx)
.await?;
let main_artists = if downloaded.metadata.artists.is_empty() {
&downloaded.metadata.album_artists
} else {
&downloaded.metadata.artists
};
let mut main_artist_ids = Vec::new();
for (position, name) in main_artists.iter().enumerate() {
let artist_id = ensure_artist(&mut tx, name).await?;
main_artist_ids.push(artist_id);
link_track_artist(&mut tx, track_id, artist_id, "main", position as i32).await?;
link_release_artist(&mut tx, release_id, artist_id, position as i32).await?;
}
for (position, name) in downloaded.metadata.featured_artists.iter().enumerate() {
let artist_id = ensure_artist(&mut tx, name).await?;
link_track_artist(&mut tx, track_id, artist_id, "featuring", position as i32).await?;
}
sqlx::query(
"UPDATE furumusic__track_ref
SET local_track_id = $2, metadata_authority = 'federation', updated_at = $3
WHERE content_id = $1",
)
.bind(content_id)
.bind(track_id)
.bind(now_iso())
.execute(&mut *tx)
.await?;
sqlx::query(
"UPDATE furumusic__fed_state_like
SET local_track_id = $2 WHERE content_id = $1",
)
.bind(content_id)
.bind(track_id)
.execute(&mut *tx)
.await?;
sqlx::query(
"UPDATE furumusic__fed_state_playlist_item
SET local_track_id = $2 WHERE content_id = $1",
)
.bind(content_id)
.bind(track_id)
.execute(&mut *tx)
.await?;
tx.commit().await?;
super::devices::materialize_content_state(pool, content_id, track_id).await?;
// Artwork is non-authoritative for identity and may be installed after
// the audio transaction. Failure does not invalidate a verified track.
if let Some((bytes, mime)) = downloaded.cover
&& let Err(err) =
install_release_artwork(pool, storage_root, release_id, &bytes, &mime).await
{
tracing::warn!(
release_id,
"failed to install federated release artwork: {err:#}"
);
}
if let Some((bytes, mime)) = downloaded.artist_image
&& let Err(err) =
install_artist_artwork(pool, storage_root, &main_artist_ids, &bytes, &mime).await
{
tracing::warn!("failed to install federated artist artwork: {err:#}");
}
Ok(track_id)
}
async fn install_release_artwork(
pool: &sqlx::PgPool,
storage_root: &Path,
release_id: i64,
bytes: &[u8],
mime: &str,
) -> Result<()> {
let media_id = persist_artwork(pool, storage_root, bytes, mime).await?;
sqlx::query(
"UPDATE furumusic__release
SET cover_file_id = $1, updated_at = $3
WHERE id = $2 AND cover_file_id IS NULL",
)
.bind(media_id)
.bind(release_id)
.bind(now_iso())
.execute(pool)
.await?;
Ok(())
}
async fn install_artist_artwork(
pool: &sqlx::PgPool,
storage_root: &Path,
artist_ids: &[i64],
bytes: &[u8],
mime: &str,
) -> Result<()> {
if artist_ids.is_empty() {
return Ok(());
}
let media_id = persist_artwork(pool, storage_root, bytes, mime).await?;
sqlx::query(
"UPDATE furumusic__artist
SET image_file_id = $1, updated_at = $3
WHERE id = ANY($2) AND image_file_id IS NULL",
)
.bind(media_id)
.bind(artist_ids)
.bind(now_iso())
.execute(pool)
.await?;
Ok(())
}
async fn persist_artwork(
pool: &sqlx::PgPool,
storage_root: &Path,
bytes: &[u8],
mime: &str,
) -> Result<i64> {
anyhow::ensure!(!bytes.is_empty() && bytes.len() as u64 <= MAX_IMAGE_BYTES);
anyhow::ensure!(mime.starts_with("image/"), "invalid artwork mime type");
let hash = format!("{:x}", Sha256::digest(bytes));
if let Some(id) = sqlx::query_scalar(
"SELECT id FROM furumusic__media_file
WHERE file_type = 'cover_art' AND sha256_hash = $1 LIMIT 1",
)
.bind(&hash)
.fetch_optional(pool)
.await?
{
return Ok(id);
}
let extension = image_extension(mime);
let filename = format!("federation-artwork-{}.{}", &hash[..16], extension);
let dir = storage_root.join("federation").join("artwork");
tokio::fs::create_dir_all(&dir).await?;
let path = dir.join(&filename);
tokio::fs::write(&path, bytes).await?;
let relative = path
.strip_prefix(storage_root)
.unwrap_or(&path)
.to_string_lossy()
.into_owned();
let id = sqlx::query_scalar(
"INSERT INTO furumusic__media_file
(file_type, file_path, original_filename, mime_type,
file_size_bytes, sha256_hash, uploaded_by_user_id,
uploader_name, created_at)
VALUES ('cover_art', $1, $2, $3, $4, $5, NULL, 'Federation', $6)
RETURNING id",
)
.bind(&relative)
.bind(&filename)
.bind(mime)
.bind(bytes.len() as i64)
.bind(&hash)
.bind(now_iso())
.fetch_one(pool)
.await?;
if let Err(err) = crate::agent::cover_variants::ensure_cover_variants(&path).await {
tracing::warn!(
media_id = id,
"failed to generate federated artwork variants: {err}"
);
}
Ok(id)
}
fn image_extension(mime: &str) -> &'static str {
match mime {
"image/png" => "png",
"image/webp" => "webp",
"image/gif" => "gif",
"image/avif" => "avif",
_ => "jpg",
}
}
async fn ensure_artist(tx: &mut sqlx::Transaction<'_, sqlx::Postgres>, name: &str) -> Result<i64> {
let normalized = music_dht::normalize_name(name);
if let Some(id) = sqlx::query_scalar(
"SELECT id FROM furumusic__artist WHERE name_sort = $1 ORDER BY id LIMIT 1",
)
.bind(&normalized)
.fetch_optional(&mut **tx)
.await?
{
return Ok(id);
}
Ok(sqlx::query_scalar(
"INSERT INTO furumusic__artist
(name, name_sort, is_hidden, model_name, created_at, updated_at)
VALUES ($1, $2, false, NULL, $3, $3) RETURNING id",
)
.bind(name)
.bind(normalized)
.bind(now_iso())
.fetch_one(&mut **tx)
.await?)
}
async fn link_track_artist(
tx: &mut sqlx::Transaction<'_, sqlx::Postgres>,
track_id: i64,
artist_id: i64,
role: &str,
position: i32,
) -> Result<()> {
sqlx::query(
"INSERT INTO furumusic__track_artist
(track_id, artist_id, role, position) VALUES ($1, $2, $3, $4)",
)
.bind(track_id)
.bind(artist_id)
.bind(role)
.bind(position)
.execute(&mut **tx)
.await?;
Ok(())
}
async fn link_release_artist(
tx: &mut sqlx::Transaction<'_, sqlx::Postgres>,
release_id: i64,
artist_id: i64,
position: i32,
) -> Result<()> {
sqlx::query(
"INSERT INTO furumusic__release_artist
(release_id, artist_id, position)
SELECT $1, $2, $3 WHERE NOT EXISTS (
SELECT 1 FROM furumusic__release_artist
WHERE release_id = $1 AND artist_id = $2
)",
)
.bind(release_id)
.bind(artist_id)
.bind(position)
.execute(&mut **tx)
.await?;
Ok(())
}
async fn local_track_id(pool: &sqlx::PgPool, content_id: &str) -> Result<Option<i64>> {
Ok(sqlx::query_scalar(
"SELECT t.id
FROM furumusic__federation_content_id_cache c
JOIN furumusic__track t ON t.audio_file_id = c.media_file_id
WHERE c.content_id = $1 AND t.is_hidden = false LIMIT 1",
)
.bind(content_id)
.fetch_optional(pool)
.await?)
}
async fn read_line<R: AsyncRead + Unpin>(reader: &mut R) -> Result<Vec<u8>> {
let mut line = Vec::new();
let mut byte = [0u8; 1];
loop {
anyhow::ensure!(
reader.read_exact(&mut byte).await.is_ok(),
"stream ended early"
);
if byte[0] == b'\n' {
return Ok(line);
}
line.push(byte[0]);
anyhow::ensure!(line.len() <= MAX_LINE, "protocol line too large");
}
}
async fn write_line<W: AsyncWrite + Unpin>(writer: &mut W, value: &impl Serialize) -> Result<()> {
let mut line = serde_json::to_vec(value)?;
line.push(b'\n');
writer.write_all(&line).await?;
Ok(())
}
async fn read_segment<R: AsyncRead + Unpin>(
reader: &mut R,
size: u64,
mime: &str,
label: &str,
) -> Result<Option<(Vec<u8>, String)>> {
if size == 0 {
return Ok(None);
}
anyhow::ensure!(size <= MAX_IMAGE_BYTES, "{label} is too large");
let mut bytes = vec![0u8; size as usize];
reader.read_exact(&mut bytes).await?;
Ok(Some((bytes, mime.to_owned())))
}
fn audio_extension(mime: &str) -> &'static str {
match mime {
"audio/mpeg" => "mp3",
"audio/flac" => "flac",
"audio/ogg" => "ogg",
"audio/opus" => "opus",
"audio/wav" => "wav",
"audio/mp4" => "m4a",
"audio/aac" => "aac",
_ => "bin",
}
}
fn nonempty(value: &str) -> Option<&str> {
(!value.trim().is_empty()).then_some(value.trim())
}