From 53b2ff29f868174dd1abbebebd86a1b0c8835d7c Mon Sep 17 00:00:00 2001 From: Ultradesu Date: Mon, 20 Jul 2026 01:59:45 +0300 Subject: [PATCH] Moved DHT state in psql --- Cargo.lock | 1 + Cargo.toml | 3 +- src/federation/mod.rs | 28 +++- src/federation/storage.rs | 316 ++++++++++++++++++++++++++++++++++++++ 4 files changed, 340 insertions(+), 8 deletions(-) create mode 100644 src/federation/storage.rs diff --git a/Cargo.lock b/Cargo.lock index 40a726d..ba11386 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1859,6 +1859,7 @@ dependencies = [ "md-5", "music-dht", "openidconnect", + "postcard", "reqwest 0.12.28", "schemars 0.9.0", "serde", diff --git a/Cargo.toml b/Cargo.toml index 4b2e9a0..ea67213 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "furumusic" -version = "0.6.0-fd" +version = "0.6.1-fd" edition = "2024" description = "Reusable web-app boilerplate: auth, OIDC/SSO, admin panel, user management, i18n, PostgreSQL" @@ -31,6 +31,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 / diff --git a/src/federation/mod.rs b/src/federation/mod.rs index a55c81c..31bb83b 100644 --- a/src/federation/mod.rs +++ b/src/federation/mod.rs @@ -13,6 +13,7 @@ //! starts, stops or re-joins the node without a server restart. mod serve; +mod storage; use std::path::PathBuf; use std::sync::{Arc, OnceLock}; @@ -27,6 +28,7 @@ use sqlx::PgPool; use sqlx::Row as _; use crate::config::AppConfig; +use storage::PostgresFederationStorage; pub use serve::{AUDIO_ALPN, CATALOG_ALPN}; @@ -40,7 +42,7 @@ struct Running { } 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, pool: tokio::sync::OnceCell, @@ -170,6 +172,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 +184,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(), @@ -276,14 +282,22 @@ impl Federation { match service.sync_library(specs).await { Ok(stats) => { *lock(&self.last_sync) = Some(format!( - "{} (+{} ~{} −{}, unchanged {})", + "{} (+{} ~{} −{}, unchanged {}, failed {})", now_iso(), stats.added, stats.updated, stats.removed, - stats.unchanged + stats.unchanged, + stats.failed )); - self.set_error(None); + 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); + } } Err(err) => { tracing::warn!("federation sync failed: {err}"); diff --git a/src/federation/storage.rs b/src/federation/storage.rs new file mode 100644 index 0000000..bc3cdf9 --- /dev/null +++ b/src/federation/storage.rs @@ -0,0 +1,316 @@ +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 + )", +]; + +#[derive(Debug, Clone)] +pub struct PostgresFederationStorage { + pool: PgPool, +} + +impl PostgresFederationStorage { + pub async fn new(pool: PgPool) -> music_dht::Result { + let storage = Self { pool }; + storage.ensure_schema().await?; + Ok(storage) + } + + pub async fn load_or_create_secret_key(&self) -> music_dht::Result { + if let Some(bytes) = sqlx::query_scalar::<_, Vec>( + "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>( + "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> { + 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::(&row.get::, _>(0)).ok()) + .collect()) + } + + async fn store_dht_record(&self, key: DhtKey, record: StoredRecord) -> music_dht::Result { + 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(&self.pool) + .await + .map_err(db_error)? + .map(|row| { + ( + row.get::(0) as u64, + row.get::(1), + row.get::(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(&self.pool) + .await + .map_err(db_error)?; + Ok(true) + } + + async fn dht_records_by_key( + &self, + key: DhtKey, + now_ms: u64, + ) -> music_dht::Result> { + 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::(&row.get::, _>(0)).ok()) + .collect()) + } + + async fn delete_expired_records(&self, now_ms: u64) -> music_dht::Result { + 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> { + 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 = 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) -> music_dht::Result { + let bytes: [u8; 32] = bytes + .as_slice() + .try_into() + .map_err(|_| MusicDhtError::Database("stored federation identity is corrupted".into()))?; + Ok(SecretKey::from_bytes(&bytes)) +} + +fn db_error(err: impl std::fmt::Display) -> MusicDhtError { + MusicDhtError::Database(err.to_string()) +}