Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
e64b61c167 | ||
|
|
b737ced3fc |
Generated
+1
-1
@@ -1845,7 +1845,7 @@ checksum = "e6d5a32815ae3f33302d95fdcb2ce17862f8c65363dcfd29360480ba1001fc9c"
|
||||
|
||||
[[package]]
|
||||
name = "furumusic"
|
||||
version = "0.7.1"
|
||||
version = "0.8.1"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"async-trait",
|
||||
|
||||
+1
-1
@@ -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
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user