Music Player & Playlist Manager

This commit is contained in:
Andrey Tkachenko
2026-04-08 15:14:03 +03:00
parent 2e8e0ef6c5
commit da066f47fa
53 changed files with 6043 additions and 23 deletions
+2
View File
@@ -10,7 +10,9 @@
pub mod calendar_storage_adapter;
pub mod contact_storage_adapter;
pub mod error_adapters;
pub mod music_storage_adapter;
pub use calendar_storage_adapter::CalendarStorageAdapter;
pub use contact_storage_adapter::ContactStorageAdapter;
pub use error_adapters::IntoDomainError;
pub use music_storage_adapter::MusicStorageAdapter;
@@ -0,0 +1,296 @@
use std::sync::Arc;
use uuid::Uuid;
use crate::application::dtos::playlist_dto::{
AudioMetadataDto, CreatePlaylistDto, PlaylistDto, PlaylistItemDto, UpdatePlaylistDto,
};
use crate::application::ports::music_ports::MusicStoragePort;
use crate::common::errors::{DomainError, ErrorKind};
use crate::domain::entities::playlist::Playlist;
use crate::domain::repositories::playlist_repository::{
AudioMetadataRepository, PlaylistItemRepository, PlaylistRepository,
};
use crate::infrastructure::repositories::pg::{
AudioMetadataPgRepository, PlaylistItemPgRepository, PlaylistPgRepository,
};
pub struct MusicStorageAdapter {
playlist_repository: Arc<PlaylistPgRepository>,
item_repository: Arc<PlaylistItemPgRepository>,
audio_metadata_repository: Arc<AudioMetadataPgRepository>,
}
impl MusicStorageAdapter {
pub fn new(
playlist_repository: Arc<PlaylistPgRepository>,
item_repository: Arc<PlaylistItemPgRepository>,
audio_metadata_repository: Arc<AudioMetadataPgRepository>,
) -> Self {
Self {
playlist_repository,
item_repository,
audio_metadata_repository,
}
}
}
impl MusicStoragePort for MusicStorageAdapter {
async fn create_playlist(
&self,
dto: CreatePlaylistDto,
user_id: Uuid,
) -> Result<PlaylistDto, DomainError> {
let playlist = Playlist::new(dto.name, user_id, dto.description)?;
let created = self.playlist_repository.create_playlist(playlist).await?;
Ok(PlaylistDto::from(created))
}
async fn update_playlist(
&self,
playlist_id: &str,
dto: UpdatePlaylistDto,
) -> Result<PlaylistDto, DomainError> {
let uuid = Uuid::parse_str(playlist_id).map_err(|_| {
DomainError::new(ErrorKind::InvalidInput, "Playlist", "Invalid playlist ID")
})?;
let mut playlist = self.playlist_repository.find_playlist_by_id(&uuid).await?;
if let Some(name) = dto.name {
playlist.update_name(name)?;
}
if let Some(description) = dto.description {
playlist.update_description(Some(description));
}
if let Some(is_public) = dto.is_public {
playlist.set_public(is_public);
}
if let Some(cover_file_id) = dto.cover_file_id {
let cover_uuid = Uuid::parse_str(&cover_file_id).map_err(|_| {
DomainError::new(ErrorKind::InvalidInput, "Playlist", "Invalid cover file ID")
})?;
playlist.set_cover(Some(cover_uuid));
}
let updated = self.playlist_repository.update_playlist(playlist).await?;
Ok(PlaylistDto::from(updated))
}
async fn delete_playlist(&self, playlist_id: &str) -> Result<(), DomainError> {
let uuid = Uuid::parse_str(playlist_id).map_err(|_| {
DomainError::new(ErrorKind::InvalidInput, "Playlist", "Invalid playlist ID")
})?;
self.playlist_repository.delete_playlist(&uuid).await
}
async fn get_playlist(&self, playlist_id: &str) -> Result<Option<PlaylistDto>, DomainError> {
let uuid = Uuid::parse_str(playlist_id).map_err(|_| {
DomainError::new(ErrorKind::InvalidInput, "Playlist", "Invalid playlist ID")
})?;
match self.playlist_repository.find_playlist_by_id(&uuid).await {
Ok(playlist) => Ok(Some(PlaylistDto::from(playlist))),
Err(e) if e.kind == ErrorKind::NotFound => Ok(None),
Err(e) => Err(e),
}
}
async fn list_playlists_by_owner(
&self,
owner_id: Uuid,
) -> Result<Vec<PlaylistDto>, DomainError> {
let playlists = self
.playlist_repository
.list_playlists_by_owner(owner_id)
.await?;
let mut result = Vec::new();
for playlist in playlists {
let dto = PlaylistDto::from(playlist);
let track_count = self
.get_track_count(&uuid::Uuid::parse_str(&dto.id).unwrap())
.await?;
result.push(dto.with_track_info(track_count, 0));
}
Ok(result)
}
async fn list_shared_with_user(&self, user_id: Uuid) -> Result<Vec<PlaylistDto>, DomainError> {
let playlists = self
.playlist_repository
.list_shared_with_user(user_id)
.await?;
let mut result = Vec::new();
for playlist in playlists {
let dto = PlaylistDto::from(playlist);
let track_count = self
.get_track_count(&uuid::Uuid::parse_str(&dto.id).unwrap())
.await?;
result.push(dto.with_track_info(track_count, 0));
}
Ok(result)
}
async fn list_public_playlists(
&self,
limit: i64,
offset: i64,
) -> Result<Vec<PlaylistDto>, DomainError> {
let playlists = self
.playlist_repository
.list_public_playlists(limit, offset)
.await?;
let mut result = Vec::new();
for playlist in playlists {
let dto = PlaylistDto::from(playlist);
let track_count = self
.get_track_count(&uuid::Uuid::parse_str(&dto.id).unwrap())
.await?;
result.push(dto.with_track_info(track_count, 0));
}
Ok(result)
}
async fn user_has_access(&self, playlist_id: &str, user_id: Uuid) -> Result<bool, DomainError> {
let uuid = Uuid::parse_str(playlist_id).map_err(|_| {
DomainError::new(ErrorKind::InvalidInput, "Playlist", "Invalid playlist ID")
})?;
self.playlist_repository
.user_has_access(&uuid, user_id)
.await
}
async fn user_can_write(&self, playlist_id: &str, user_id: Uuid) -> Result<bool, DomainError> {
let uuid = Uuid::parse_str(playlist_id).map_err(|_| {
DomainError::new(ErrorKind::InvalidInput, "Playlist", "Invalid playlist ID")
})?;
let playlist = self.playlist_repository.find_playlist_by_id(&uuid).await?;
if playlist.owner_id() == &user_id {
return Ok(true);
}
let shares = self.playlist_repository.get_shares(&uuid).await?;
Ok(shares
.iter()
.any(|(uid, can_write)| uid == &user_id && *can_write))
}
async fn add_tracks(
&self,
playlist_id: &Uuid,
file_ids: &[Uuid],
) -> Result<Vec<PlaylistItemDto>, DomainError> {
let mut max_position = self.item_repository.get_max_position(playlist_id).await?;
let mut items = Vec::new();
for file_id in file_ids {
max_position += 1;
let item = crate::domain::entities::playlist::PlaylistItem::new(
*playlist_id,
*file_id,
max_position,
)?;
let created = self.item_repository.add_item(item).await?;
items.push(PlaylistItemDto::from(created));
}
Ok(items)
}
async fn remove_track(&self, playlist_id: &Uuid, file_id: &Uuid) -> Result<(), DomainError> {
self.item_repository
.remove_item_by_playlist_and_file(playlist_id, file_id)
.await
}
async fn reorder_tracks(
&self,
playlist_id: &Uuid,
item_ids: &[Uuid],
) -> Result<(), DomainError> {
self.item_repository
.reorder_items(playlist_id, item_ids)
.await
}
async fn list_playlist_tracks(
&self,
playlist_id: &Uuid,
) -> Result<Vec<PlaylistItemDto>, DomainError> {
let enriched_items = self
.item_repository
.list_items_in_playlist_enriched(playlist_id)
.await?;
let items: Vec<PlaylistItemDto> = enriched_items
.into_iter()
.map(
|row| crate::application::dtos::playlist_dto::PlaylistItemDto {
id: row.id.to_string(),
playlist_id: row.playlist_id.to_string(),
file_id: row.file_id.to_string(),
position: row.position,
added_at: row.added_at,
file_name: row.file_name,
file_size: row.file_size,
mime_type: row.mime_type,
title: row.title,
artist: row.artist,
album: row.album,
duration_secs: row.duration_secs,
},
)
.collect();
Ok(items)
}
async fn share_playlist(
&self,
playlist_id: &Uuid,
user_id: Uuid,
can_write: bool,
) -> Result<(), DomainError> {
self.playlist_repository
.share_playlist(playlist_id, user_id, can_write)
.await
}
async fn remove_share(&self, playlist_id: &Uuid, user_id: Uuid) -> Result<(), DomainError> {
self.playlist_repository
.remove_share(playlist_id, user_id)
.await
}
async fn get_shares(&self, playlist_id: &Uuid) -> Result<Vec<(Uuid, bool)>, DomainError> {
self.playlist_repository.get_shares(playlist_id).await
}
async fn get_audio_metadata(
&self,
file_id: &Uuid,
) -> Result<Option<AudioMetadataDto>, DomainError> {
match self
.audio_metadata_repository
.find_by_file_id(file_id)
.await
{
Ok(Some(metadata)) => Ok(Some(AudioMetadataDto::from(metadata))),
Ok(None) => Ok(None),
Err(e) => Err(e),
}
}
}
impl MusicStorageAdapter {
async fn get_track_count(&self, playlist_id: &Uuid) -> Result<i64, DomainError> {
let count: (i64,) =
sqlx::query_as("SELECT COUNT(*) FROM audio.playlist_items WHERE playlist_id = $1")
.bind(playlist_id)
.fetch_one(self.playlist_repository.pool())
.await
.map_err(|e| {
DomainError::database_error(format!("Failed to get track count: {}", e))
})?;
Ok(count.0)
}
}
@@ -9,6 +9,7 @@ mod device_code_pg_repository;
mod favorites_pg_repository;
pub mod file_metadata_repository;
mod nextcloud_object_id_repository;
pub mod playlist_pg_repository;
mod recent_items_pg_repository;
mod session_pg_repository;
mod settings_pg_repository;
@@ -36,6 +37,9 @@ pub use file_blob_write_repository::FileBlobWriteRepository;
pub use file_metadata_repository::FileMetadataRepository;
pub use folder_db_repository::FolderDbRepository;
pub use nextcloud_object_id_repository::NextcloudObjectIdRepository;
pub use playlist_pg_repository::{
AudioMetadataPgRepository, PlaylistItemPgRepository, PlaylistPgRepository,
};
pub use recent_items_pg_repository::RecentItemsPgRepository;
pub use session_pg_repository::SessionPgRepository;
pub use settings_pg_repository::SettingsPgRepository;
@@ -0,0 +1,787 @@
use chrono::{DateTime, Utc};
use sqlx::{FromRow, PgPool};
use std::sync::Arc;
use uuid::Uuid;
use crate::common::errors::{DomainError, ErrorKind};
use crate::domain::entities::playlist::{AudioFileMetadata, Playlist, PlaylistItem};
use crate::domain::repositories::playlist_repository::{
AudioMetadataRepository, AudioMetadataRepositoryResult, PlaylistItemRepository,
PlaylistItemRepositoryResult, PlaylistRepository, PlaylistRepositoryResult,
};
#[derive(FromRow)]
struct PlaylistRow {
id: Uuid,
name: String,
description: Option<String>,
owner_id: Uuid,
is_public: bool,
cover_file_id: Option<Uuid>,
created_at: DateTime<Utc>,
updated_at: DateTime<Utc>,
}
#[derive(FromRow)]
struct PlaylistItemRow {
id: Uuid,
playlist_id: Uuid,
file_id: Uuid,
position: i32,
added_at: DateTime<Utc>,
}
#[derive(FromRow)]
struct AudioMetadataRow {
file_id: Uuid,
title: Option<String>,
artist: Option<String>,
album: Option<String>,
album_artist: Option<String>,
genre: Option<String>,
track_number: Option<i32>,
disc_number: Option<i32>,
year: Option<i32>,
duration_secs: i32,
bitrate: Option<i32>,
sample_rate: Option<i32>,
channels: Option<i16>,
format: Option<String>,
created_at: DateTime<Utc>,
updated_at: DateTime<Utc>,
}
#[derive(FromRow)]
pub struct PlaylistItemEnrichedRow {
pub id: Uuid,
pub playlist_id: Uuid,
pub file_id: Uuid,
pub position: i32,
pub added_at: DateTime<Utc>,
pub file_name: Option<String>,
pub file_size: Option<i64>,
pub mime_type: Option<String>,
pub title: Option<String>,
pub artist: Option<String>,
pub album: Option<String>,
pub duration_secs: Option<i32>,
}
pub struct PlaylistPgRepository {
pool: Arc<PgPool>,
}
impl PlaylistPgRepository {
pub fn new(pool: Arc<PgPool>) -> Self {
Self { pool }
}
pub fn pool(&self) -> &PgPool {
&self.pool
}
}
impl PlaylistRepository for PlaylistPgRepository {
async fn create_playlist(&self, playlist: Playlist) -> PlaylistRepositoryResult<Playlist> {
let row = sqlx::query_as::<_, PlaylistRow>(
r#"
INSERT INTO audio.playlists (id, name, description, owner_id, is_public, cover_file_id, created_at, updated_at)
VALUES ($1, $2, $3, $4, $5, $6, $7, $8)
RETURNING id, name, description, owner_id, is_public, cover_file_id, created_at, updated_at
"#,
)
.bind(playlist.id())
.bind(playlist.name())
.bind(playlist.description())
.bind(playlist.owner_id())
.bind(playlist.is_public())
.bind(playlist.cover_file_id())
.bind(playlist.created_at())
.bind(playlist.updated_at())
.fetch_one(&*self.pool)
.await
.map_err(|e| DomainError::database_error(format!("Failed to create playlist: {}", e)))?;
Playlist::with_id(
row.id,
row.name,
row.description,
row.owner_id,
row.is_public,
row.cover_file_id,
row.created_at,
row.updated_at,
)
.map_err(|e| DomainError::new(ErrorKind::InternalError, "Playlist", e.to_string()))
}
async fn update_playlist(&self, playlist: Playlist) -> PlaylistRepositoryResult<Playlist> {
let row = sqlx::query_as::<_, PlaylistRow>(
r#"
UPDATE audio.playlists
SET name = $2, description = $3, is_public = $4, cover_file_id = $5, updated_at = NOW()
WHERE id = $1
RETURNING id, name, description, owner_id, is_public, cover_file_id, created_at, updated_at
"#,
)
.bind(playlist.id())
.bind(playlist.name())
.bind(playlist.description())
.bind(playlist.is_public())
.bind(playlist.cover_file_id())
.fetch_one(&*self.pool)
.await
.map_err(|e| DomainError::database_error(format!("Failed to update playlist: {}", e)))?;
Playlist::with_id(
row.id,
row.name,
row.description,
row.owner_id,
row.is_public,
row.cover_file_id,
row.created_at,
row.updated_at,
)
.map_err(|e| DomainError::new(ErrorKind::InternalError, "Playlist", e.to_string()))
}
async fn delete_playlist(&self, id: &Uuid) -> PlaylistRepositoryResult<()> {
sqlx::query("DELETE FROM audio.playlists WHERE id = $1")
.bind(id)
.execute(&*self.pool)
.await
.map_err(|e| {
DomainError::database_error(format!("Failed to delete playlist: {}", e))
})?;
Ok(())
}
async fn find_playlist_by_id(&self, id: &Uuid) -> PlaylistRepositoryResult<Playlist> {
let row = sqlx::query_as::<_, PlaylistRow>(
"SELECT id, name, description, owner_id, is_public, cover_file_id, created_at, updated_at FROM audio.playlists WHERE id = $1",
)
.bind(id)
.fetch_optional(&*self.pool)
.await
.map_err(|e| DomainError::database_error(format!("Failed to find playlist: {}", e)))?
.ok_or_else(|| DomainError::new(ErrorKind::NotFound, "Playlist", "Playlist not found"))?;
Playlist::with_id(
row.id,
row.name,
row.description,
row.owner_id,
row.is_public,
row.cover_file_id,
row.created_at,
row.updated_at,
)
.map_err(|e| DomainError::new(ErrorKind::InternalError, "Playlist", e.to_string()))
}
async fn list_playlists_by_owner(
&self,
owner_id: Uuid,
) -> PlaylistRepositoryResult<Vec<Playlist>> {
let rows = sqlx::query_as::<_, PlaylistRow>(
"SELECT id, name, description, owner_id, is_public, cover_file_id, created_at, updated_at FROM audio.playlists WHERE owner_id = $1 ORDER BY updated_at DESC",
)
.bind(owner_id)
.fetch_all(&*self.pool)
.await
.map_err(|e| DomainError::database_error(format!("Failed to list playlists: {}", e)))?;
rows.into_iter()
.map(|row| {
Playlist::with_id(
row.id,
row.name,
row.description,
row.owner_id,
row.is_public,
row.cover_file_id,
row.created_at,
row.updated_at,
)
.map_err(|e| DomainError::new(ErrorKind::InternalError, "Playlist", e.to_string()))
})
.collect()
}
async fn list_public_playlists(
&self,
limit: i64,
offset: i64,
) -> PlaylistRepositoryResult<Vec<Playlist>> {
let rows = sqlx::query_as::<_, PlaylistRow>(
"SELECT id, name, description, owner_id, is_public, cover_file_id, created_at, updated_at FROM audio.playlists WHERE is_public = TRUE ORDER BY updated_at DESC LIMIT $1 OFFSET $2",
)
.bind(limit)
.bind(offset)
.fetch_all(&*self.pool)
.await
.map_err(|e| DomainError::database_error(format!("Failed to list public playlists: {}", e)))?;
rows.into_iter()
.map(|row| {
Playlist::with_id(
row.id,
row.name,
row.description,
row.owner_id,
row.is_public,
row.cover_file_id,
row.created_at,
row.updated_at,
)
.map_err(|e| DomainError::new(ErrorKind::InternalError, "Playlist", e.to_string()))
})
.collect()
}
async fn list_shared_with_user(
&self,
user_id: Uuid,
) -> PlaylistRepositoryResult<Vec<Playlist>> {
let rows = sqlx::query_as::<_, PlaylistRow>(
r#"
SELECT p.id, p.name, p.description, p.owner_id, p.is_public, p.cover_file_id, p.created_at, p.updated_at
FROM audio.playlists p
JOIN audio.playlist_shares ps ON p.id = ps.playlist_id
WHERE ps.user_id = $1
ORDER BY p.updated_at DESC
"#,
)
.bind(user_id)
.fetch_all(&*self.pool)
.await
.map_err(|e| DomainError::database_error(format!("Failed to list shared playlists: {}", e)))?;
rows.into_iter()
.map(|row| {
Playlist::with_id(
row.id,
row.name,
row.description,
row.owner_id,
row.is_public,
row.cover_file_id,
row.created_at,
row.updated_at,
)
.map_err(|e| DomainError::new(ErrorKind::InternalError, "Playlist", e.to_string()))
})
.collect()
}
async fn user_has_access(
&self,
playlist_id: &Uuid,
user_id: Uuid,
) -> PlaylistRepositoryResult<bool> {
let row = sqlx::query_scalar::<_, bool>(
r#"
SELECT EXISTS(
SELECT 1 FROM audio.playlists p
WHERE p.id = $1 AND (p.owner_id = $2 OR p.is_public = TRUE)
UNION
SELECT 1 FROM audio.playlist_shares ps
WHERE ps.playlist_id = $1 AND ps.user_id = $2
)
"#,
)
.bind(playlist_id)
.bind(user_id)
.fetch_one(&*self.pool)
.await
.map_err(|e| {
DomainError::database_error(format!("Failed to check playlist access: {}", e))
})?;
Ok(row)
}
async fn share_playlist(
&self,
playlist_id: &Uuid,
user_id: Uuid,
can_write: bool,
) -> PlaylistRepositoryResult<()> {
sqlx::query(
r#"
INSERT INTO audio.playlist_shares (playlist_id, user_id, can_write)
VALUES ($1, $2, $3)
ON CONFLICT (playlist_id, user_id) DO UPDATE SET can_write = $3
"#,
)
.bind(playlist_id)
.bind(user_id)
.bind(can_write)
.execute(&*self.pool)
.await
.map_err(|e| DomainError::database_error(format!("Failed to share playlist: {}", e)))?;
Ok(())
}
async fn remove_share(
&self,
playlist_id: &Uuid,
user_id: Uuid,
) -> PlaylistRepositoryResult<()> {
sqlx::query("DELETE FROM audio.playlist_shares WHERE playlist_id = $1 AND user_id = $2")
.bind(playlist_id)
.bind(user_id)
.execute(&*self.pool)
.await
.map_err(|e| {
DomainError::database_error(format!("Failed to remove playlist share: {}", e))
})?;
Ok(())
}
async fn get_shares(&self, playlist_id: &Uuid) -> PlaylistRepositoryResult<Vec<(Uuid, bool)>> {
let rows = sqlx::query_as::<_, (Uuid, bool)>(
"SELECT user_id, can_write FROM audio.playlist_shares WHERE playlist_id = $1",
)
.bind(playlist_id)
.fetch_all(&*self.pool)
.await
.map_err(|e| {
DomainError::database_error(format!("Failed to get playlist shares: {}", e))
})?;
Ok(rows)
}
}
pub struct PlaylistItemPgRepository {
pool: Arc<PgPool>,
}
impl PlaylistItemPgRepository {
pub fn new(pool: Arc<PgPool>) -> Self {
Self { pool }
}
pub async fn list_items_in_playlist_enriched(
&self,
playlist_id: &Uuid,
) -> PlaylistItemRepositoryResult<Vec<PlaylistItemEnrichedRow>> {
let rows = sqlx::query_as::<_, PlaylistItemEnrichedRow>(
r#"
SELECT
pi.id, pi.playlist_id, pi.file_id, pi.position, pi.added_at,
f.name as file_name, f.size as file_size, f.mime_type,
m.title, m.artist, m.album, m.duration_secs
FROM audio.playlist_items pi
LEFT JOIN storage.files f ON pi.file_id = f.id
LEFT JOIN audio.file_metadata m ON pi.file_id = m.file_id
WHERE pi.playlist_id = $1
ORDER BY pi.position ASC
"#,
)
.bind(playlist_id)
.fetch_all(&*self.pool)
.await
.map_err(|e| {
DomainError::database_error(format!("Failed to list playlist items: {}", e))
})?;
Ok(rows)
}
}
impl PlaylistItemRepository for PlaylistItemPgRepository {
async fn add_item(&self, item: PlaylistItem) -> PlaylistItemRepositoryResult<PlaylistItem> {
let row = sqlx::query_as::<_, PlaylistItemRow>(
r#"
INSERT INTO audio.playlist_items (id, playlist_id, file_id, position, added_at)
VALUES ($1, $2, $3, $4, $5)
ON CONFLICT (playlist_id, file_id) DO UPDATE SET position = $4
RETURNING id, playlist_id, file_id, position, added_at
"#,
)
.bind(item.id())
.bind(item.playlist_id())
.bind(item.file_id())
.bind(item.position())
.bind(item.added_at())
.fetch_one(&*self.pool)
.await
.map_err(|e| DomainError::database_error(format!("Failed to add playlist item: {}", e)))?;
PlaylistItem::with_id(
row.id,
row.playlist_id,
row.file_id,
row.position,
row.added_at,
)
.map_err(|e| DomainError::new(ErrorKind::InternalError, "PlaylistItem", e.to_string()))
}
async fn remove_item(&self, id: &Uuid) -> PlaylistItemRepositoryResult<()> {
sqlx::query("DELETE FROM audio.playlist_items WHERE id = $1")
.bind(id)
.execute(&*self.pool)
.await
.map_err(|e| {
DomainError::database_error(format!("Failed to remove playlist item: {}", e))
})?;
Ok(())
}
async fn remove_item_by_playlist_and_file(
&self,
playlist_id: &Uuid,
file_id: &Uuid,
) -> PlaylistItemRepositoryResult<()> {
sqlx::query("DELETE FROM audio.playlist_items WHERE playlist_id = $1 AND file_id = $2")
.bind(playlist_id)
.bind(file_id)
.execute(&*self.pool)
.await
.map_err(|e| {
DomainError::database_error(format!("Failed to remove playlist item: {}", e))
})?;
Ok(())
}
async fn find_item_by_id(&self, id: &Uuid) -> PlaylistItemRepositoryResult<PlaylistItem> {
let row = sqlx::query_as::<_, PlaylistItemRow>(
"SELECT id, playlist_id, file_id, position, added_at FROM audio.playlist_items WHERE id = $1",
)
.bind(id)
.fetch_optional(&*self.pool)
.await
.map_err(|e| DomainError::database_error(format!("Failed to find playlist item: {}", e)))?
.ok_or_else(|| DomainError::new(ErrorKind::NotFound, "PlaylistItem", "Item not found"))?;
PlaylistItem::with_id(
row.id,
row.playlist_id,
row.file_id,
row.position,
row.added_at,
)
.map_err(|e| DomainError::new(ErrorKind::InternalError, "PlaylistItem", e.to_string()))
}
async fn list_items_in_playlist(
&self,
playlist_id: &Uuid,
) -> PlaylistItemRepositoryResult<Vec<PlaylistItem>> {
let rows = sqlx::query_as::<_, PlaylistItemRow>(
"SELECT id, playlist_id, file_id, position, added_at FROM audio.playlist_items WHERE playlist_id = $1 ORDER BY position ASC",
)
.bind(playlist_id)
.fetch_all(&*self.pool)
.await
.map_err(|e| DomainError::database_error(format!("Failed to list playlist items: {}", e)))?;
rows.into_iter()
.map(|row| {
PlaylistItem::with_id(
row.id,
row.playlist_id,
row.file_id,
row.position,
row.added_at,
)
.map_err(|e| {
DomainError::new(ErrorKind::InternalError, "PlaylistItem", e.to_string())
})
})
.collect()
}
async fn update_position(
&self,
id: &Uuid,
new_position: i32,
) -> PlaylistItemRepositoryResult<()> {
sqlx::query("UPDATE audio.playlist_items SET position = $2 WHERE id = $1")
.bind(id)
.bind(new_position)
.execute(&*self.pool)
.await
.map_err(|e| {
DomainError::database_error(format!("Failed to update position: {}", e))
})?;
Ok(())
}
async fn reorder_items(
&self,
playlist_id: &Uuid,
item_ids: &[Uuid],
) -> PlaylistItemRepositoryResult<()> {
for (index, item_id) in item_ids.iter().enumerate() {
sqlx::query(
"UPDATE audio.playlist_items SET position = $2 WHERE id = $1 AND playlist_id = $3",
)
.bind(item_id)
.bind(index as i32)
.bind(playlist_id)
.execute(&*self.pool)
.await
.map_err(|e| DomainError::database_error(format!("Failed to reorder: {}", e)))?;
}
Ok(())
}
async fn get_item_count(&self, playlist_id: &Uuid) -> PlaylistItemRepositoryResult<i64> {
let count = sqlx::query_scalar::<_, i64>(
"SELECT COUNT(*) FROM audio.playlist_items WHERE playlist_id = $1",
)
.bind(playlist_id)
.fetch_one(&*self.pool)
.await
.map_err(|e| DomainError::database_error(format!("Failed to count items: {}", e)))?;
Ok(count)
}
async fn get_max_position(&self, playlist_id: &Uuid) -> PlaylistItemRepositoryResult<i32> {
let max_pos = sqlx::query_scalar::<_, Option<i32>>(
"SELECT MAX(position) FROM audio.playlist_items WHERE playlist_id = $1",
)
.bind(playlist_id)
.fetch_one(&*self.pool)
.await
.map_err(|e| DomainError::database_error(format!("Failed to get max position: {}", e)))?;
Ok(max_pos.unwrap_or(-1))
}
}
pub struct AudioMetadataPgRepository {
pool: Arc<PgPool>,
}
impl AudioMetadataPgRepository {
pub fn new(pool: Arc<PgPool>) -> Self {
Self { pool }
}
}
impl AudioMetadataRepository for AudioMetadataPgRepository {
async fn create_or_update(
&self,
metadata: AudioFileMetadata,
) -> AudioMetadataRepositoryResult<AudioFileMetadata> {
let row = sqlx::query_as::<_, AudioMetadataRow>(
r#"
INSERT INTO audio.file_metadata (file_id, title, artist, album, album_artist, genre, track_number, disc_number, year, duration_secs, bitrate, sample_rate, channels, format, created_at, updated_at)
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16)
ON CONFLICT (file_id) DO UPDATE SET
title = COALESCE($2, audio.file_metadata.title),
artist = COALESCE($3, audio.file_metadata.artist),
album = COALESCE($4, audio.file_metadata.album),
album_artist = COALESCE($5, audio.file_metadata.album_artist),
genre = COALESCE($6, audio.file_metadata.genre),
track_number = COALESCE($7, audio.file_metadata.track_number),
disc_number = COALESCE($8, audio.file_metadata.disc_number),
year = COALESCE($9, audio.file_metadata.year),
duration_secs = $10,
bitrate = COALESCE($11, audio.file_metadata.bitrate),
sample_rate = COALESCE($12, audio.file_metadata.sample_rate),
channels = COALESCE($13, audio.file_metadata.channels),
format = COALESCE($14, audio.file_metadata.format),
updated_at = NOW()
RETURNING file_id, title, artist, album, album_artist, genre, track_number, disc_number, year, duration_secs, bitrate, sample_rate, channels, format, created_at, updated_at
"#,
)
.bind(metadata.file_id())
.bind(metadata.title())
.bind(metadata.artist())
.bind(metadata.album())
.bind(metadata.album_artist())
.bind(metadata.genre())
.bind(metadata.track_number())
.bind(metadata.disc_number())
.bind(metadata.year())
.bind(metadata.duration_secs())
.bind(metadata.bitrate())
.bind(metadata.sample_rate())
.bind(metadata.channels())
.bind(metadata.format())
.bind(metadata.created_at())
.bind(metadata.updated_at())
.fetch_one(&*self.pool)
.await
.map_err(|e| DomainError::database_error(format!("Failed to create audio metadata: {}", e)))?;
Ok(AudioFileMetadata::with_all_fields(
row.file_id,
row.title,
row.artist,
row.album,
row.album_artist,
row.genre,
row.track_number,
row.disc_number,
row.year,
row.duration_secs,
row.bitrate,
row.sample_rate,
row.channels,
row.format,
row.created_at,
row.updated_at,
))
}
async fn find_by_file_id(
&self,
file_id: &Uuid,
) -> AudioMetadataRepositoryResult<Option<AudioFileMetadata>> {
let row = sqlx::query_as::<_, AudioMetadataRow>(
"SELECT file_id, title, artist, album, album_artist, genre, track_number, disc_number, year, duration_secs, bitrate, sample_rate, channels, format, created_at, updated_at FROM audio.file_metadata WHERE file_id = $1",
)
.bind(file_id)
.fetch_optional(&*self.pool)
.await
.map_err(|e| DomainError::database_error(format!("Failed to find audio metadata: {}", e)))?;
Ok(row.map(|r| {
AudioFileMetadata::with_all_fields(
r.file_id,
r.title,
r.artist,
r.album,
r.album_artist,
r.genre,
r.track_number,
r.disc_number,
r.year,
r.duration_secs,
r.bitrate,
r.sample_rate,
r.channels,
r.format,
r.created_at,
r.updated_at,
)
}))
}
async fn delete(&self, file_id: &Uuid) -> AudioMetadataRepositoryResult<()> {
sqlx::query("DELETE FROM audio.file_metadata WHERE file_id = $1")
.bind(file_id)
.execute(&*self.pool)
.await
.map_err(|e| {
DomainError::database_error(format!("Failed to delete audio metadata: {}", e))
})?;
Ok(())
}
async fn list_by_artist(
&self,
artist: &str,
) -> AudioMetadataRepositoryResult<Vec<AudioFileMetadata>> {
let rows = sqlx::query_as::<_, AudioMetadataRow>(
"SELECT file_id, title, artist, album, album_artist, genre, track_number, disc_number, year, duration_secs, bitrate, sample_rate, channels, format, created_at, updated_at FROM audio.file_metadata WHERE artist ILIKE $1 ORDER BY album, track_number",
)
.bind(format!("%{}%", artist))
.fetch_all(&*self.pool)
.await
.map_err(|e| DomainError::database_error(format!("Failed to list by artist: {}", e)))?;
Ok(rows
.into_iter()
.map(|r| {
AudioFileMetadata::with_all_fields(
r.file_id,
r.title,
r.artist,
r.album,
r.album_artist,
r.genre,
r.track_number,
r.disc_number,
r.year,
r.duration_secs,
r.bitrate,
r.sample_rate,
r.channels,
r.format,
r.created_at,
r.updated_at,
)
})
.collect())
}
async fn list_by_album(
&self,
album: &str,
) -> AudioMetadataRepositoryResult<Vec<AudioFileMetadata>> {
let rows = sqlx::query_as::<_, AudioMetadataRow>(
"SELECT file_id, title, artist, album, album_artist, genre, track_number, disc_number, year, duration_secs, bitrate, sample_rate, channels, format, created_at, updated_at FROM audio.file_metadata WHERE album ILIKE $1 ORDER BY disc_number, track_number",
)
.bind(format!("%{}%", album))
.fetch_all(&*self.pool)
.await
.map_err(|e| DomainError::database_error(format!("Failed to list by album: {}", e)))?;
Ok(rows
.into_iter()
.map(|r| {
AudioFileMetadata::with_all_fields(
r.file_id,
r.title,
r.artist,
r.album,
r.album_artist,
r.genre,
r.track_number,
r.disc_number,
r.year,
r.duration_secs,
r.bitrate,
r.sample_rate,
r.channels,
r.format,
r.created_at,
r.updated_at,
)
})
.collect())
}
async fn list_by_genre(
&self,
genre: &str,
) -> AudioMetadataRepositoryResult<Vec<AudioFileMetadata>> {
let rows = sqlx::query_as::<_, AudioMetadataRow>(
"SELECT file_id, title, artist, album, album_artist, genre, track_number, disc_number, year, duration_secs, bitrate, sample_rate, channels, format, created_at, updated_at FROM audio.file_metadata WHERE genre ILIKE $1 ORDER BY artist, album, track_number",
)
.bind(format!("%{}%", genre))
.fetch_all(&*self.pool)
.await
.map_err(|e| DomainError::database_error(format!("Failed to list by genre: {}", e)))?;
Ok(rows
.into_iter()
.map(|r| {
AudioFileMetadata::with_all_fields(
r.file_id,
r.title,
r.artist,
r.album,
r.album_artist,
r.genre,
r.track_number,
r.disc_number,
r.year,
r.duration_secs,
r.bitrate,
r.sample_rate,
r.channels,
r.format,
r.created_at,
r.updated_at,
)
})
.collect())
}
}
@@ -0,0 +1,224 @@
use id3::{Tag, TagLike};
use sqlx::{FromRow, PgPool};
use std::path::{Path, PathBuf};
use std::sync::Arc;
use tracing::{info, warn};
use uuid::Uuid;
use crate::common::errors::DomainError;
#[derive(Debug, FromRow)]
pub struct AudioFileRow {
pub file_id: Uuid,
pub blob_hash: String,
}
pub struct AudioMetadataService {
pool: Arc<PgPool>,
blob_root: PathBuf,
}
impl AudioMetadataService {
pub fn new(pool: Arc<PgPool>, blob_root: PathBuf) -> Self {
Self { pool, blob_root }
}
pub fn is_audio_file(mime_type: &str) -> bool {
mime_type.starts_with("audio/")
}
pub fn spawn_extraction_background(service: Arc<Self>, file_id: Uuid, file_path: PathBuf) {
tokio::spawn(async move {
tracing::info!("🎵 Extracting audio metadata for: {}", file_id);
if let Err(e) = service.extract_and_save(&file_id, &file_path).await {
tracing::warn!("Failed to extract audio metadata: {}", e);
}
});
}
pub fn spawn_extraction_with_delete_background(
service: Arc<Self>,
file_id: Uuid,
file_path: PathBuf,
) {
tokio::spawn(async move {
tracing::info!("🎵 Updating audio metadata for: {}", file_id);
let _ = service.delete_metadata(&file_id).await;
if let Err(e) = service.extract_and_save(&file_id, &file_path).await {
tracing::warn!("Failed to update audio metadata: {}", e);
}
});
}
fn blob_path(&self, hash: &str) -> PathBuf {
let prefix = &hash[0..2];
self.blob_root.join(prefix).join(format!("{}.blob", hash))
}
fn get_duration_secs(file_path: &Path) -> i32 {
match mp3_duration::from_path(file_path) {
Ok(dur) => dur.as_secs_f64().round() as i32,
Err(_) => {
if let Ok(tag) = Tag::read_from_path(file_path) {
tag.duration().unwrap_or(0) as i32
} else {
0
}
}
}
}
pub async fn extract_and_save(
&self,
file_id: &Uuid,
file_path: &Path,
) -> Result<(), DomainError> {
info!(
"AudioMetadataService: blob_root={:?}, file_id={}, file_path={:?}, exists={}",
self.blob_root,
file_id,
file_path,
file_path.exists()
);
if !file_path.exists() {
warn!("File does not exist: {:?}", file_path);
return Ok(());
}
let tag = match Tag::read_from_path(file_path) {
Ok(t) => t,
Err(e) => {
warn!("Failed to read ID3 tag from {:?}: {}", file_path, e);
return Ok(());
}
};
let title = tag.title().map(|s| s.to_string());
let artist = tag.artist().map(|s| s.to_string());
let album = tag.album().map(|s| s.to_string());
let genre = tag.genre().map(|s| s.to_string());
let track_number: Option<i32> = tag.track().map(|n| n as i32);
let disc_number: Option<i32> = tag.disc().map(|n| n as i32);
let year: Option<i32> = tag.year();
let duration_secs = Self::get_duration_secs(file_path);
let album_artist =
tag.frames()
.find(|f| f.id() == "TPE2")
.and_then(|f| match f.content() {
id3::frame::Content::Text(t) => Some(t.clone()),
_ => None,
});
info!(
"Extracted audio metadata for file {}: title={:?}, artist={:?}, album={:?}, duration={}s",
file_id, title, artist, album, duration_secs
);
info!("Saving metadata to database for file_id={}", file_id);
sqlx::query(
r#"
INSERT INTO audio.file_metadata
(file_id, title, artist, album, album_artist, genre, track_number, disc_number,
year, duration_secs, format)
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11)
ON CONFLICT (file_id) DO UPDATE SET
title = EXCLUDED.title,
artist = EXCLUDED.artist,
album = EXCLUDED.album,
album_artist = EXCLUDED.album_artist,
genre = EXCLUDED.genre,
track_number = EXCLUDED.track_number,
disc_number = EXCLUDED.disc_number,
year = EXCLUDED.year,
duration_secs = EXCLUDED.duration_secs,
format = EXCLUDED.format,
updated_at = CURRENT_TIMESTAMP
"#,
)
.bind(file_id)
.bind(&title)
.bind(&artist)
.bind(&album)
.bind(&album_artist)
.bind(&genre)
.bind(track_number)
.bind(disc_number)
.bind(year)
.bind(duration_secs)
.bind("MPEG")
.execute(&*self.pool)
.await
.map_err(|e| {
DomainError::database_error(format!("Failed to save audio metadata: {}", e))
})?;
Ok(())
}
pub async fn delete_metadata(&self, file_id: &Uuid) -> Result<(), DomainError> {
sqlx::query("DELETE FROM audio.file_metadata WHERE file_id = $1")
.bind(file_id)
.execute(&*self.pool)
.await
.map_err(|e| {
DomainError::database_error(format!("Failed to delete audio metadata: {}", e))
})?;
Ok(())
}
pub async fn reextract_all_audio_metadata(
&self,
) -> Result<MetadataExtractionResult, DomainError> {
let audio_files = sqlx::query_as::<_, AudioFileRow>(
r#"
SELECT id as file_id, blob_hash
FROM storage.files
WHERE mime_type LIKE 'audio/%'
"#,
)
.fetch_all(&*self.pool)
.await
.map_err(|e| DomainError::database_error(format!("Failed to fetch audio files: {}", e)))?;
let total = audio_files.len();
let mut processed = 0;
let mut failed = 0;
info!("Starting metadata extraction for {} audio files", total);
for audio_file in audio_files {
let file_path = self.blob_path(&audio_file.blob_hash);
match self.extract_and_save(&audio_file.file_id, &file_path).await {
Ok(()) => processed += 1,
Err(e) => {
warn!(
"Failed to extract metadata for file {}: {}",
audio_file.file_id, e
);
failed += 1;
}
}
}
info!(
"Metadata extraction complete: {} processed, {} failed out of {} total",
processed, failed, total
);
Ok(MetadataExtractionResult {
total,
processed,
failed,
})
}
}
#[derive(Debug, serde::Serialize)]
pub struct MetadataExtractionResult {
pub total: usize,
pub processed: usize,
pub failed: usize,
}
+1
View File
@@ -1,3 +1,4 @@
pub mod audio_metadata_service;
pub mod chunked_upload_service;
pub mod compression_service;
pub mod dedup_service;