Compare commits

..
2 Commits
Author SHA1 Message Date
Ultradesu e64b61c167 Added FED devices
Build and Publish / Build and Publish Docker Image (push) Successful in 4m11s
2026-07-24 17:08:00 +03:00
Ultradesu b737ced3fc Added FED devices
Build and Publish / Build and Publish Docker Image (push) Successful in 4m1s
2026-07-24 16:54:41 +03:00
3 changed files with 201 additions and 43 deletions
Generated
+1 -1
View File
@@ -1845,7 +1845,7 @@ checksum = "e6d5a32815ae3f33302d95fdcb2ce17862f8c65363dcfd29360480ba1001fc9c"
[[package]]
name = "furumusic"
version = "0.7.1"
version = "0.8.1"
dependencies = [
"anyhow",
"async-trait",
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "furumusic"
version = "0.8.0"
version = "0.8.2"
edition = "2024"
description = "Reusable web-app boilerplate: auth, OIDC/SSO, admin panel, user management, i18n, PostgreSQL"
+199 -41
View File
@@ -27,6 +27,7 @@ const PAIRING_WAIT_MS: i64 = 5 * 60 * 1000;
const PAIRING_RETRY_DELAY: Duration = Duration::from_secs(1);
const RESPONSE_DRAIN_TIMEOUT: Duration = Duration::from_secs(2);
const DEVICE_SYNC_INTERVAL: Duration = Duration::from_secs(2);
const LOCAL_SEED_RECHECK_MS: i64 = 60 * 1000;
const MAX_LINE: usize = 8 * 1024 * 1024;
const MAX_OPS_PER_BATCH: i64 = 1000;
@@ -596,6 +597,7 @@ pub async fn revoke_device(pool: &sqlx::PgPool, user_id: i64, device_id: &str) -
pub async fn status(pool: &sqlx::PgPool, user_id: i64, user_name: &str) -> Result<FedDeviceStatus> {
let identity = ensure_identity(pool, user_id, user_name).await?;
maybe_seed_local_user_state(pool, user_id).await?;
let active_devices: i64 = sqlx::query_scalar(
"SELECT COUNT(*) FROM furumusic__fed_device
WHERE user_id = $1 AND trusted_at_ms IS NOT NULL AND revoked_at_ms IS NULL",
@@ -746,13 +748,18 @@ pub async fn record_track_like(
liked: bool,
) -> Result<()> {
if let Some(content_id) = track_content_id(pool, track_id).await? {
let fed = if liked {
synced_fed_track_for_track(pool, track_id, &content_id).await?
} else {
None
};
record_local_op(
pool,
user_id,
SyncOpPayload::TrackLikeSet {
content_id,
liked,
fed: None,
fed,
},
)
.await?;
@@ -828,7 +835,8 @@ pub async fn record_playlist_tracks_added(
let title = playlist_title(pool, playlist_id).await?.unwrap_or_default();
let sync_id = ensure_local_playlist_sync_id(pool, user_id, playlist_id, &title).await?;
let rows = playlist_track_content_positions(pool, playlist_id, track_ids).await?;
for (content_id, position) in rows {
for (track_id, content_id, position) in rows {
let fed = synced_fed_track_for_track(pool, track_id, &content_id).await?;
record_local_op(
pool,
user_id,
@@ -836,7 +844,7 @@ pub async fn record_playlist_tracks_added(
playlist_id: sync_id.clone(),
content_id,
position,
fed: None,
fed,
},
)
.await?;
@@ -1467,6 +1475,48 @@ async fn ensure_identity(pool: &sqlx::PgPool, user_id: i64, user_name: &str) ->
})
}
async fn maybe_seed_local_user_state(pool: &sqlx::PgPool, user_id: i64) -> Result<()> {
let last_seeded_at_ms: Option<i64> = sqlx::query_scalar(
"SELECT local_seeded_at_ms
FROM furumusic__fed_device_identity
WHERE user_id = $1",
)
.bind(user_id)
.fetch_optional(pool)
.await?;
let Some(last_seeded_at_ms) = last_seeded_at_ms else {
return Ok(());
};
if now_ms().saturating_sub(last_seeded_at_ms) < LOCAL_SEED_RECHECK_MS
&& !local_seed_needs_metadata_backfill(pool, user_id).await?
{
return Ok(());
}
seed_local_user_state(pool, user_id).await
}
async fn local_seed_needs_metadata_backfill(pool: &sqlx::PgPool, user_id: i64) -> Result<bool> {
sqlx::query_scalar(
"SELECT EXISTS (
SELECT 1 FROM furumusic__fed_state_like
WHERE user_id = $1
AND liked = true
AND fed_json IS NULL
AND local_track_id IS NOT NULL
) OR EXISTS (
SELECT 1 FROM furumusic__fed_state_playlist_item
WHERE user_id = $1
AND present = true
AND fed_json IS NULL
AND local_track_id IS NOT NULL
)",
)
.bind(user_id)
.fetch_one(pool)
.await
.map_err(Into::into)
}
async fn seed_local_user_state(pool: &sqlx::PgPool, user_id: i64) -> Result<()> {
let now = now_ms();
let like_rows = sqlx::query(
@@ -1487,17 +1537,24 @@ async fn seed_local_user_state(pool: &sqlx::PgPool, user_id: i64) -> Result<()>
continue;
};
let track_id: i64 = row.get("track_id");
let fed_json = synced_fed_track_for_track(pool, track_id, &content_id)
.await?
.map(serde_json::to_value)
.transpose()?;
sqlx::query(
"INSERT INTO furumusic__fed_state_like
(user_id, content_id, liked, hlc_ms, op_id, local_track_id, fed_json)
VALUES ($1, $2, true, $3, $4, $5, NULL)
ON CONFLICT (user_id, content_id) DO NOTHING",
VALUES ($1, $2, true, $3, $4, $5, $6)
ON CONFLICT (user_id, content_id) DO UPDATE SET
local_track_id = COALESCE(furumusic__fed_state_like.local_track_id, EXCLUDED.local_track_id),
fed_json = COALESCE(furumusic__fed_state_like.fed_json, EXCLUDED.fed_json)",
)
.bind(user_id)
.bind(&content_id)
.bind(now)
.bind(format!("local_seed:like:{track_id}"))
.bind(track_id)
.bind(fed_json)
.execute(pool)
.await?;
}
@@ -1512,22 +1569,8 @@ async fn seed_local_user_state(pool: &sqlx::PgPool, user_id: i64) -> Result<()>
.await?;
for playlist in playlists {
let playlist_id: i64 = playlist.get("id");
let sync_id = format!("webpl_{user_id}_{playlist_id}");
let title: String = playlist.get("title");
sqlx::query(
"INSERT INTO furumusic__fed_state_playlist
(user_id, playlist_id, local_playlist_id, title, deleted, hlc_ms, op_id)
VALUES ($1, $2, $3, $4, false, $5, $6)
ON CONFLICT (user_id, playlist_id) DO NOTHING",
)
.bind(user_id)
.bind(&sync_id)
.bind(playlist_id)
.bind(&title)
.bind(now)
.bind(format!("local_seed:playlist:{playlist_id}"))
.execute(pool)
.await?;
let sync_id = ensure_seed_playlist_sync_id(pool, user_id, playlist_id, &title, now).await?;
let items = sqlx::query(
"SELECT pt.track_id, pt.position, c.content_id
@@ -1549,12 +1592,18 @@ async fn seed_local_user_state(pool: &sqlx::PgPool, user_id: i64) -> Result<()>
};
let track_id: i64 = item.get("track_id");
let position: i32 = item.get("position");
let fed_json = synced_fed_track_for_track(pool, track_id, &content_id)
.await?
.map(serde_json::to_value)
.transpose()?;
sqlx::query(
"INSERT INTO furumusic__fed_state_playlist_item
(user_id, playlist_id, content_id, present, position, hlc_ms, op_id,
local_track_id, fed_json)
VALUES ($1, $2, $3, true, $4, $5, $6, $7, NULL)
ON CONFLICT (user_id, playlist_id, content_id) DO NOTHING",
VALUES ($1, $2, $3, true, $4, $5, $6, $7, $8)
ON CONFLICT (user_id, playlist_id, content_id) DO UPDATE SET
local_track_id = COALESCE(furumusic__fed_state_playlist_item.local_track_id, EXCLUDED.local_track_id),
fed_json = COALESCE(furumusic__fed_state_playlist_item.fed_json, EXCLUDED.fed_json)",
)
.bind(user_id)
.bind(&sync_id)
@@ -1563,6 +1612,7 @@ async fn seed_local_user_state(pool: &sqlx::PgPool, user_id: i64) -> Result<()>
.bind(now)
.bind(format!("local_seed:playlist_item:{playlist_id}:{track_id}"))
.bind(track_id)
.bind(fed_json)
.execute(pool)
.await?;
}
@@ -1593,17 +1643,16 @@ async fn seed_local_user_state(pool: &sqlx::PgPool, user_id: i64) -> Result<()>
.bind(user_id)
.fetch_one(pool)
.await?;
if missing_likes + missing_playlist_items == 0 {
sqlx::query(
"UPDATE furumusic__fed_device_identity
SET local_seeded_at_ms = $2
WHERE user_id = $1 AND local_seeded_at_ms = 0",
)
.bind(user_id)
.bind(now)
.execute(pool)
.await?;
} else {
sqlx::query(
"UPDATE furumusic__fed_device_identity
SET local_seeded_at_ms = $2
WHERE user_id = $1",
)
.bind(user_id)
.bind(now)
.execute(pool)
.await?;
if missing_likes + missing_playlist_items != 0 {
tracing::debug!(
user_id,
missing_likes,
@@ -3274,6 +3323,21 @@ async fn track_artist_json(
.collect())
}
async fn track_artist_names(pool: &sqlx::PgPool, track_id: i64, role: &str) -> Result<Vec<String>> {
let rows = sqlx::query(
"SELECT a.name::text AS name
FROM furumusic__track_artist ta
JOIN furumusic__artist a ON a.id = ta.artist_id
WHERE ta.track_id = $1 AND ta.role = $2
ORDER BY ta.position",
)
.bind(track_id)
.bind(role)
.fetch_all(pool)
.await?;
Ok(rows.into_iter().map(|row| row.get("name")).collect())
}
fn placeholder_playback_track(
track: &PlaybackTrack,
index: usize,
@@ -3387,6 +3451,45 @@ async fn track_content_id(pool: &sqlx::PgPool, track_id: i64) -> Result<Option<S
Ok(value.and_then(|value| music_dht::normalize_content_id(&value)))
}
async fn synced_fed_track_for_track(
pool: &sqlx::PgPool,
track_id: i64,
content_id: &str,
) -> Result<Option<SyncedFedTrack>> {
let Some(content_id) = music_dht::normalize_content_id(content_id) else {
return Ok(None);
};
let row = sqlx::query(
"SELECT t.title::text AS title, t.track_number, t.disc_number,
t.duration_seconds, r.title::text AS release_title, r.year AS release_year
FROM furumusic__track t
JOIN furumusic__release r ON r.id = t.release_id
WHERE t.id = $1",
)
.bind(track_id)
.fetch_optional(pool)
.await?;
let Some(row) = row else {
return Ok(None);
};
let duration = row.get::<f64, _>("duration_seconds");
Ok(Some(SyncedFedTrack {
item_id: format!("web:{track_id}"),
owner: "WEB".to_string(),
title: row.get("title"),
artist_names: track_artist_names(pool, track_id, "main").await?,
featured_artist_names: track_artist_names(pool, track_id, "featuring").await?,
year: row.get("release_year"),
duration_seconds: duration
.is_finite()
.then(|| duration.round().max(0.0) as i64),
content_id,
release_title: row.get("release_title"),
track_number: row.get("track_number"),
disc_number: row.get("disc_number"),
}))
}
async fn ensure_materialized_playlist(
pool: &sqlx::PgPool,
user_id: i64,
@@ -3414,17 +3517,36 @@ async fn ensure_materialized_playlist(
.bind(&now)
.fetch_one(pool)
.await?;
upsert_materialized_playlist_link(pool, user_id, playlist_sync_id, id, title).await?;
Ok(id)
}
async fn upsert_materialized_playlist_link(
pool: &sqlx::PgPool,
user_id: i64,
playlist_sync_id: &str,
local_playlist_id: i64,
title: &str,
) -> Result<()> {
sqlx::query(
"UPDATE furumusic__fed_state_playlist
SET local_playlist_id = $3
WHERE user_id = $1 AND playlist_id = $2",
"INSERT INTO furumusic__fed_state_playlist
(user_id, playlist_id, local_playlist_id, title, deleted, hlc_ms, op_id)
VALUES ($1, $2, $3, $4, false, 0, $5)
ON CONFLICT (user_id, playlist_id) DO UPDATE SET
local_playlist_id = COALESCE(furumusic__fed_state_playlist.local_playlist_id, EXCLUDED.local_playlist_id),
title = CASE
WHEN furumusic__fed_state_playlist.hlc_ms = 0 THEN EXCLUDED.title
ELSE furumusic__fed_state_playlist.title
END",
)
.bind(user_id)
.bind(playlist_sync_id)
.bind(id)
.bind(local_playlist_id)
.bind(title)
.bind(format!("local_materialized:{local_playlist_id}"))
.execute(pool)
.await?;
Ok(id)
Ok(())
}
async fn ensure_playlist_for_item(
@@ -3449,6 +3571,40 @@ async fn ensure_playlist_for_item(
ensure_materialized_playlist(pool, user_id, playlist_sync_id, None, &title).await
}
async fn ensure_seed_playlist_sync_id(
pool: &sqlx::PgPool,
user_id: i64,
playlist_id: i64,
title: &str,
hlc_ms: i64,
) -> Result<String> {
if let Some(sync_id) = local_playlist_sync_id(pool, user_id, playlist_id).await? {
return Ok(sync_id);
}
let sync_id = format!("webpl_{user_id}_{playlist_id}");
sqlx::query(
"INSERT INTO furumusic__fed_state_playlist
(user_id, playlist_id, local_playlist_id, title, deleted, hlc_ms, op_id)
VALUES ($1, $2, $3, $4, false, $5, $6)
ON CONFLICT (user_id, playlist_id) DO UPDATE SET
local_playlist_id = COALESCE(furumusic__fed_state_playlist.local_playlist_id, EXCLUDED.local_playlist_id),
title = CASE
WHEN furumusic__fed_state_playlist.hlc_ms = 0 THEN EXCLUDED.title
ELSE furumusic__fed_state_playlist.title
END,
deleted = false",
)
.bind(user_id)
.bind(&sync_id)
.bind(playlist_id)
.bind(title)
.bind(hlc_ms)
.bind(format!("local_seed:playlist:{playlist_id}"))
.execute(pool)
.await?;
Ok(sync_id)
}
async fn ensure_local_playlist_sync_id(
pool: &sqlx::PgPool,
user_id: i64,
@@ -3524,9 +3680,9 @@ async fn playlist_track_content_positions(
pool: &sqlx::PgPool,
playlist_id: i64,
track_ids: &[i64],
) -> Result<Vec<(String, i64)>> {
) -> Result<Vec<(i64, String, i64)>> {
let rows = sqlx::query(
"SELECT c.content_id, pt.position
"SELECT pt.track_id, c.content_id, pt.position
FROM furumusic__playlist_track pt
JOIN furumusic__track t ON t.id = pt.track_id
JOIN furumusic__media_file m ON m.id = t.audio_file_id
@@ -3542,9 +3698,10 @@ async fn playlist_track_content_positions(
Ok(rows
.into_iter()
.filter_map(|row| {
let track_id: i64 = row.get("track_id");
let content_id: String = row.get("content_id");
music_dht::normalize_content_id(&content_id)
.map(|content_id| (content_id, row.get::<i32, _>("position") as i64))
.map(|content_id| (track_id, content_id, row.get::<i32, _>("position") as i64))
})
.collect())
}
@@ -3707,6 +3864,7 @@ async fn ops_for_peer(
}
async fn snapshot(pool: &sqlx::PgPool, user_id: i64) -> Result<SyncSnapshot> {
maybe_seed_local_user_state(pool, user_id).await?;
let like_rows = sqlx::query(
"SELECT content_id, liked, hlc_ms, op_id, fed_json
FROM furumusic__fed_state_like