Compare commits

...
4 Commits
Author SHA1 Message Date
Ultradesu 42c772f735 Fixed blake3 computation
Build and Publish / Build and Publish Docker Image (push) Successful in 5m5s
2026-07-20 18:24:39 +03:00
Ultradesu e738086573 Extend DHT scheme with content_id
Build and Publish / Build and Publish Docker Image (push) Successful in 5m6s
2026-07-20 18:05:31 +03:00
Ultradesu 4b7756c36e Extend DHT scheme with content_id 2026-07-20 18:05:16 +03:00
Ultradesu 4381750c6e Extend DHT scheme
Build and Publish / Build and Publish Docker Image (push) Successful in 5m5s
2026-07-20 16:29:21 +03:00
5 changed files with 198 additions and 49 deletions
Generated
+10 -9
View File
@@ -690,9 +690,9 @@ dependencies = [
[[package]]
name = "clap"
version = "4.6.2"
version = "4.6.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "dd059f9da4f5c36b3787f65d38ccaab1cc315f07b01f89abc8359ee6a8205011"
checksum = "0fb99565819980999fb7b4a1796046a5c949e6d4ff132cf5fadf5a641e20d776"
dependencies = [
"clap_builder",
"clap_derive",
@@ -712,9 +712,9 @@ dependencies = [
[[package]]
name = "clap_derive"
version = "4.6.1"
version = "4.6.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "f2ce8604710f6733aa641a2b3731eaa1e8b3d9973d5e3565da11800813f997a9"
checksum = "32f2392eae7f16557a3d727ef3a12e57b2b2ca6f98566a5f4fb41ffe305df077"
dependencies = [
"heck",
"proc-macro2",
@@ -1717,7 +1717,7 @@ dependencies = [
[[package]]
name = "federation-net"
version = "0.1.0"
source = "git+https://gt.hexor.cy/ab/frid.git#3f15e9aebd574f7dd922e623e724eb7fde741b59"
source = "git+https://gt.hexor.cy/ab/frid.git#8ee1db9cf89ea604c6b8a8c3fe089a714c1e321f"
dependencies = [
"blake3",
"data-encoding",
@@ -1844,11 +1844,12 @@ checksum = "e6d5a32815ae3f33302d95fdcb2ce17862f8c65363dcfd29360480ba1001fc9c"
[[package]]
name = "furumusic"
version = "0.6.2-fd"
version = "0.6.4-fd"
dependencies = [
"anyhow",
"async-trait",
"base64 0.22.1",
"blake3",
"chrono",
"cot",
"croner",
@@ -2470,9 +2471,9 @@ dependencies = [
[[package]]
name = "hyper"
version = "1.10.1"
version = "1.11.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "55281c53a1894c864990125767da440a4e630446785086f52523b20033b74498"
checksum = "d22053281f852e11534f5198498373cbb59295120a20771d90f7ed1897490a72"
dependencies = [
"atomic-waker",
"bytes",
@@ -3621,7 +3622,7 @@ dependencies = [
[[package]]
name = "music-dht"
version = "0.1.0"
source = "git+https://gt.hexor.cy/ab/frid.git#3f15e9aebd574f7dd922e623e724eb7fde741b59"
source = "git+https://gt.hexor.cy/ab/frid.git#8ee1db9cf89ea604c6b8a8c3fe089a714c1e321f"
dependencies = [
"async-trait",
"blake3",
+2 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "furumusic"
version = "0.6.2-fd"
version = "0.6.5-fd"
edition = "2024"
description = "Reusable web-app boilerplate: auth, OIDC/SSO, admin panel, user management, i18n, PostgreSQL"
@@ -16,6 +16,7 @@ reqwest = { version = "0.12", default-features = false, features = ["rustls-tls"
tokio = { version = "1", features = ["sync", "fs", "io-util"] }
tower = "0.5"
base64 = "0.22"
blake3 = "1"
serde_json = "1"
tracing = "0.1"
tracing-subscriber = { version = "0.3", features = ["env-filter"] }
+115 -10
View File
@@ -15,6 +15,7 @@
mod serve;
mod storage;
use std::collections::{HashMap, HashSet};
use std::path::PathBuf;
use std::sync::{Arc, OnceLock};
use std::time::Duration;
@@ -35,6 +36,9 @@ pub use serve::{AUDIO_ALPN, CATALOG_ALPN};
/// How often the published library is re-synchronized with the database.
const SYNC_INTERVAL: Duration = Duration::from_secs(60);
/// Content IDs are file hashes and can be expensive to warm for large
/// libraries. Keep DHT publishing responsive and let later syncs fill them in.
const MAX_CONTENT_HASH_JOBS_PER_SYNC: usize = 512;
struct Running {
service: Arc<MusicDhtService>,
@@ -46,6 +50,9 @@ pub struct Federation {
/// Transport data directory; server-side DHT state and identity live in PostgreSQL.
data_dir: PathBuf,
database_url: std::sync::Mutex<String>,
storage_dir: std::sync::Mutex<String>,
content_cache: std::sync::Mutex<HashMap<i64, (String, String)>>,
content_pending: std::sync::Mutex<HashSet<i64>>,
pool: tokio::sync::OnceCell<PgPool>,
running: tokio::sync::Mutex<Option<Running>>,
last_sync: std::sync::Mutex<Option<String>>,
@@ -69,6 +76,9 @@ pub fn handle() -> Arc<Federation> {
Arc::new(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()),
content_cache: std::sync::Mutex::new(Default::default()),
content_pending: std::sync::Mutex::new(Default::default()),
pool: tokio::sync::OnceCell::new(),
running: tokio::sync::Mutex::new(None),
last_sync: std::sync::Mutex::new(None),
@@ -149,6 +159,7 @@ impl Federation {
/// node. Called at boot and every time the admin settings are saved.
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();
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 {
@@ -281,7 +292,7 @@ impl Federation {
Ok(())
}
async fn sync_once(&self, service: &MusicDhtService) -> Result<SyncStats> {
async fn sync_once(self: &Arc<Self>, service: &MusicDhtService) -> Result<SyncStats> {
let specs = match self.collect_specs().await {
Ok(specs) => specs,
Err(err) => {
@@ -341,7 +352,7 @@ impl Federation {
/// Everything the regular player shows, as DHT item specs: non-hidden
/// artists, releases and tracks (a track also hides with its release).
async fn collect_specs(&self) -> Result<Vec<ItemSpec>> {
async fn collect_specs(self: &Arc<Self>) -> Result<Vec<ItemSpec>> {
let pool = self.pool().await?;
let mut specs = Vec::new();
@@ -355,9 +366,14 @@ impl Federation {
kind: ItemKind::Artist,
name: row.get(1),
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,
});
}
@@ -389,52 +405,130 @@ impl Federation {
kind: ItemKind::Release,
name: row.get(1),
artist_names: artists_of_release.remove(&id).unwrap_or_default(),
featured_artist_names: Vec::new(),
year: row.get(2),
release_type: row.get(3),
release_title: None,
track_number: None,
disc_number: None,
duration_seconds: None,
content_id: None,
});
}
let track_artists = sqlx::query(
"SELECT ta.track_id, a.name FROM furumusic__track_artist ta
"SELECT ta.track_id, a.name, ta.role FROM furumusic__track_artist ta
JOIN furumusic__artist a ON a.id = ta.artist_id
WHERE ta.role IN ('main', 'featuring')
ORDER BY ta.track_id, ta.position",
ORDER BY ta.track_id,
CASE ta.role WHEN 'main' THEN 0 ELSE 1 END,
ta.position",
)
.fetch_all(&pool)
.await?;
let mut artists_of_track: std::collections::HashMap<i64, Vec<String>> = Default::default();
let mut featured_of_track: std::collections::HashMap<i64, Vec<String>> = Default::default();
for row in &track_artists {
artists_of_track
.entry(row.get(0))
.or_default()
.push(row.get(1));
let id: i64 = row.get(0);
let name: String = row.get(1);
if row.get::<String, _>(2) == "featuring" {
featured_of_track.entry(id).or_default().push(name);
} else {
artists_of_track.entry(id).or_default().push(name);
}
}
let tracks = sqlx::query(
"SELECT t.id, t.title, COALESCE(t.year, r.year), t.duration_seconds
"SELECT t.id, t.title, COALESCE(t.year, r.year), t.duration_seconds,
r.title, r.release_type, t.track_number, t.disc_number,
t.audio_file_id, m.file_path, m.sha256_hash
FROM furumusic__track t
JOIN furumusic__release r ON r.id = t.release_id
JOIN furumusic__media_file m ON m.id = t.audio_file_id
WHERE t.is_hidden = false AND r.is_hidden = false",
)
.fetch_all(&pool)
.await?;
let storage_dir = lock(&self.storage_dir).clone();
let mut hash_jobs_scheduled = 0usize;
for row in &tracks {
let id: i64 = row.get(0);
let duration: f64 = row.get(3);
let media_file_id: i64 = row.get(8);
let file_path: String = row.get(9);
let sha256_hash: String = row.get(10);
let (content_id, scheduled_hash) = self.content_id_for_media(
media_file_id,
sha256_hash,
file_path,
&storage_dir,
hash_jobs_scheduled < MAX_CONTENT_HASH_JOBS_PER_SYNC,
);
if scheduled_hash {
hash_jobs_scheduled += 1;
}
specs.push(ItemSpec {
local_key: format!("track:{id}"),
kind: ItemKind::Track,
name: row.get(1),
artist_names: artists_of_track.remove(&id).unwrap_or_default(),
featured_artist_names: featured_of_track.remove(&id).unwrap_or_default(),
year: row.get(2),
release_type: None,
release_type: row.get(5),
release_title: Some(row.get(4)),
track_number: row.get(6),
disc_number: row.get(7),
duration_seconds: (duration > 0.0).then_some(duration),
content_id,
});
}
Ok(specs)
}
fn content_id_for_media(
self: &Arc<Self>,
media_file_id: i64,
sha256_hash: String,
file_path: String,
storage_dir: &str,
schedule_missing: bool,
) -> (Option<String>, bool) {
if storage_dir.trim().is_empty() {
return (None, false);
}
if let Some((cached_hash, content_id)) = lock(&self.content_cache).get(&media_file_id)
&& cached_hash == &sha256_hash
{
return (Some(content_id.clone()), false);
}
if !schedule_missing {
return (None, false);
}
{
let mut pending = lock(&self.content_pending);
if !pending.insert(media_file_id) {
return (None, false);
}
}
let storage_dir = storage_dir.to_string();
let fed = Arc::clone(self);
tokio::spawn(async move {
let content_id =
tokio::task::spawn_blocking(move || audio_content_id(&storage_dir, &file_path))
.await
.ok()
.flatten();
lock(&fed.content_pending).remove(&media_file_id);
if let Some(content_id) = content_id {
lock(&fed.content_cache).insert(media_file_id, (sha256_hash, content_id));
}
});
(None, true)
}
/// Live status for the admin page.
pub async fn status(&self) -> Value {
let guard = self.running.lock().await;
@@ -492,6 +586,17 @@ impl Federation {
}
}
fn audio_content_id(storage_dir: &str, file_path: &str) -> Option<String> {
if storage_dir.trim().is_empty() {
return None;
}
let path = crate::media_paths::resolve_media_file_path(storage_dir, file_path);
let mut file = std::fs::File::open(path).ok()?;
let mut hasher = blake3::Hasher::new();
std::io::copy(&mut file, &mut hasher).ok()?;
Some(format!("b3:{}", hasher.finalize().to_hex()))
}
async fn stop_running(running: Option<Running>) {
let Some(running) = running else { return };
for task in &running.tasks {
+70 -28
View File
@@ -101,6 +101,8 @@ struct CatalogTrack {
track_number: Option<i32>,
disc_number: Option<i32>,
duration_seconds: Option<f64>,
#[serde(skip_serializing_if = "Option::is_none")]
content_id: Option<String>,
item_id: String,
}
@@ -150,7 +152,10 @@ async fn read_line<R: AsyncRead + Unpin>(reader: &mut R) -> Result<Vec<u8>> {
}
}
async fn write_line<W: AsyncWriteExt + Unpin>(writer: &mut W, value: &impl Serialize) -> Result<()> {
async fn write_line<W: AsyncWriteExt + 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?;
@@ -186,7 +191,10 @@ fn guess_mime(path: &Path) -> &'static str {
}
/// Reads an image media file from disk, bounded by [`MAX_IMAGE_BYTES`].
async fn read_image(storage_dir: &str, media: Option<(String, String)>) -> Option<(Vec<u8>, String)> {
async fn read_image(
storage_dir: &str,
media: Option<(String, String)>,
) -> Option<(Vec<u8>, String)> {
let (file_path, mime) = media?;
let path = resolve_media_path(storage_dir, &file_path);
let size = tokio::fs::metadata(&path).await.ok()?.len();
@@ -367,7 +375,9 @@ async fn serve_audio_one(
Some(item_id) => match resolve_track_id(&pool, &own, item_id).await {
Ok(Some(track_id)) => track_id,
Ok(None) => return refuse_audio(stream, "track not found in the library").await,
Err(err) => return refuse_audio(stream, &format!("library lookup failed: {err:#}")).await,
Err(err) => {
return refuse_audio(stream, &format!("library lookup failed: {err:#}")).await;
}
},
None => return refuse_audio(stream, "malformed item_id").await,
};
@@ -378,7 +388,9 @@ async fn serve_audio_one(
let path = resolve_media_path(&storage_dir, &file_path);
let mut file = match tokio::fs::File::open(&path).await {
Ok(file) => file,
Err(err) => return refuse_audio(stream, &format!("audio file is not readable: {err}")).await,
Err(err) => {
return refuse_audio(stream, &format!("audio file is not readable: {err}")).await;
}
};
let total_size = file.metadata().await?.len();
let offset = request.offset.min(total_size);
@@ -395,10 +407,17 @@ async fn serve_audio_one(
};
let (cover, artist_image) = if request.want_cover {
(
read_image(&storage_dir, track_cover_file(&pool, track_id).await.ok().flatten()).await,
read_image(
&storage_dir,
track_artist_image_file(&pool, track_id).await.ok().flatten(),
track_cover_file(&pool, track_id).await.ok().flatten(),
)
.await,
read_image(
&storage_dir,
track_artist_image_file(&pool, track_id)
.await
.ok()
.flatten(),
)
.await,
)
@@ -503,7 +522,7 @@ async fn serve_catalog_one(
match request.want.as_deref() {
None | Some("catalog") => {
let response = match build_catalog(&pool, &own, &request.artist).await {
let response = match build_catalog(&pool, &own, &storage_dir, &request.artist).await {
Ok(response) => response,
Err(err) => CatalogResponse {
ok: false,
@@ -511,7 +530,10 @@ async fn serve_catalog_one(
artist: None,
},
};
stream.send.write_all(&serde_json::to_vec(&response)?).await?;
stream
.send
.write_all(&serde_json::to_vec(&response)?)
.await?;
}
Some(want @ ("artist_image" | "release_cover")) => {
let media = if want == "release_cover" {
@@ -549,7 +571,10 @@ async fn serve_catalog_one(
error: Some(format!("unknown request kind '{other}'")),
artist: None,
};
stream.send.write_all(&serde_json::to_vec(&response)?).await?;
stream
.send
.write_all(&serde_json::to_vec(&response)?)
.await?;
}
}
stream.send.finish()?;
@@ -557,7 +582,12 @@ async fn serve_catalog_one(
Ok(())
}
async fn build_catalog(pool: &PgPool, own: &EndpointId, artist: &str) -> Result<CatalogResponse> {
async fn build_catalog(
pool: &PgPool,
own: &EndpointId,
storage_dir: &str,
artist: &str,
) -> Result<CatalogResponse> {
let Some(artist_row) = sqlx::query(
"SELECT id, name FROM furumusic__artist
WHERE LOWER(name) = LOWER($1) AND is_hidden = false
@@ -589,32 +619,36 @@ async fn build_catalog(pool: &PgPool, own: &EndpointId, artist: &str) -> Result<
for release_row in release_rows {
let release_id: i64 = release_row.get(0);
let track_rows = sqlx::query(
"SELECT id, title, track_number, disc_number, duration_seconds
FROM furumusic__track
WHERE release_id = $1 AND is_hidden = false
ORDER BY disc_number NULLS FIRST, track_number NULLS LAST, title",
"SELECT t.id, t.title, t.track_number, t.disc_number, t.duration_seconds,
m.file_path
FROM furumusic__track t
JOIN furumusic__media_file m ON m.id = t.audio_file_id
WHERE t.release_id = $1 AND t.is_hidden = false
ORDER BY t.disc_number NULLS FIRST, t.track_number NULLS LAST, t.title",
)
.bind(release_id)
.fetch_all(pool)
.await?;
let mut tracks = Vec::with_capacity(track_rows.len());
for row in track_rows {
let track_id: i64 = row.get(0);
let duration: f64 = row.get(4);
let file_path: String = row.get(5);
let content_id = catalog_content_id(storage_dir, file_path).await;
tracks.push(CatalogTrack {
title: row.get(1),
track_number: row.get(2),
disc_number: row.get(3),
duration_seconds: (duration > 0.0).then_some(duration),
content_id,
item_id: item_id_of(own, track_id),
});
}
releases.push(CatalogRelease {
title: release_row.get(1),
release_type: release_row.get(2),
year: release_row.get(3),
tracks: track_rows
.into_iter()
.map(|row| {
let track_id: i64 = row.get(0);
let duration: f64 = row.get(4);
CatalogTrack {
title: row.get(1),
track_number: row.get(2),
disc_number: row.get(3),
duration_seconds: (duration > 0.0).then_some(duration),
item_id: item_id_of(own, track_id),
}
})
.collect(),
tracks,
});
}
@@ -628,6 +662,14 @@ async fn build_catalog(pool: &PgPool, own: &EndpointId, artist: &str) -> Result<
})
}
async fn catalog_content_id(storage_dir: &str, file_path: String) -> Option<String> {
let storage_dir = storage_dir.to_string();
tokio::task::spawn_blocking(move || super::audio_content_id(&storage_dir, &file_path))
.await
.ok()
.flatten()
}
async fn artist_image_by_name(pool: &PgPool, artist: &str) -> Result<Option<(String, String)>> {
let row = sqlx::query(
"SELECT m.file_path, m.mime_type FROM furumusic__artist a
+1 -1
View File
@@ -3,6 +3,7 @@ mod agent;
mod api;
mod auth;
mod config;
mod federation;
mod i18n;
mod jobs;
mod lastfm;
@@ -12,7 +13,6 @@ mod music;
mod oidc;
mod player;
mod scheduler;
mod federation;
mod torrents;
mod user;