Compare commits

...
11 Commits
Author SHA1 Message Date
Ultradesu 2fc5fd7960 Improved DHT publication. Adjust batching, added validation per key
Build and Publish / Build and Publish Docker Image (push) Successful in 4m9s
2026-07-23 16:44:35 +03:00
Ultradesu 63506e3af2 Aupdated FRID to v4 protocol
Build and Publish / Build and Publish Docker Image (push) Successful in 5m4s
2026-07-23 15:52:42 +03:00
Ultradesu 5b339aa921 Fixed content_id persistence
Build and Publish / Build and Publish Docker Image (push) Successful in 7m30s
2026-07-20 18:40:32 +03:00
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
Ultradesu 3485f643f4 Fixed DHT republish-after-ttl
Build and Publish / Build and Publish Docker Image (push) Successful in 4m58s
2026-07-20 14:19:11 +03:00
Ultradesu bca0f5e2f0 added psql support
Build and Publish / Build and Publish Docker Image (push) Successful in 9m24s
2026-07-20 02:03:44 +03:00
Ultradesu 53b2ff29f8 Moved DHT state in psql 2026-07-20 01:59:45 +03:00
Ultradesu c349512fb0 fixed gitignore 2026-07-17 14:01:13 +03:00
9 changed files with 864 additions and 316 deletions
+1
View File
@@ -2,3 +2,4 @@
/nul
/.claude
/media
/federation
Generated
+223 -253
View File
File diff suppressed because it is too large Load Diff
+3 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "furumusic"
version = "0.6.0-fd"
version = "0.7.0"
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"] }
@@ -31,6 +32,7 @@ anyhow = "1.0"
tokio-cron-scheduler = "0.15"
croner = "3"
async-trait = "0.1"
postcard = { version = "1", features = ["alloc"] }
uuid = "1"
librqbit = { version = "8.1.1", features = ["disable-upload"] }
# P2P federation: publishes the library into a shared DHT and serves audio /
-1
View File
@@ -1 +0,0 @@
f48e5b2499a4ce1be021b2132401de29a1c23e84163d22e778874a8950757ff1
Binary file not shown.
+232 -34
View File
@@ -13,20 +13,24 @@
//! starts, stops or re-joins the node without a server restart.
mod serve;
mod storage;
use std::collections::{HashMap, HashSet};
use std::path::PathBuf;
use std::sync::{Arc, OnceLock};
use std::time::Duration;
use anyhow::{Context, Result};
use music_dht::{
ItemKind, ItemSpec, MusicDhtConfig, MusicDhtService, NetworkId, PeerTicket, RendezvousConfig,
ItemKind, ItemSpec, MusicDhtConfig, MusicDhtService, NetworkId, PeerTicket, PublishStats,
RendezvousConfig, SyncStats,
};
use serde_json::{Value, json};
use sqlx::PgPool;
use sqlx::Row as _;
use crate::config::AppConfig;
use storage::PostgresFederationStorage;
pub use serve::{AUDIO_ALPN, CATALOG_ALPN};
@@ -39,10 +43,19 @@ struct Running {
tasks: Vec<tokio::task::JoinHandle<()>>,
}
struct ContentHashJob {
media_file_id: i64,
sha256_hash: String,
file_path: String,
}
pub struct Federation {
/// Directory for the peer identity and DHT replica state.
/// 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>>,
@@ -66,6 +79,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),
@@ -134,8 +150,7 @@ impl Federation {
}
"federation_network_id" => effective.federation_network_id = value,
"agent_storage_dir" => {
effective.agent_storage_dir =
crate::media_paths::resolve_config_path(&value);
effective.agent_storage_dir = crate::media_paths::resolve_config_path(&value);
}
_ => {}
}
@@ -147,6 +162,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 {
@@ -170,6 +186,9 @@ impl Federation {
stop_running(guard.take()).await;
}
let dht_storage = Arc::new(PostgresFederationStorage::new(pool.clone()).await?);
let secret_key = dht_storage.load_or_create_secret_key().await?;
let config = MusicDhtConfig::builder()
.data_dir(&self.data_dir)
.network_id(NetworkId::from_name(&network_name))
@@ -179,9 +198,10 @@ impl Federation {
.stream_protocol(CATALOG_ALPN)
.build()
.map_err(|err| anyhow::anyhow!("invalid federation config: {err}"))?;
let (service, mut events) = MusicDhtService::start(config)
.await
.map_err(|err| anyhow::anyhow!("failed to start the DHT node: {err}"))?;
let (service, mut events) =
MusicDhtService::start_with_storage_and_secret_key(config, dht_storage, secret_key)
.await
.map_err(|err| anyhow::anyhow!("failed to start the DHT node: {err}"))?;
let service = Arc::new(service);
tracing::info!(
endpoint_id = %service.endpoint_id(),
@@ -202,7 +222,7 @@ impl Federation {
let mut interval = tokio::time::interval(SYNC_INTERVAL);
loop {
interval.tick().await;
sync_self.sync_once(&sync_service).await;
let _ = sync_self.sync_once(&sync_service).await;
}
});
// Serve audio and catalog requests from other peers.
@@ -254,54 +274,94 @@ impl Federation {
async fn spawn_sync_soon(self: &Arc<Self>) {
if let Ok(service) = self.service().await {
let fed = Arc::clone(self);
tokio::spawn(async move { fed.sync_once(&service).await });
tokio::spawn(async move {
let _ = fed.sync_once(&service).await;
});
}
}
pub async fn sync_now(self: &Arc<Self>) -> Result<()> {
let service = self.service().await?;
self.sync_once(&service).await;
let sync_stats = self.sync_once(&service).await?;
let publish_stats = match service.republish().await {
Ok(stats) => stats,
Err(err) => {
tracing::warn!("federation republish failed: {err}");
self.set_error(Some(format!("republish failed: {err}")));
anyhow::bail!("republish failed: {err}");
}
};
self.record_publish_success(sync_stats, publish_stats);
Ok(())
}
async fn sync_once(&self, service: &MusicDhtService) {
async fn sync_once(self: &Arc<Self>, service: &MusicDhtService) -> Result<SyncStats> {
let specs = match self.collect_specs().await {
Ok(specs) => specs,
Err(err) => {
tracing::warn!("federation sync: library read failed: {err:#}");
self.set_error(Some(format!("library read failed: {err}")));
return;
anyhow::bail!("library read failed: {err}");
}
};
match service.sync_library(specs).await {
Ok(stats) => {
*lock(&self.last_sync) = Some(format!(
"{} (+{} ~{} {}, unchanged {})",
now_iso(),
stats.added,
stats.updated,
stats.removed,
stats.unchanged
));
self.set_error(None);
self.record_sync_success(stats);
if stats.failed > 0 {
self.set_error(Some(format!(
"{} item(s) failed to publish in the last sync",
stats.failed
)));
} else {
self.set_error(None);
}
Ok(stats)
}
Err(err) => {
tracing::warn!("federation sync failed: {err}");
self.set_error(Some(format!("sync failed: {err}")));
Err(anyhow::anyhow!("sync failed: {err}"))
}
}
}
fn record_sync_success(&self, stats: SyncStats) {
*lock(&self.last_sync) = Some(format!(
"{} (+{} ~{} {}, unchanged {}, failed {})",
now_iso(),
stats.added,
stats.updated,
stats.removed,
stats.unchanged,
stats.failed
));
}
fn record_publish_success(&self, sync_stats: SyncStats, publish_stats: PublishStats) {
*lock(&self.last_sync) = Some(format!(
"{} (+{} ~{} {}, unchanged {}, failed {}; republished {} records, {} keys, remote nodes {})",
now_iso(),
sync_stats.added,
sync_stats.updated,
sync_stats.removed,
sync_stats.unchanged,
sync_stats.failed,
publish_stats.records,
publish_stats.keys,
publish_stats.remote_nodes,
));
self.set_error(None);
}
/// 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();
let artists =
sqlx::query("SELECT id, name FROM furumusic__artist WHERE is_hidden = false")
.fetch_all(&pool)
.await?;
let artists = sqlx::query("SELECT id, name FROM furumusic__artist WHERE is_hidden = false")
.fetch_all(&pool)
.await?;
for row in &artists {
let id: i64 = row.get(0);
specs.push(ItemSpec {
@@ -309,9 +369,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,
});
}
@@ -343,52 +408,150 @@ 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, c.content_id
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
LEFT JOIN furumusic__federation_content_id_cache c
ON c.media_file_id = m.id AND c.sha256_hash = m.sha256_hash
WHERE t.is_hidden = false AND r.is_hidden = false",
)
.fetch_all(&pool)
.await?;
let storage_dir = lock(&self.storage_dir).clone();
let mut content_hash_jobs = Vec::new();
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 cached_content_id: Option<String> = row.get(11);
let content_id = cached_content_id
.or_else(|| self.cached_content_id_for_media(media_file_id, &sha256_hash));
if content_id.is_none()
&& !storage_dir.trim().is_empty()
&& self.mark_content_hash_pending(media_file_id)
{
content_hash_jobs.push(ContentHashJob {
media_file_id,
sha256_hash,
file_path,
});
}
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,
});
}
self.spawn_content_warmer(pool.clone(), storage_dir, content_hash_jobs);
Ok(specs)
}
fn cached_content_id_for_media(&self, media_file_id: i64, sha256_hash: &str) -> Option<String> {
if let Some((cached_hash, content_id)) = lock(&self.content_cache).get(&media_file_id)
&& cached_hash == sha256_hash
{
return Some(content_id.clone());
}
None
}
fn mark_content_hash_pending(&self, media_file_id: i64) -> bool {
lock(&self.content_pending).insert(media_file_id)
}
fn spawn_content_warmer(
self: &Arc<Self>,
pool: PgPool,
storage_dir: String,
jobs: Vec<ContentHashJob>,
) {
if jobs.is_empty() {
return;
}
let fed = Arc::clone(self);
tokio::spawn(async move {
let total = jobs.len();
let mut stored = 0usize;
for job in jobs {
let job_storage_dir = storage_dir.clone();
let job_file_path = job.file_path.clone();
let content_id = tokio::task::spawn_blocking(move || {
audio_content_id(&job_storage_dir, &job_file_path)
})
.await
.ok()
.flatten();
lock(&fed.content_pending).remove(&job.media_file_id);
if let Some(content_id) = content_id {
lock(&fed.content_cache).insert(
job.media_file_id,
(job.sha256_hash.clone(), content_id.clone()),
);
if let Err(err) =
persist_content_id(&pool, job.media_file_id, &job.sha256_hash, &content_id)
.await
{
tracing::warn!(
media_file_id = job.media_file_id,
"federation content-id cache write failed: {err:#}"
);
} else {
stored += 1;
}
}
}
tracing::info!(total, stored, "federation content-id cache warm finished");
});
}
/// Live status for the admin page.
pub async fn status(&self) -> Value {
let guard = self.running.lock().await;
@@ -446,6 +609,41 @@ impl Federation {
}
}
async fn persist_content_id(
pool: &PgPool,
media_file_id: i64,
sha256_hash: &str,
content_id: &str,
) -> Result<()> {
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_file_id)
.bind(sha256_hash)
.bind(content_id)
.bind(now_iso())
.execute(pool)
.await?;
Ok(())
}
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 {
+55 -26
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,
)
@@ -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()?;
@@ -589,32 +614,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,
c.content_id
FROM furumusic__track t
JOIN furumusic__media_file m ON m.id = t.audio_file_id
LEFT JOIN furumusic__federation_content_id_cache c
ON c.media_file_id = m.id AND c.sha256_hash = m.sha256_hash
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);
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: row.get(5),
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,
});
}
+349
View File
@@ -0,0 +1,349 @@
use std::str::FromStr;
use async_trait::async_trait;
use music_dht::{
DhtKey, EndpointId, LibraryItem, MAX_RECORDS_PER_RESPONSE, MusicDhtError, MusicDhtStorage,
NodeContact, NodeId, SecretKey, StoreDecision, StoredRecord, decide_store,
};
use sqlx::{PgPool, Row as _};
const IDENTITY_NAME: &str = "default";
const SCHEMA: &[&str] = &[
"CREATE TABLE IF NOT EXISTS furumusic__federation_identity (
name TEXT PRIMARY KEY,
secret_key BYTEA NOT NULL,
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL
)",
"CREATE TABLE IF NOT EXISTS furumusic__federation_local_item (
id BYTEA PRIMARY KEY,
normalized_name TEXT NOT NULL,
revision BIGINT NOT NULL,
deleted BOOLEAN NOT NULL DEFAULT false,
updated_at_ms BIGINT NOT NULL,
payload BYTEA NOT NULL
)",
"CREATE INDEX IF NOT EXISTS idx_furumusic_federation_local_item_normalized_name
ON furumusic__federation_local_item(normalized_name)",
"CREATE TABLE IF NOT EXISTS furumusic__federation_dht_record (
dht_key BYTEA NOT NULL,
item_id BYTEA NOT NULL,
owner_peer_id TEXT NOT NULL,
payload BYTEA NOT NULL,
revision BIGINT NOT NULL,
deleted BOOLEAN NOT NULL,
expires_at_ms BIGINT NOT NULL,
PRIMARY KEY (dht_key, item_id, owner_peer_id)
)",
"CREATE INDEX IF NOT EXISTS idx_furumusic_federation_dht_record_expires_at
ON furumusic__federation_dht_record(expires_at_ms)",
"CREATE TABLE IF NOT EXISTS furumusic__federation_known_peer (
peer_id TEXT PRIMARY KEY,
node_id BYTEA NOT NULL,
ticket TEXT NOT NULL,
last_seen_ms BIGINT NOT NULL
)",
"CREATE TABLE IF NOT EXISTS furumusic__federation_content_id_cache (
media_file_id BIGINT PRIMARY KEY,
sha256_hash TEXT NOT NULL,
content_id TEXT NOT NULL,
updated_at TEXT NOT NULL
)",
"CREATE INDEX IF NOT EXISTS idx_furumusic_federation_content_id_cache_content_id
ON furumusic__federation_content_id_cache(content_id)",
];
#[derive(Debug, Clone)]
pub struct PostgresFederationStorage {
pool: PgPool,
}
impl PostgresFederationStorage {
pub async fn new(pool: PgPool) -> music_dht::Result<Self> {
let storage = Self { pool };
storage.ensure_schema().await?;
Ok(storage)
}
pub async fn load_or_create_secret_key(&self) -> music_dht::Result<SecretKey> {
if let Some(bytes) = sqlx::query_scalar::<_, Vec<u8>>(
"SELECT secret_key FROM furumusic__federation_identity WHERE name = $1",
)
.bind(IDENTITY_NAME)
.fetch_optional(&self.pool)
.await
.map_err(db_error)?
{
return secret_from_bytes(bytes);
}
let key = SecretKey::generate();
let key_bytes = key.to_bytes();
let now = now_iso();
let inserted = sqlx::query(
"INSERT INTO furumusic__federation_identity
(name, secret_key, created_at, updated_at)
VALUES ($1, $2, $3, $4)
ON CONFLICT (name) DO NOTHING",
)
.bind(IDENTITY_NAME)
.bind(key_bytes.as_slice())
.bind(&now)
.bind(&now)
.execute(&self.pool)
.await
.map_err(db_error)?
.rows_affected();
if inserted == 1 {
return Ok(key);
}
let bytes = sqlx::query_scalar::<_, Vec<u8>>(
"SELECT secret_key FROM furumusic__federation_identity WHERE name = $1",
)
.bind(IDENTITY_NAME)
.fetch_one(&self.pool)
.await
.map_err(db_error)?;
secret_from_bytes(bytes)
}
async fn ensure_schema(&self) -> music_dht::Result<()> {
for sql in SCHEMA {
sqlx::query(sql)
.execute(&self.pool)
.await
.map_err(db_error)?;
}
Ok(())
}
}
#[async_trait]
impl MusicDhtStorage for PostgresFederationStorage {
async fn upsert_local_item(&self, item: &LibraryItem) -> music_dht::Result<()> {
let payload = postcard::to_stdvec(item).map_err(db_error)?;
sqlx::query(
"INSERT INTO furumusic__federation_local_item
(id, normalized_name, revision, deleted, updated_at_ms, payload)
VALUES ($1, $2, $3, $4, $5, $6)
ON CONFLICT (id) DO UPDATE SET
normalized_name = EXCLUDED.normalized_name,
revision = EXCLUDED.revision,
deleted = EXCLUDED.deleted,
updated_at_ms = EXCLUDED.updated_at_ms,
payload = EXCLUDED.payload",
)
.bind(item.id.as_bytes().as_slice())
.bind(&item.normalized_name)
.bind(item.revision as i64)
.bind(item.deleted)
.bind(item.updated_at_ms as i64)
.bind(payload)
.execute(&self.pool)
.await
.map_err(db_error)?;
Ok(())
}
async fn list_local_items(&self, include_deleted: bool) -> music_dht::Result<Vec<LibraryItem>> {
let rows = sqlx::query(
"SELECT payload
FROM furumusic__federation_local_item
WHERE $1 OR deleted = false
ORDER BY normalized_name",
)
.bind(include_deleted)
.fetch_all(&self.pool)
.await
.map_err(db_error)?;
Ok(rows
.into_iter()
.filter_map(|row| postcard::from_bytes::<LibraryItem>(&row.get::<Vec<u8>, _>(0)).ok())
.collect())
}
async fn store_dht_record(&self, key: DhtKey, record: StoredRecord) -> music_dht::Result<bool> {
let mut conn = self.pool.acquire().await.map_err(db_error)?;
store_record_in_conn(&mut conn, &key, &record).await
}
async fn store_dht_records(
&self,
entries: Vec<(DhtKey, StoredRecord)>,
) -> music_dht::Result<Vec<bool>> {
let mut tx = self.pool.begin().await.map_err(db_error)?;
let mut stored = Vec::with_capacity(entries.len());
for (key, record) in &entries {
stored.push(store_record_in_conn(&mut tx, key, record).await?);
}
tx.commit().await.map_err(db_error)?;
Ok(stored)
}
async fn dht_records_by_key(
&self,
key: DhtKey,
now_ms: u64,
) -> music_dht::Result<Vec<StoredRecord>> {
let rows = sqlx::query(
"SELECT payload
FROM furumusic__federation_dht_record
WHERE dht_key = $1 AND expires_at_ms > $2
ORDER BY expires_at_ms DESC, item_id
LIMIT $3",
)
.bind(key.as_bytes().as_slice())
.bind(now_ms as i64)
.bind(MAX_RECORDS_PER_RESPONSE as i64)
.fetch_all(&self.pool)
.await
.map_err(db_error)?;
Ok(rows
.into_iter()
.filter_map(|row| postcard::from_bytes::<StoredRecord>(&row.get::<Vec<u8>, _>(0)).ok())
.collect())
}
async fn delete_expired_records(&self, now_ms: u64) -> music_dht::Result<usize> {
let result =
sqlx::query("DELETE FROM furumusic__federation_dht_record WHERE expires_at_ms <= $1")
.bind(now_ms as i64)
.execute(&self.pool)
.await
.map_err(db_error)?;
Ok(result.rows_affected() as usize)
}
async fn upsert_known_peer(&self, contact: &NodeContact) -> music_dht::Result<()> {
sqlx::query(
"INSERT INTO furumusic__federation_known_peer
(peer_id, node_id, ticket, last_seen_ms)
VALUES ($1, $2, $3, $4)
ON CONFLICT (peer_id) DO UPDATE SET
node_id = EXCLUDED.node_id,
ticket = EXCLUDED.ticket,
last_seen_ms = EXCLUDED.last_seen_ms",
)
.bind(contact.peer_id.to_string())
.bind(contact.node_id.as_bytes().as_slice())
.bind(&contact.ticket)
.bind(contact.last_seen_ms as i64)
.execute(&self.pool)
.await
.map_err(db_error)?;
Ok(())
}
async fn delete_known_peer(&self, peer_id: EndpointId) -> music_dht::Result<()> {
sqlx::query("DELETE FROM furumusic__federation_known_peer WHERE peer_id = $1")
.bind(peer_id.to_string())
.execute(&self.pool)
.await
.map_err(db_error)?;
Ok(())
}
async fn load_known_peers(&self) -> music_dht::Result<Vec<NodeContact>> {
let rows = sqlx::query(
"SELECT peer_id, node_id, ticket, last_seen_ms
FROM furumusic__federation_known_peer",
)
.fetch_all(&self.pool)
.await
.map_err(db_error)?;
let mut contacts = Vec::new();
for row in rows {
let peer_id: String = row.get(0);
let node_id: Vec<u8> = row.get(1);
let ticket: String = row.get(2);
let last_seen_ms: i64 = row.get(3);
let Ok(peer_id) = EndpointId::from_str(&peer_id) else {
continue;
};
let Ok(node_id) = <[u8; 32]>::try_from(node_id.as_slice()) else {
continue;
};
contacts.push(NodeContact {
node_id: NodeId::from_bytes(node_id),
peer_id,
ticket,
last_seen_ms: last_seen_ms as u64,
});
}
Ok(contacts)
}
}
fn now_iso() -> String {
chrono::Utc::now().format("%Y-%m-%dT%H:%M:%SZ").to_string()
}
fn secret_from_bytes(bytes: Vec<u8>) -> music_dht::Result<SecretKey> {
let bytes: [u8; 32] = bytes
.as_slice()
.try_into()
.map_err(|_| MusicDhtError::Database("stored federation identity is corrupted".into()))?;
Ok(SecretKey::from_bytes(&bytes))
}
/// Applies one validated record following the revision/tombstone rules.
/// Returns `true` if the record was written or refreshed. Runs against a
/// pooled connection or an open transaction.
async fn store_record_in_conn(
conn: &mut sqlx::PgConnection,
key: &DhtKey,
record: &StoredRecord,
) -> music_dht::Result<bool> {
let existing = sqlx::query(
"SELECT revision, deleted, expires_at_ms
FROM furumusic__federation_dht_record
WHERE dht_key = $1 AND item_id = $2 AND owner_peer_id = $3",
)
.bind(key.as_bytes().as_slice())
.bind(record.item.id.as_bytes().as_slice())
.bind(record.item.owner.to_string())
.fetch_optional(&mut *conn)
.await
.map_err(db_error)?
.map(|row| {
(
row.get::<i64, _>(0) as u64,
row.get::<bool, _>(1),
row.get::<i64, _>(2) as u64,
)
});
match decide_store(existing, record) {
StoreDecision::Ignore => return Ok(false),
StoreDecision::Write | StoreDecision::RefreshExpiry(_) => {}
}
let payload = postcard::to_stdvec(record).map_err(db_error)?;
sqlx::query(
"INSERT INTO furumusic__federation_dht_record
(dht_key, item_id, owner_peer_id, payload, revision, deleted, expires_at_ms)
VALUES ($1, $2, $3, $4, $5, $6, $7)
ON CONFLICT (dht_key, item_id, owner_peer_id) DO UPDATE SET
payload = EXCLUDED.payload,
revision = EXCLUDED.revision,
deleted = EXCLUDED.deleted,
expires_at_ms = EXCLUDED.expires_at_ms",
)
.bind(key.as_bytes().as_slice())
.bind(record.item.id.as_bytes().as_slice())
.bind(record.item.owner.to_string())
.bind(payload)
.bind(record.item.revision as i64)
.bind(record.item.deleted)
.bind(record.expires_at_ms as i64)
.execute(&mut *conn)
.await
.map_err(db_error)?;
Ok(true)
}
fn db_error(err: impl std::fmt::Display) -> MusicDhtError {
MusicDhtError::Database(err.to_string())
}
+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;