From 2fc5fd79602a3d2badb8d5ec60d6d2c62161c15a Mon Sep 17 00:00:00 2001 From: Ultradesu Date: Thu, 23 Jul 2026 16:44:35 +0300 Subject: [PATCH] Improved DHT publication. Adjust batching, added validation per key --- Cargo.lock | 4 +- Cargo.toml | 2 +- src/federation/storage.rs | 113 +++++++++++++++++++++++--------------- 3 files changed, 72 insertions(+), 47 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 461dd8c..176d20b 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1718,7 +1718,7 @@ dependencies = [ [[package]] name = "federation-net" version = "0.1.0" -source = "git+https://gt.hexor.cy/ab/frid.git#aaaad780f2bc3f4ba42bd3eef6aca871a2ff68da" +source = "git+https://gt.hexor.cy/ab/frid.git#512a818a6a52ec713678e9a4e1cf0f50bb1e34ab" dependencies = [ "blake3", "data-encoding", @@ -3622,7 +3622,7 @@ dependencies = [ [[package]] name = "music-dht" version = "0.1.0" -source = "git+https://gt.hexor.cy/ab/frid.git#aaaad780f2bc3f4ba42bd3eef6aca871a2ff68da" +source = "git+https://gt.hexor.cy/ab/frid.git#512a818a6a52ec713678e9a4e1cf0f50bb1e34ab" dependencies = [ "async-trait", "blake3", diff --git a/Cargo.toml b/Cargo.toml index 31722ef..55b1a2e 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "furumusic" -version = "0.6.7-fd" +version = "0.7.0" edition = "2024" description = "Reusable web-app boilerplate: auth, OIDC/SSO, admin panel, user management, i18n, PostgreSQL" diff --git a/src/federation/storage.rs b/src/federation/storage.rs index 070e8d8..7fac005 100644 --- a/src/federation/storage.rs +++ b/src/federation/storage.rs @@ -165,52 +165,21 @@ impl MusicDhtStorage for PostgresFederationStorage { } 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, - ) - }); + let mut conn = self.pool.acquire().await.map_err(db_error)?; + store_record_in_conn(&mut conn, &key, &record).await + } - match decide_store(existing, &record) { - StoreDecision::Ignore => return Ok(false), - StoreDecision::Write | StoreDecision::RefreshExpiry(_) => {} + async fn store_dht_records( + &self, + entries: Vec<(DhtKey, StoredRecord)>, + ) -> music_dht::Result> { + 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?); } - - 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) + tx.commit().await.map_err(db_error)?; + Ok(stored) } async fn dht_records_by_key( @@ -319,6 +288,62 @@ fn secret_from_bytes(bytes: Vec) -> music_dht::Result { 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 { + 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::(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(&mut *conn) + .await + .map_err(db_error)?; + Ok(true) +} + fn db_error(err: impl std::fmt::Display) -> MusicDhtError { MusicDhtError::Database(err.to_string()) }