refactor: extract adapter-common crate — shared sqlx error mapping,

row-to-domain conversions, date/uuid utils

Eliminates 39 map_err copies, consolidates parse_uuid/parse_datetime/
datetime_to_str/format_year_month, extracts 7 shared row-to-domain
conversion functions (movie, review, watchlist, stats, user_summary).
Row structs stay per-adapter (FromRow is db-specific), only conversion
logic is shared. Net -281 lines.
This commit is contained in:
2026-07-10 05:44:42 +02:00
parent 6a9b4e5c00
commit c224cc6bd2
66 changed files with 924 additions and 971 deletions

View File

@@ -11,6 +11,7 @@ sqlx = { version = "0.8.6", features = [
"macros",
"chrono",
] }
adapter-common = { workspace = true }
domain = { workspace = true }
postgres-federation = { workspace = true }
anyhow = { workspace = true }

View File

@@ -22,23 +22,19 @@ impl PostgresDiaryRepository {
Self { pool }
}
fn map_err(e: sqlx::Error) -> DomainError {
tracing::error!("Database error: {:?}", e);
DomainError::InfrastructureError("Database operation failed".into())
}
async fn count_diary_entries(&self, movie_id: Option<&str>) -> Result<i64, DomainError> {
match movie_id {
None => sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM reviews")
.fetch_one(&self.pool)
.await
.map_err(Self::map_err),
.map_err(adapter_common::map_sqlx_error),
Some(id) => {
sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM reviews WHERE movie_id = $1")
.bind(id)
.fetch_one(&self.pool)
.await
.map_err(Self::map_err)
.map_err(adapter_common::map_sqlx_error)
}
}
}
@@ -73,7 +69,7 @@ impl PostgresDiaryRepository {
.bind(offset)
.fetch_all(&self.pool)
.await
.map_err(Self::map_err)
.map_err(adapter_common::map_sqlx_error)
}
async fn fetch_movie_diary_rows(
@@ -109,7 +105,7 @@ impl PostgresDiaryRepository {
.bind(offset)
.fetch_all(&self.pool)
.await
.map_err(Self::map_err)
.map_err(adapter_common::map_sqlx_error)
}
async fn count_user_diary_entries(
@@ -138,7 +134,7 @@ impl PostgresDiaryRepository {
if has_search {
q = q.bind(search.unwrap());
}
q.fetch_one(&self.pool).await.map_err(Self::map_err)
q.fetch_one(&self.pool).await.map_err(adapter_common::map_sqlx_error)
}
async fn fetch_user_diary_rows(
@@ -197,7 +193,7 @@ impl PostgresDiaryRepository {
.bind(offset)
.fetch_all(&self.pool)
.await
.map_err(Self::map_err)
.map_err(adapter_common::map_sqlx_error)
}
}
@@ -375,7 +371,7 @@ impl DiaryQuery for PostgresDiaryRepository {
}
let count_q = bind_filter_params!(sqlx::query_scalar::<_, i64>(&count_sql));
let total = count_q.fetch_one(&self.pool).await.map_err(Self::map_err)?;
let total = count_q.fetch_one(&self.pool).await.map_err(adapter_common::map_sqlx_error)?;
let rows_q = bind_filter_params!(sqlx::query_as::<_, FeedRow>(&select_sql));
let rows = rows_q
@@ -383,7 +379,7 @@ impl DiaryQuery for PostgresDiaryRepository {
.bind(offset)
.fetch_all(&self.pool)
.await
.map_err(Self::map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
let items = rows
.into_iter()
@@ -408,7 +404,7 @@ impl DiaryQuery for PostgresDiaryRepository {
.bind(&id_str)
.fetch_optional(&self.pool)
.await
.map_err(Self::map_err)?
.map_err(adapter_common::map_sqlx_error)?
.ok_or_else(|| DomainError::NotFound(format!("Movie {}", id_str)))?
.into_domain()?;
@@ -423,7 +419,7 @@ impl DiaryQuery for PostgresDiaryRepository {
.bind(&id_str)
.fetch_all(&self.pool)
.await
.map_err(Self::map_err)?
.map_err(adapter_common::map_sqlx_error)?
.into_iter()
.map(ReviewRow::into_domain)
.collect::<Result<Vec<_>, _>>()?;
@@ -448,7 +444,7 @@ impl DiaryQuery for PostgresDiaryRepository {
.bind(&uid)
.fetch_all(&self.pool)
.await
.map_err(Self::map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
rows.into_iter().map(DiaryRow::into_domain).collect()
}
@@ -477,7 +473,7 @@ impl DiaryQuery for PostgresDiaryRepository {
while let Some(row) = futures::StreamExt::next(&mut rows).await {
yield match row {
Ok(r) => r.into_domain(),
Err(e) => Err(Self::map_err(e)),
Err(e) => Err(adapter_common::map_sqlx_error(e)),
};
}
})
@@ -500,7 +496,7 @@ impl DiaryQuery for PostgresDiaryRepository {
.bind(id_str)
.fetch_one(&self.pool)
.await
.map_err(Self::map_err)
.map_err(adapter_common::map_sqlx_error)
.map(MovieStatsRow::into_domain)
}
@@ -517,7 +513,7 @@ impl DiaryQuery for PostgresDiaryRepository {
.bind(&id_str)
.fetch_one(&self.pool)
.await
.map_err(Self::map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
let rows = sqlx::query_as::<_, FeedRow>(
"SELECT m.id, m.external_metadata_id, m.title, m.release_year, m.director, m.poster_path,
@@ -541,7 +537,7 @@ impl DiaryQuery for PostgresDiaryRepository {
.bind(offset)
.fetch_all(&self.pool)
.await
.map_err(Self::map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
let items = rows
.into_iter()
@@ -561,7 +557,7 @@ impl DiaryQuery for PostgresDiaryRepository {
sqlx::query_scalar("SELECT COUNT(*) FROM reviews WHERE remote_actor_url IS NULL")
.fetch_one(&self.pool)
.await
.map_err(Self::map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
Ok(count as u64)
}
}

View File

@@ -7,7 +7,7 @@ use domain::{
};
use sqlx::{PgPool, Row};
use crate::models::{datetime_to_str, parse_datetime, parse_uuid};
use adapter_common::{datetime_to_str, parse_datetime, parse_uuid};
pub struct PostgresGoalRepository {
pool: PgPool,
@@ -18,10 +18,6 @@ impl PostgresGoalRepository {
Self { pool }
}
fn map_err(e: sqlx::Error) -> DomainError {
tracing::error!("Database error: {:?}", e);
DomainError::InfrastructureError("Database operation failed".into())
}
}
#[async_trait]
@@ -46,7 +42,7 @@ impl GoalCommand for PostgresGoalRepository {
.bind(&created_at)
.execute(&self.pool)
.await
.map_err(Self::map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
Ok(())
}
@@ -60,7 +56,7 @@ impl GoalCommand for PostgresGoalRepository {
.bind(&id)
.execute(&self.pool)
.await
.map_err(Self::map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
if result.rows_affected() == 0 {
return Err(DomainError::NotFound("Goal not found".into()));
@@ -77,7 +73,7 @@ impl GoalCommand for PostgresGoalRepository {
.bind(&uid)
.execute(&self.pool)
.await
.map_err(Self::map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
if result.rows_affected() == 0 {
return Err(DomainError::NotFound("Goal not found".into()));
@@ -105,7 +101,7 @@ impl GoalQuery for PostgresGoalRepository {
.bind(y)
.fetch_optional(&self.pool)
.await
.map_err(Self::map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
row.map(|r| row_to_goal(&r)).transpose()
}
@@ -121,7 +117,7 @@ impl GoalQuery for PostgresGoalRepository {
.bind(&uid)
.fetch_all(&self.pool)
.await
.map_err(Self::map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
rows.iter().map(row_to_goal).collect()
}

View File

@@ -97,10 +97,6 @@ impl PostgresImportProfileRepository {
Self { pool }
}
fn map_err(e: sqlx::Error) -> DomainError {
tracing::error!("DB error: {:?}", e);
DomainError::InfrastructureError("Database operation failed".into())
}
}
#[async_trait]
@@ -118,7 +114,7 @@ impl ImportProfileRepository for PostgresImportProfileRepository {
.execute(&self.pool)
.await
.map(|_| ())
.map_err(Self::map_err)
.map_err(adapter_common::map_sqlx_error)
}
async fn list_for_user(&self, user_id: &UserId) -> Result<Vec<ImportProfile>, DomainError> {
@@ -139,7 +135,7 @@ impl ImportProfileRepository for PostgresImportProfileRepository {
.bind(&uid)
.fetch_all(&self.pool)
.await
.map_err(Self::map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
rows.into_iter()
.map(|r| {
@@ -184,7 +180,7 @@ impl ImportProfileRepository for PostgresImportProfileRepository {
.bind(&id_str).bind(&uid_str)
.fetch_optional(&self.pool)
.await
.map_err(Self::map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
row.map(|r| {
Ok(ImportProfile {
@@ -212,6 +208,6 @@ impl ImportProfileRepository for PostgresImportProfileRepository {
.execute(&self.pool)
.await
.map(|_| ())
.map_err(Self::map_err)
.map_err(adapter_common::map_sqlx_error)
}
}

View File

@@ -202,10 +202,6 @@ impl PostgresImportSessionRepository {
Self { pool }
}
fn map_err(e: sqlx::Error) -> DomainError {
tracing::error!("DB error: {:?}", e);
DomainError::InfrastructureError("Database operation failed".into())
}
fn serialize_session(
s: &ImportSession,
@@ -301,7 +297,7 @@ impl ImportSessionRepository for PostgresImportSessionRepository {
.execute(&self.pool)
.await
.map(|_| ())
.map_err(Self::map_err)
.map_err(adapter_common::map_sqlx_error)
}
async fn get(
@@ -331,7 +327,7 @@ impl ImportSessionRepository for PostgresImportSessionRepository {
.bind(&uid_str)
.fetch_optional(&self.pool)
.await
.map_err(Self::map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
row.map(|r| {
Self::deserialize_session(
@@ -359,7 +355,7 @@ impl ImportSessionRepository for PostgresImportSessionRepository {
.execute(&self.pool)
.await
.map(|_| ())
.map_err(Self::map_err)
.map_err(adapter_common::map_sqlx_error)
}
async fn delete(&self, id: &ImportSessionId) -> Result<(), DomainError> {
@@ -369,14 +365,14 @@ impl ImportSessionRepository for PostgresImportSessionRepository {
.execute(&self.pool)
.await
.map(|_| ())
.map_err(Self::map_err)
.map_err(adapter_common::map_sqlx_error)
}
async fn delete_expired(&self) -> Result<u64, DomainError> {
let result = sqlx::query("DELETE FROM import_sessions WHERE expires_at < NOW()")
.execute(&self.pool)
.await
.map_err(Self::map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
Ok(result.rows_affected())
}
@@ -387,6 +383,6 @@ impl ImportSessionRepository for PostgresImportSessionRepository {
.execute(&self.pool)
.await
.map(|_| ())
.map_err(Self::map_err)
.map_err(adapter_common::map_sqlx_error)
}
}

View File

@@ -39,30 +39,6 @@ pub use watch_event::{PostgresWatchEventRepository, PostgresWebhookTokenReposito
pub use watchlist::PostgresWatchlistRepository;
pub use wrapup::{PostgresWrapUpRepository, PostgresWrapUpStatsQuery};
pub(crate) fn format_year_month(ym: &str) -> String {
let parts: Vec<&str> = ym.splitn(2, '-').collect();
if parts.len() != 2 {
return ym.to_string();
}
let year = parts[0].get(2..).unwrap_or(parts[0]);
let month = match parts[1] {
"01" => "Jan",
"02" => "Feb",
"03" => "Mar",
"04" => "Apr",
"05" => "May",
"06" => "Jun",
"07" => "Jul",
"08" => "Aug",
"09" => "Sep",
"10" => "Oct",
"11" => "Nov",
"12" => "Dec",
_ => parts[1],
};
format!("{} '{}", month, year)
}
pub async fn migrate(pool: &PgPool) -> Result<(), DomainError> {
sqlx::migrate!("./migrations")
.set_ignore_missing(true)

View File

@@ -1,16 +1,11 @@
use chrono::NaiveDateTime;
use domain::{
errors::DomainError,
models::{
DiaryEntry, FeedEntry, Movie, MovieSummary, PersistedReview, Review, ReviewSource,
UserSummary,
},
value_objects::{
Comment, Email, ExternalMetadataId, MovieId, MovieTitle, PosterPath, Rating, ReleaseYear,
ReviewId, UserId, Username,
},
models::{DiaryEntry, FeedEntry, Movie, MovieSummary, Review},
};
use adapter_common::{
movie_row_to_domain, movie_stats_to_domain, movie_summary_to_domain, review_row_to_domain,
user_summary_to_domain,
};
use uuid::Uuid;
#[derive(sqlx::FromRow)]
pub(crate) struct MovieRow {
@@ -24,22 +19,14 @@ pub(crate) struct MovieRow {
impl MovieRow {
pub fn into_domain(self) -> Result<Movie, DomainError> {
let id = MovieId::from_uuid(parse_uuid(&self.id)?);
let external_metadata_id = self
.external_metadata_id
.map(ExternalMetadataId::new)
.transpose()?;
let title = MovieTitle::new(self.title)?;
let release_year = ReleaseYear::new(self.release_year as u16)?;
let poster_path = self.poster_path.map(PosterPath::new).transpose()?;
Ok(Movie::from_persistence(
id,
external_metadata_id,
title,
release_year,
movie_row_to_domain(
self.id,
self.external_metadata_id,
self.title,
self.release_year,
self.director,
poster_path,
))
self.poster_path,
)
}
}
@@ -60,23 +47,22 @@ pub(crate) struct MovieSummaryRow {
impl MovieSummaryRow {
pub fn into_domain(self) -> Result<MovieSummary, DomainError> {
let movie = MovieRow {
id: self.id,
external_metadata_id: self.external_metadata_id,
title: self.title,
release_year: self.release_year,
director: self.director,
poster_path: self.poster_path,
}
.into_domain()?;
Ok(MovieSummary {
let movie = movie_row_to_domain(
self.id,
self.external_metadata_id,
self.title,
self.release_year,
self.director,
self.poster_path,
)?;
Ok(movie_summary_to_domain(
movie,
genres: self.genres.unwrap_or_default(),
runtime_minutes: self.runtime_minutes.map(|v| v as u32),
original_language: self.original_language,
overview: self.overview,
collection_name: self.collection_name,
})
self.genres.unwrap_or_default(),
self.runtime_minutes,
self.original_language,
self.overview,
self.collection_name,
))
}
}
@@ -95,29 +81,17 @@ pub(crate) struct ReviewRow {
impl ReviewRow {
pub fn into_domain(self) -> Result<Review, DomainError> {
let id = ReviewId::from_uuid(parse_uuid(&self.id)?);
let movie_id = MovieId::from_uuid(parse_uuid(&self.movie_id)?);
let user_id = UserId::from_uuid(parse_uuid(&self.user_id)?);
let rating = Rating::new(self.rating as u8)?;
let comment = self.comment.map(Comment::new).transpose()?;
let watched_at = parse_datetime(&self.watched_at)?;
let created_at = parse_datetime(&self.created_at)?;
let source = match self.remote_actor_url {
None => ReviewSource::Local,
Some(url) => ReviewSource::Remote { actor_url: url },
};
let watch_medium = self.watch_medium.map(|s| s.parse()).transpose()?;
Ok(Review::from_persistence(PersistedReview {
id,
movie_id,
user_id,
rating,
comment,
watched_at,
created_at,
source,
watch_medium,
}))
review_row_to_domain(
self.id,
self.movie_id,
self.user_id,
self.rating,
self.comment,
self.watched_at,
self.created_at,
self.remote_actor_url,
self.watch_medium,
)
}
}
@@ -142,27 +116,25 @@ pub(crate) struct DiaryRow {
impl DiaryRow {
pub fn into_domain(self) -> Result<DiaryEntry, DomainError> {
let movie = MovieRow {
id: self.id,
external_metadata_id: self.external_metadata_id,
title: self.title,
release_year: self.release_year,
director: self.director,
poster_path: self.poster_path,
}
.into_domain()?;
let review = ReviewRow {
id: self.review_id,
movie_id: self.movie_id,
user_id: self.user_id,
rating: self.rating,
comment: self.comment,
watched_at: self.watched_at,
created_at: self.created_at,
remote_actor_url: self.remote_actor_url,
watch_medium: self.watch_medium,
}
.into_domain()?;
let movie = movie_row_to_domain(
self.id,
self.external_metadata_id,
self.title,
self.release_year,
self.director,
self.poster_path,
)?;
let review = review_row_to_domain(
self.review_id,
self.movie_id,
self.user_id,
self.rating,
self.comment,
self.watched_at,
self.created_at,
self.remote_actor_url,
self.watch_medium,
)?;
Ok(DiaryEntry::new(movie, review))
}
}
@@ -189,24 +161,26 @@ pub(crate) struct FeedRow {
impl FeedRow {
pub fn into_domain(self) -> Result<FeedEntry, DomainError> {
let diary = DiaryRow {
id: self.id,
external_metadata_id: self.external_metadata_id,
title: self.title,
release_year: self.release_year,
director: self.director,
poster_path: self.poster_path,
review_id: self.review_id,
movie_id: self.movie_id,
user_id: self.user_id,
rating: self.rating,
comment: self.comment,
watched_at: self.watched_at,
created_at: self.created_at,
remote_actor_url: self.remote_actor_url,
watch_medium: self.watch_medium,
}
.into_domain()?;
let movie = movie_row_to_domain(
self.id,
self.external_metadata_id,
self.title,
self.release_year,
self.director,
self.poster_path,
)?;
let review = review_row_to_domain(
self.review_id,
self.movie_id,
self.user_id,
self.rating,
self.comment,
self.watched_at,
self.created_at,
self.remote_actor_url,
self.watch_medium,
)?;
let diary = DiaryEntry::new(movie, review);
Ok(FeedEntry::new(diary, self.user_email))
}
}
@@ -225,18 +199,18 @@ pub(crate) struct MovieStatsRow {
impl MovieStatsRow {
pub fn into_domain(self) -> domain::models::MovieStats {
domain::models::MovieStats {
total_count: self.total_count as u64,
avg_rating: self.avg_rating,
federated_count: self.federated_count as u64,
rating_histogram: [
self.rating_1 as u64,
self.rating_2 as u64,
self.rating_3 as u64,
self.rating_4 as u64,
self.rating_5 as u64,
movie_stats_to_domain(
self.total_count,
self.avg_rating,
self.federated_count,
[
self.rating_1,
self.rating_2,
self.rating_3,
self.rating_4,
self.rating_5,
],
}
)
}
}
@@ -252,16 +226,16 @@ pub(crate) struct UserSummaryRow {
}
impl UserSummaryRow {
pub fn into_domain(self) -> Result<UserSummary, DomainError> {
Ok(UserSummary::new(
UserId::from_uuid(parse_uuid(&self.id)?),
Email::new(self.email)?,
Username::new(self.username)?,
pub fn into_domain(self) -> Result<domain::models::UserSummary, DomainError> {
user_summary_to_domain(
self.id,
self.email,
self.username,
self.display_name,
self.total_movies,
self.avg_rating,
self.avatar_path,
))
)
}
}
@@ -283,17 +257,3 @@ pub(crate) struct MonthlyRatingRow {
pub avg_rating: f64,
pub count: i64,
}
pub(crate) fn parse_uuid(s: &str) -> Result<Uuid, DomainError> {
Uuid::parse_str(s)
.map_err(|e| DomainError::InfrastructureError(format!("Invalid UUID '{}': {}", s, e)))
}
pub(crate) fn datetime_to_str(dt: &NaiveDateTime) -> String {
dt.format("%Y-%m-%d %H:%M:%S").to_string()
}
pub(crate) fn parse_datetime(s: &str) -> Result<NaiveDateTime, DomainError> {
NaiveDateTime::parse_from_str(s, "%Y-%m-%d %H:%M:%S")
.map_err(|e| DomainError::InfrastructureError(format!("Invalid datetime '{}': {}", s, e)))
}

View File

@@ -21,10 +21,6 @@ impl PostgresMovieRepository {
Self { pool }
}
fn map_err(e: sqlx::Error) -> DomainError {
tracing::error!("Database error: {:?}", e);
DomainError::InfrastructureError("Database operation failed".into())
}
}
#[async_trait]
@@ -55,7 +51,7 @@ impl MovieCommand for PostgresMovieRepository {
.bind(&poster_path)
.execute(&self.pool)
.await
.map_err(Self::map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
Ok(())
}
@@ -66,7 +62,7 @@ impl MovieCommand for PostgresMovieRepository {
.bind(&id)
.execute(&self.pool)
.await
.map_err(Self::map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
Ok(())
}
}
@@ -85,7 +81,7 @@ impl MovieQuery for PostgresMovieRepository {
.bind(id)
.fetch_optional(&self.pool)
.await
.map_err(Self::map_err)?
.map_err(adapter_common::map_sqlx_error)?
.map(MovieRow::into_domain)
.transpose()
}
@@ -99,7 +95,7 @@ impl MovieQuery for PostgresMovieRepository {
.bind(&id)
.fetch_optional(&self.pool)
.await
.map_err(Self::map_err)?
.map_err(adapter_common::map_sqlx_error)?
.map(MovieRow::into_domain)
.transpose()
}
@@ -119,7 +115,7 @@ impl MovieQuery for PostgresMovieRepository {
.bind(year)
.fetch_all(&self.pool)
.await
.map_err(Self::map_err)?
.map_err(adapter_common::map_sqlx_error)?
.into_iter()
.map(MovieRow::into_domain)
.collect()
@@ -139,7 +135,7 @@ impl MovieQuery for PostgresMovieRepository {
.bind(&vals)
.fetch_all(&self.pool)
.await
.map_err(Self::map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
Ok(rows.into_iter().map(|(id,)| id).collect())
}
@@ -162,7 +158,7 @@ impl MovieQuery for PostgresMovieRepository {
.bind(&years)
.fetch_all(&self.pool)
.await
.map_err(Self::map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
Ok(rows
.into_iter()
.map(|r| {
@@ -211,7 +207,7 @@ impl MovieQuery for PostgresMovieRepository {
.bind(offset)
.fetch_all(&self.pool)
.await
.map_err(Self::map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
let total: i64 = sqlx::query(
"SELECT COUNT(DISTINCT m.id) \
@@ -226,7 +222,7 @@ impl MovieQuery for PostgresMovieRepository {
.bind(genre)
.fetch_one(&self.pool)
.await
.map_err(Self::map_err)?
.map_err(adapter_common::map_sqlx_error)?
.try_get(0)
.unwrap_or(0);
@@ -250,7 +246,7 @@ impl MovieQuery for PostgresMovieRepository {
)
.fetch_all(&self.pool)
.await
.map_err(Self::map_err)?
.map_err(adapter_common::map_sqlx_error)?
.into_iter()
.map(|r| r.into_domain())
.collect()

View File

@@ -13,10 +13,6 @@ impl PostgresMovieDeduplicator {
Self { pool }
}
fn map_err(e: sqlx::Error) -> DomainError {
tracing::error!("Database error: {:?}", e);
DomainError::InfrastructureError("Database operation failed".into())
}
}
#[async_trait]
@@ -36,7 +32,7 @@ impl MovieDeduplicator for PostgresMovieDeduplicator {
let director = canonical.director().map(str::to_string);
let poster = canonical.poster_path().map(|p| p.value().to_string());
let mut tx = self.pool.begin().await.map_err(Self::map_err)?;
let mut tx = self.pool.begin().await.map_err(adapter_common::map_sqlx_error)?;
// 1. Upsert canonical movie record
sqlx::query(
@@ -47,7 +43,7 @@ impl MovieDeduplicator for PostgresMovieDeduplicator {
poster_path = COALESCE(EXCLUDED.poster_path, movies.poster_path)",
)
.bind(&new).bind(&ext_id).bind(&title).bind(year).bind(&director).bind(&poster)
.execute(&mut *tx).await.map_err(Self::map_err)?;
.execute(&mut *tx).await.map_err(adapter_common::map_sqlx_error)?;
// 2. Re-point simple FK tables
let reviews = sqlx::query("UPDATE reviews SET movie_id = $1 WHERE movie_id = $2")
@@ -55,7 +51,7 @@ impl MovieDeduplicator for PostgresMovieDeduplicator {
.bind(&old)
.execute(&mut *tx)
.await
.map_err(Self::map_err)?
.map_err(adapter_common::map_sqlx_error)?
.rows_affected();
let watchlist =
@@ -64,7 +60,7 @@ impl MovieDeduplicator for PostgresMovieDeduplicator {
.bind(&old)
.execute(&mut *tx)
.await
.map_err(Self::map_err)?
.map_err(adapter_common::map_sqlx_error)?
.rows_affected();
let watch_events = sqlx::query("UPDATE watch_events SET movie_id = $1 WHERE movie_id = $2")
@@ -72,7 +68,7 @@ impl MovieDeduplicator for PostgresMovieDeduplicator {
.bind(&old)
.execute(&mut *tx)
.await
.map_err(Self::map_err)?
.map_err(adapter_common::map_sqlx_error)?
.rows_affected();
// 3. Re-point movie_profiles (PK — move only if canonical has none)
@@ -81,7 +77,7 @@ impl MovieDeduplicator for PostgresMovieDeduplicator {
.bind(&old)
.execute(&mut *tx)
.await
.map_err(Self::map_err)?
.map_err(adapter_common::map_sqlx_error)?
.rows_affected();
// 4. Re-point enrichment tables with composite PKs (INSERT … ON CONFLICT DO NOTHING + DELETE)
@@ -95,12 +91,12 @@ impl MovieDeduplicator for PostgresMovieDeduplicator {
.bind(&old)
.execute(&mut *tx)
.await
.map_err(Self::map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
sqlx::query("DELETE FROM movie_genres WHERE movie_id = $1")
.bind(&old)
.execute(&mut *tx)
.await
.map_err(Self::map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
sqlx::query(
"INSERT INTO movie_keywords (movie_id, tmdb_id, name)
@@ -111,43 +107,43 @@ impl MovieDeduplicator for PostgresMovieDeduplicator {
.bind(&old)
.execute(&mut *tx)
.await
.map_err(Self::map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
sqlx::query("DELETE FROM movie_keywords WHERE movie_id = $1")
.bind(&old)
.execute(&mut *tx)
.await
.map_err(Self::map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
sqlx::query(
"INSERT INTO movie_cast (movie_id, tmdb_person_id, name, character, billing_order, profile_path)
SELECT $1, tmdb_person_id, name, character, billing_order, profile_path FROM movie_cast WHERE movie_id = $2
ON CONFLICT DO NOTHING",
).bind(&new).bind(&old).execute(&mut *tx).await.map_err(Self::map_err)?;
).bind(&new).bind(&old).execute(&mut *tx).await.map_err(adapter_common::map_sqlx_error)?;
sqlx::query("DELETE FROM movie_cast WHERE movie_id = $1")
.bind(&old)
.execute(&mut *tx)
.await
.map_err(Self::map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
sqlx::query(
"INSERT INTO movie_crew (movie_id, tmdb_person_id, name, job, department, profile_path)
SELECT $1, tmdb_person_id, name, job, department, profile_path FROM movie_crew WHERE movie_id = $2
ON CONFLICT DO NOTHING",
).bind(&new).bind(&old).execute(&mut *tx).await.map_err(Self::map_err)?;
).bind(&new).bind(&old).execute(&mut *tx).await.map_err(adapter_common::map_sqlx_error)?;
sqlx::query("DELETE FROM movie_crew WHERE movie_id = $1")
.bind(&old)
.execute(&mut *tx)
.await
.map_err(Self::map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
// 5. Delete the now-empty old movie record (remaining cascades are safe: all FKs cleared above)
sqlx::query("DELETE FROM movies WHERE id = $1")
.bind(&old)
.execute(&mut *tx)
.await
.map_err(Self::map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
tx.commit().await.map_err(Self::map_err)?;
tx.commit().await.map_err(adapter_common::map_sqlx_error)?;
Ok(reviews + watchlist + watch_events + profiles)
}

View File

@@ -29,9 +29,6 @@ pub fn create_person_adapter(pool: PgPool) -> (Arc<dyn PersonCommand>, Arc<dyn P
)
}
fn map_err(e: sqlx::Error) -> DomainError {
DomainError::InfrastructureError(e.to_string())
}
#[async_trait]
impl PersonCommand for PostgresPersonAdapter {
@@ -56,7 +53,7 @@ impl PersonCommand for PostgresPersonAdapter {
.bind(person.profile_path())
.execute(&self.pool)
.await
.map_err(map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
}
Ok(())
}
@@ -89,7 +86,7 @@ impl PersonCommand for PostgresPersonAdapter {
.bind(batch_size as i64)
.fetch_all(&self.pool)
.await
.map_err(map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
let has_more = rows.len() as u32 >= batch_size;
let mut count = 0u64;
@@ -109,7 +106,7 @@ impl PersonCommand for PostgresPersonAdapter {
.bind(&row.profile_path)
.execute(&self.pool)
.await
.map_err(map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
count += 1;
}
Ok((count, has_more))
@@ -137,7 +134,7 @@ impl PersonCommand for PostgresPersonAdapter {
.bind(id.value().to_string())
.execute(&self.pool)
.await
.map_err(map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
Ok(())
}
}
@@ -151,7 +148,7 @@ impl PersonQuery for PostgresPersonAdapter {
.bind(id.value().to_string())
.fetch_optional(&self.pool)
.await
.map_err(map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
Ok(row.map(PersonRow::into_person))
}
@@ -166,7 +163,7 @@ impl PersonQuery for PostgresPersonAdapter {
.bind(id.value())
.fetch_optional(&self.pool)
.await
.map_err(map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
Ok(row.map(PersonRow::into_person))
}
@@ -182,7 +179,7 @@ impl PersonQuery for PostgresPersonAdapter {
.bind(id.value().to_string())
.fetch_optional(&self.pool)
.await
.map_err(map_err)?
.map_err(adapter_common::map_sqlx_error)?
.flatten();
let Some(tmdb_id) = tmdb_id else {
@@ -219,7 +216,7 @@ impl PersonQuery for PostgresPersonAdapter {
.bind(tmdb_id)
.fetch_all(&self.pool)
.await
.map_err(map_err)?
.map_err(adapter_common::map_sqlx_error)?
.into_iter()
.map(|r| CastCredit {
movie_id: MovieId::from_uuid(uuid::Uuid::parse_str(&r.id).unwrap_or_default()),
@@ -238,7 +235,7 @@ impl PersonQuery for PostgresPersonAdapter {
.bind(tmdb_id)
.fetch_all(&self.pool)
.await
.map_err(map_err)?
.map_err(adapter_common::map_sqlx_error)?
.into_iter()
.map(|r| CrewCredit {
movie_id: MovieId::from_uuid(uuid::Uuid::parse_str(&r.id).unwrap_or_default()),
@@ -261,7 +258,7 @@ impl PersonQuery for PostgresPersonAdapter {
.bind(offset as i64)
.fetch_all(&self.pool)
.await
.map_err(map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
Ok(rows.into_iter().map(PersonRow::into_person).collect())
}
@@ -279,7 +276,7 @@ impl PersonQuery for PostgresPersonAdapter {
)
.fetch_all(&self.pool)
.await
.map_err(map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
Ok(rows
.into_iter()

View File

@@ -17,10 +17,6 @@ impl PostgresMovieProfileRepository {
Self { pool }
}
fn map_err(e: sqlx::Error) -> DomainError {
tracing::error!("Database error: {:?}", e);
DomainError::InfrastructureError("Database operation failed".into())
}
}
#[async_trait]
@@ -28,7 +24,7 @@ impl MovieProfileRepository for PostgresMovieProfileRepository {
async fn upsert(&self, p: &MovieProfile) -> Result<(), DomainError> {
let movie_id = p.movie_id.value().to_string();
let mut tx = self.pool.begin().await.map_err(Self::map_err)?;
let mut tx = self.pool.begin().await.map_err(adapter_common::map_sqlx_error)?;
sqlx::query(
r#"INSERT INTO movie_profiles
@@ -61,35 +57,35 @@ impl MovieProfileRepository for PostgresMovieProfileRepository {
.bind(p.enriched_at)
.execute(&mut *tx)
.await
.map_err(Self::map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
sqlx::query("DELETE FROM movie_genres WHERE movie_id = $1")
.bind(&movie_id)
.execute(&mut *tx)
.await
.map_err(Self::map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
for g in &p.genres {
sqlx::query("INSERT INTO movie_genres (movie_id, tmdb_id, name) VALUES ($1,$2,$3) ON CONFLICT DO NOTHING")
.bind(&movie_id).bind(g.tmdb_id as i32).bind(&g.name)
.execute(&mut *tx).await.map_err(Self::map_err)?;
.execute(&mut *tx).await.map_err(adapter_common::map_sqlx_error)?;
}
sqlx::query("DELETE FROM movie_keywords WHERE movie_id = $1")
.bind(&movie_id)
.execute(&mut *tx)
.await
.map_err(Self::map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
for k in &p.keywords {
sqlx::query("INSERT INTO movie_keywords (movie_id, tmdb_id, name) VALUES ($1,$2,$3) ON CONFLICT DO NOTHING")
.bind(&movie_id).bind(k.tmdb_id as i32).bind(&k.name)
.execute(&mut *tx).await.map_err(Self::map_err)?;
.execute(&mut *tx).await.map_err(adapter_common::map_sqlx_error)?;
}
sqlx::query("DELETE FROM movie_cast WHERE movie_id = $1")
.bind(&movie_id)
.execute(&mut *tx)
.await
.map_err(Self::map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
for c in &p.cast {
sqlx::query(
"INSERT INTO movie_cast \
@@ -104,14 +100,14 @@ impl MovieProfileRepository for PostgresMovieProfileRepository {
.bind(&c.profile_path)
.execute(&mut *tx)
.await
.map_err(Self::map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
}
sqlx::query("DELETE FROM movie_crew WHERE movie_id = $1")
.bind(&movie_id)
.execute(&mut *tx)
.await
.map_err(Self::map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
for cr in &p.crew {
sqlx::query(
"INSERT INTO movie_crew \
@@ -126,10 +122,10 @@ impl MovieProfileRepository for PostgresMovieProfileRepository {
.bind(&cr.profile_path)
.execute(&mut *tx)
.await
.map_err(Self::map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
}
tx.commit().await.map_err(Self::map_err)
tx.commit().await.map_err(adapter_common::map_sqlx_error)
}
async fn get_by_movie_id(&self, id: &MovieId) -> Result<Option<MovieProfile>, DomainError> {
@@ -144,7 +140,7 @@ impl MovieProfileRepository for PostgresMovieProfileRepository {
.bind(&movie_id)
.fetch_optional(&self.pool)
.await
.map_err(Self::map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
let row = match row {
Some(r) => r,
@@ -159,7 +155,7 @@ impl MovieProfileRepository for PostgresMovieProfileRepository {
.bind(&movie_id)
.fetch_all(&self.pool)
.await
.map_err(Self::map_err)?
.map_err(adapter_common::map_sqlx_error)?
.into_iter()
.map(|r| Genre {
tmdb_id: r.try_get::<i32, _>("tmdb_id").unwrap_or(0) as u32,
@@ -171,7 +167,7 @@ impl MovieProfileRepository for PostgresMovieProfileRepository {
.bind(&movie_id)
.fetch_all(&self.pool)
.await
.map_err(Self::map_err)?
.map_err(adapter_common::map_sqlx_error)?
.into_iter()
.map(|r| Keyword {
tmdb_id: r.try_get::<i32, _>("tmdb_id").unwrap_or(0) as u32,
@@ -186,7 +182,7 @@ impl MovieProfileRepository for PostgresMovieProfileRepository {
.bind(&movie_id)
.fetch_all(&self.pool)
.await
.map_err(Self::map_err)?
.map_err(adapter_common::map_sqlx_error)?
.into_iter()
.map(|r| CastMember {
tmdb_person_id: r.try_get::<i64, _>("tmdb_person_id").unwrap_or(0) as u64,
@@ -204,7 +200,7 @@ impl MovieProfileRepository for PostgresMovieProfileRepository {
.bind(&movie_id)
.fetch_all(&self.pool)
.await
.map_err(Self::map_err)?
.map_err(adapter_common::map_sqlx_error)?
.into_iter()
.map(|r| CrewMember {
tmdb_person_id: r.try_get::<i64, _>("tmdb_person_id").unwrap_or(0) as u64,
@@ -257,7 +253,7 @@ impl MovieProfileRepository for PostgresMovieProfileRepository {
.bind(threshold)
.fetch_all(&self.pool)
.await
.map_err(Self::map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
Ok(rows
.into_iter()

View File

@@ -16,9 +16,6 @@ impl PostgresRefreshSessionAdapter {
}
}
fn map_err(e: sqlx::Error) -> DomainError {
DomainError::InfrastructureError(e.to_string())
}
#[async_trait]
impl RefreshSessionRepository for PostgresRefreshSessionAdapter {
@@ -34,7 +31,7 @@ impl RefreshSessionRepository for PostgresRefreshSessionAdapter {
.bind(session.created_at)
.execute(&self.pool)
.await
.map_err(map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
Ok(())
}
@@ -48,7 +45,7 @@ impl RefreshSessionRepository for PostgresRefreshSessionAdapter {
.bind(token)
.fetch_optional(&self.pool)
.await
.map_err(map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
row.map(RefreshSessionRow::into_domain).transpose()
}
@@ -58,7 +55,7 @@ impl RefreshSessionRepository for PostgresRefreshSessionAdapter {
.bind(token)
.execute(&self.pool)
.await
.map_err(map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
Ok(())
}
@@ -67,7 +64,7 @@ impl RefreshSessionRepository for PostgresRefreshSessionAdapter {
.bind(user_id.value().to_string())
.execute(&self.pool)
.await
.map_err(map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
Ok(())
}
@@ -75,7 +72,7 @@ impl RefreshSessionRepository for PostgresRefreshSessionAdapter {
let result = sqlx::query("DELETE FROM refresh_sessions WHERE expires_at < NOW()")
.execute(&self.pool)
.await
.map_err(map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
Ok(result.rows_affected())
}
}

View File

@@ -7,7 +7,8 @@ use domain::{
};
use sqlx::PgPool;
use crate::models::{ReviewRow, datetime_to_str};
use adapter_common::datetime_to_str;
use crate::models::ReviewRow;
pub struct PostgresReviewRepository {
pool: PgPool,
@@ -18,10 +19,6 @@ impl PostgresReviewRepository {
Self { pool }
}
fn map_err(e: sqlx::Error) -> DomainError {
tracing::error!("Database error: {:?}", e);
DomainError::InfrastructureError("Database operation failed".into())
}
}
#[async_trait]
@@ -54,7 +51,7 @@ impl ReviewRepository for PostgresReviewRepository {
.bind(review.watch_medium().map(|wm| wm.to_string()))
.execute(&self.pool)
.await
.map_err(Self::map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
Ok(())
}
@@ -72,7 +69,7 @@ impl ReviewRepository for PostgresReviewRepository {
.bind(&id)
.fetch_optional(&self.pool)
.await
.map_err(Self::map_err)?
.map_err(adapter_common::map_sqlx_error)?
.map(ReviewRow::into_domain)
.transpose()
}
@@ -94,7 +91,7 @@ impl ReviewRepository for PostgresReviewRepository {
.bind(&id)
.execute(&self.pool)
.await
.map_err(Self::map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
Ok(())
}
@@ -105,7 +102,7 @@ impl ReviewRepository for PostgresReviewRepository {
.bind(&id)
.execute(&self.pool)
.await
.map_err(Self::map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
Ok(())
}
@@ -122,7 +119,7 @@ impl ReviewRepository for PostgresReviewRepository {
.bind(&uid)
.fetch_all(&self.pool)
.await
.map_err(Self::map_err)?
.map_err(adapter_common::map_sqlx_error)?
.into_iter()
.map(ReviewRow::into_domain)
.collect()

View File

@@ -7,7 +7,7 @@ use domain::{
};
use sqlx::PgPool;
use crate::format_year_month;
use adapter_common::format_year_month;
use crate::models::{DirectorCountRow, MonthlyRatingRow, UserTotalsRow};
pub struct PostgresStatsRepository {
@@ -19,10 +19,6 @@ impl PostgresStatsRepository {
Self { pool }
}
fn map_err(e: sqlx::Error) -> DomainError {
tracing::error!("Database error: {:?}", e);
DomainError::InfrastructureError("Database operation failed".into())
}
async fn fetch_user_totals(&self, user_id: &str) -> Result<UserTotalsRow, DomainError> {
sqlx::query_as::<_, UserTotalsRow>(
@@ -33,7 +29,7 @@ impl PostgresStatsRepository {
.bind(user_id)
.fetch_one(&self.pool)
.await
.map_err(Self::map_err)
.map_err(adapter_common::map_sqlx_error)
}
async fn fetch_user_favorite_director(
@@ -52,7 +48,7 @@ impl PostgresStatsRepository {
.bind(user_id)
.fetch_optional(&self.pool)
.await
.map_err(Self::map_err)
.map_err(adapter_common::map_sqlx_error)
}
async fn fetch_user_most_active_month(
@@ -70,7 +66,7 @@ impl PostgresStatsRepository {
.bind(user_id)
.fetch_optional(&self.pool)
.await
.map_err(Self::map_err)
.map_err(adapter_common::map_sqlx_error)
}
}
@@ -126,7 +122,7 @@ impl StatsRepository for PostgresStatsRepository {
.bind(&uid)
.fetch_all(&self.pool)
)
.map_err(Self::map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
let max_director_count = director_rows.iter().map(|d| d.count).max().unwrap_or(1);

View File

@@ -16,10 +16,6 @@ impl PostgresUserSettingsRepository {
Self { pool }
}
fn map_err(e: sqlx::Error) -> DomainError {
tracing::error!("Database error: {:?}", e);
DomainError::InfrastructureError("Database operation failed".into())
}
}
#[async_trait]
@@ -33,7 +29,7 @@ impl UserSettingsRepository for PostgresUserSettingsRepository {
.bind(&uid)
.fetch_optional(&self.pool)
.await
.map_err(Self::map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
match row {
Some(r) => {
@@ -65,7 +61,7 @@ impl UserSettingsRepository for PostgresUserSettingsRepository {
.bind(settings.federate_watchlist())
.execute(&self.pool)
.await
.map_err(Self::map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
Ok(())
}
}
@@ -81,7 +77,7 @@ impl UserFederationSettingsQuery for PostgresUserSettingsRepository {
.bind(&uid)
.fetch_optional(&self.pool)
.await
.map_err(Self::map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
match row {
Some(r) => {

View File

@@ -20,10 +20,6 @@ impl PostgresUserRepository {
Self { pool }
}
fn map_err(e: sqlx::Error) -> DomainError {
tracing::error!("Database error: {:?}", e);
DomainError::InfrastructureError("Database operation failed".into())
}
fn parse_role(s: &str) -> UserRole {
match s {
@@ -76,7 +72,7 @@ impl UserRepository for PostgresUserRepository {
.bind(email_str)
.fetch_optional(&self.pool)
.await
.map_err(Self::map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
row.as_ref()
.map(|r| Self::row_to_user(r, vec![]))
.transpose()
@@ -90,7 +86,7 @@ impl UserRepository for PostgresUserRepository {
.bind(username_str)
.fetch_optional(&self.pool)
.await
.map_err(Self::map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
row.as_ref()
.map(|r| Self::row_to_user(r, vec![]))
.transpose()
@@ -130,7 +126,7 @@ impl UserRepository for PostgresUserRepository {
.bind(role)
.execute(&self.pool)
.await
.map_err(Self::map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
Ok(())
}
@@ -140,7 +136,7 @@ impl UserRepository for PostgresUserRepository {
.bind(&id_str)
.fetch_optional(&self.pool)
.await
.map_err(Self::map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
let Some(r) = row else { return Ok(None) };
@@ -197,7 +193,7 @@ impl UserRepository for PostgresUserRepository {
)
.fetch_all(&self.pool)
.await
.map_err(Self::map_err)?
.map_err(adapter_common::map_sqlx_error)?
.into_iter()
.map(UserSummaryRow::into_domain)
.collect()

View File

@@ -7,12 +7,8 @@ use domain::{
};
use sqlx::{PgPool, Row};
use crate::models::{parse_datetime, parse_uuid};
use adapter_common::{parse_datetime, parse_uuid};
fn map_err(e: sqlx::Error) -> DomainError {
tracing::error!("Database error: {:?}", e);
DomainError::InfrastructureError("Database operation failed".into())
}
// ── WatchEventRepository ──────────────────────────────────────────────────────
@@ -52,7 +48,7 @@ impl WatchEventCommand for PostgresWatchEventRepository {
.bind(event.created_at())
.execute(&self.pool)
.await
.map_err(map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
Ok(())
}
@@ -70,7 +66,7 @@ impl WatchEventCommand for PostgresWatchEventRepository {
.bind(&id_str)
.execute(&self.pool)
.await
.map_err(map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
Ok(())
}
@@ -90,7 +86,7 @@ impl WatchEventCommand for PostgresWatchEventRepository {
.bind(&id_strs)
.execute(&self.pool)
.await
.map_err(map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
Ok(result.rows_affected())
}
@@ -103,7 +99,7 @@ impl WatchEventCommand for PostgresWatchEventRepository {
.bind(before)
.execute(&self.pool)
.await
.map_err(map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
Ok(result.rows_affected())
}
}
@@ -126,7 +122,7 @@ impl WatchEventQuery for PostgresWatchEventRepository {
.bind(&uid)
.fetch_all(&self.pool)
.await
.map_err(map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
rows.iter().map(row_to_watch_event).collect()
}
@@ -145,7 +141,7 @@ impl WatchEventQuery for PostgresWatchEventRepository {
.bind(&id_str)
.fetch_optional(&self.pool)
.await
.map_err(map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
row.as_ref().map(row_to_watch_event).transpose()
}
@@ -166,7 +162,7 @@ impl WatchEventQuery for PostgresWatchEventRepository {
.bind(&id_strs)
.fetch_all(&self.pool)
.await
.map_err(map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
rows.iter().map(row_to_watch_event).collect()
}
@@ -187,23 +183,23 @@ impl WatchEventQuery for PostgresWatchEventRepository {
.bind(after)
.fetch_one(&self.pool)
.await
.map_err(map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
Ok(count > 0)
}
}
fn row_to_watch_event(row: &sqlx::postgres::PgRow) -> Result<WatchEvent, DomainError> {
let id_str: String = row.try_get("id").map_err(map_err)?;
let user_id_str: String = row.try_get("user_id").map_err(map_err)?;
let movie_id_str: Option<String> = row.try_get("movie_id").map_err(map_err)?;
let title: String = row.try_get("title").map_err(map_err)?;
let year: Option<i32> = row.try_get("year").map_err(map_err)?;
let ext_id: Option<String> = row.try_get("external_metadata_id").map_err(map_err)?;
let source_str: String = row.try_get("source").map_err(map_err)?;
let watched_at_str: String = row.try_get("watched_at").map_err(map_err)?;
let status_str: String = row.try_get("status").map_err(map_err)?;
let created_at_str: String = row.try_get("created_at").map_err(map_err)?;
let id_str: String = row.try_get("id").map_err(adapter_common::map_sqlx_error)?;
let user_id_str: String = row.try_get("user_id").map_err(adapter_common::map_sqlx_error)?;
let movie_id_str: Option<String> = row.try_get("movie_id").map_err(adapter_common::map_sqlx_error)?;
let title: String = row.try_get("title").map_err(adapter_common::map_sqlx_error)?;
let year: Option<i32> = row.try_get("year").map_err(adapter_common::map_sqlx_error)?;
let ext_id: Option<String> = row.try_get("external_metadata_id").map_err(adapter_common::map_sqlx_error)?;
let source_str: String = row.try_get("source").map_err(adapter_common::map_sqlx_error)?;
let watched_at_str: String = row.try_get("watched_at").map_err(adapter_common::map_sqlx_error)?;
let status_str: String = row.try_get("status").map_err(adapter_common::map_sqlx_error)?;
let created_at_str: String = row.try_get("created_at").map_err(adapter_common::map_sqlx_error)?;
let source: WatchEventSource = source_str
.parse()
@@ -265,7 +261,7 @@ impl WebhookTokenRepository for PostgresWebhookTokenRepository {
.bind(token.last_used_at())
.execute(&self.pool)
.await
.map_err(map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
Ok(())
}
@@ -280,7 +276,7 @@ impl WebhookTokenRepository for PostgresWebhookTokenRepository {
.bind(hash)
.fetch_optional(&self.pool)
.await
.map_err(map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
row.as_ref().map(row_to_webhook_token).transpose()
}
@@ -297,7 +293,7 @@ impl WebhookTokenRepository for PostgresWebhookTokenRepository {
.bind(&uid)
.fetch_all(&self.pool)
.await
.map_err(map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
rows.iter().map(row_to_webhook_token).collect()
}
@@ -311,7 +307,7 @@ impl WebhookTokenRepository for PostgresWebhookTokenRepository {
.bind(&uid)
.execute(&self.pool)
.await
.map_err(map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
if result.rows_affected() == 0 {
return Err(DomainError::NotFound(format!("Webhook token {id_str}")));
@@ -326,20 +322,20 @@ impl WebhookTokenRepository for PostgresWebhookTokenRepository {
.bind(&id_str)
.execute(&self.pool)
.await
.map_err(map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
Ok(())
}
}
fn row_to_webhook_token(row: &sqlx::postgres::PgRow) -> Result<WebhookToken, DomainError> {
let id_str: String = row.try_get("id").map_err(map_err)?;
let user_id_str: String = row.try_get("user_id").map_err(map_err)?;
let token_hash: String = row.try_get("token_hash").map_err(map_err)?;
let provider_str: String = row.try_get("provider").map_err(map_err)?;
let label: Option<String> = row.try_get("label").map_err(map_err)?;
let created_at_str: String = row.try_get("created_at").map_err(map_err)?;
let last_used_str: Option<String> = row.try_get("last_used_at").map_err(map_err)?;
let id_str: String = row.try_get("id").map_err(adapter_common::map_sqlx_error)?;
let user_id_str: String = row.try_get("user_id").map_err(adapter_common::map_sqlx_error)?;
let token_hash: String = row.try_get("token_hash").map_err(adapter_common::map_sqlx_error)?;
let provider_str: String = row.try_get("provider").map_err(adapter_common::map_sqlx_error)?;
let label: Option<String> = row.try_get("label").map_err(adapter_common::map_sqlx_error)?;
let created_at_str: String = row.try_get("created_at").map_err(adapter_common::map_sqlx_error)?;
let last_used_str: Option<String> = row.try_get("last_used_at").map_err(adapter_common::map_sqlx_error)?;
let provider: WatchEventSource = provider_str
.parse()

View File

@@ -10,7 +10,8 @@ use domain::{
};
use sqlx::{PgPool, Row};
use crate::models::{MovieRow, parse_datetime, parse_uuid};
use adapter_common::{parse_datetime, parse_uuid};
use crate::models::MovieRow;
pub struct PostgresWatchlistRepository {
pool: PgPool,
@@ -21,10 +22,6 @@ impl PostgresWatchlistRepository {
Self { pool }
}
fn map_err(e: sqlx::Error) -> DomainError {
tracing::error!("Database error: {:?}", e);
DomainError::InfrastructureError("Database operation failed".into())
}
}
#[async_trait]
@@ -46,7 +43,7 @@ impl WatchlistRepository for PostgresWatchlistRepository {
.bind(added_at)
.execute(&self.pool)
.await
.map_err(Self::map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
Ok(())
}
@@ -61,7 +58,7 @@ impl WatchlistRepository for PostgresWatchlistRepository {
.bind(&mid)
.execute(&self.pool)
.await
.map_err(Self::map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
if result.rows_affected() == 0 {
return Err(DomainError::NotFound(format!(
@@ -85,7 +82,7 @@ impl WatchlistRepository for PostgresWatchlistRepository {
.bind(&mid)
.execute(&self.pool)
.await
.map_err(Self::map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
Ok(result.rows_affected() > 0)
}
@@ -114,14 +111,14 @@ impl WatchlistRepository for PostgresWatchlistRepository {
.bind(offset)
.fetch_all(&self.pool)
.await
.map_err(Self::map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
let total: i64 =
sqlx::query_scalar("SELECT COUNT(*) FROM watchlist_entries WHERE user_id = $1")
.bind(&uid)
.fetch_one(&self.pool)
.await
.map_err(Self::map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
let items = rows
.into_iter()
@@ -187,7 +184,7 @@ impl WatchlistRepository for PostgresWatchlistRepository {
.bind(&mid)
.fetch_one(&self.pool)
.await
.map_err(Self::map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
Ok(count > 0)
}
}

View File

@@ -13,12 +13,8 @@ use domain::{
use sqlx::{PgPool, Row};
use uuid::Uuid;
use crate::models::{parse_datetime, parse_uuid};
use adapter_common::{parse_datetime, parse_uuid};
fn map_err(e: sqlx::Error) -> DomainError {
tracing::error!("Database error: {:?}", e);
DomainError::InfrastructureError("Database operation failed".into())
}
fn status_to_str(s: &WrapUpStatus) -> &'static str {
match s {
@@ -76,7 +72,7 @@ impl WrapUpRepository for PostgresWrapUpRepository {
.bind(record.completed_at)
.execute(&self.pool)
.await
.map_err(map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
Ok(())
}
@@ -96,7 +92,7 @@ impl WrapUpRepository for PostgresWrapUpRepository {
.bind(&id_str)
.execute(&self.pool)
.await
.map_err(map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
Ok(())
}
@@ -115,7 +111,7 @@ impl WrapUpRepository for PostgresWrapUpRepository {
.bind(&id_str)
.execute(&self.pool)
.await
.map_err(map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
Ok(())
}
@@ -132,7 +128,7 @@ impl WrapUpRepository for PostgresWrapUpRepository {
.bind(&id_str)
.fetch_optional(&self.pool)
.await
.map_err(map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
row.as_ref().map(row_to_record).transpose()
}
@@ -149,7 +145,7 @@ impl WrapUpRepository for PostgresWrapUpRepository {
.bind(&uid)
.fetch_all(&self.pool)
.await
.map_err(map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
rows.iter().map(row_to_record).collect()
}
@@ -163,7 +159,7 @@ impl WrapUpRepository for PostgresWrapUpRepository {
)
.fetch_all(&self.pool)
.await
.map_err(map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
rows.iter().map(row_to_record).collect()
}
@@ -190,7 +186,7 @@ impl WrapUpRepository for PostgresWrapUpRepository {
.bind(end)
.fetch_optional(&self.pool)
.await
.map_err(map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
row.as_ref().map(row_to_record).transpose()
}
@@ -200,7 +196,7 @@ impl WrapUpRepository for PostgresWrapUpRepository {
.bind(id.value().to_string())
.execute(&self.pool)
.await
.map_err(map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
Ok(())
}
@@ -213,21 +209,21 @@ impl WrapUpRepository for PostgresWrapUpRepository {
.bind(before)
.execute(&self.pool)
.await
.map_err(map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
Ok(result.rows_affected())
}
}
fn row_to_record(row: &sqlx::postgres::PgRow) -> Result<WrapUpRecord, DomainError> {
let id_str: String = row.try_get("id").map_err(map_err)?;
let user_id_str: Option<String> = row.try_get("user_id").map_err(map_err)?;
let start_date: NaiveDate = row.try_get("start_date").map_err(map_err)?;
let end_date: NaiveDate = row.try_get("end_date").map_err(map_err)?;
let status_str: String = row.try_get("status").map_err(map_err)?;
let report_json: Option<String> = row.try_get("report_json").map_err(map_err)?;
let error_message: Option<String> = row.try_get("error_message").map_err(map_err)?;
let created_at_str: String = row.try_get("created_at").map_err(map_err)?;
let completed_at_str: Option<String> = row.try_get("completed_at").map_err(map_err)?;
let id_str: String = row.try_get("id").map_err(adapter_common::map_sqlx_error)?;
let user_id_str: Option<String> = row.try_get("user_id").map_err(adapter_common::map_sqlx_error)?;
let start_date: NaiveDate = row.try_get("start_date").map_err(adapter_common::map_sqlx_error)?;
let end_date: NaiveDate = row.try_get("end_date").map_err(adapter_common::map_sqlx_error)?;
let status_str: String = row.try_get("status").map_err(adapter_common::map_sqlx_error)?;
let report_json: Option<String> = row.try_get("report_json").map_err(adapter_common::map_sqlx_error)?;
let error_message: Option<String> = row.try_get("error_message").map_err(adapter_common::map_sqlx_error)?;
let created_at_str: String = row.try_get("created_at").map_err(adapter_common::map_sqlx_error)?;
let completed_at_str: Option<String> = row.try_get("completed_at").map_err(adapter_common::map_sqlx_error)?;
let user_id = user_id_str.as_deref().map(parse_uuid).transpose()?;
@@ -292,7 +288,7 @@ impl WrapUpStatsQuery for PostgresWrapUpStatsQuery {
q = q.bind(uid);
}
let rows = q.fetch_all(&self.pool).await.map_err(map_err)?;
let rows = q.fetch_all(&self.pool).await.map_err(adapter_common::map_sqlx_error)?;
if rows.is_empty() {
return Ok(vec![]);
@@ -302,7 +298,7 @@ impl WrapUpStatsQuery for PostgresWrapUpStatsQuery {
let mut movie_ids: Vec<String> = Vec::new();
let mut seen = std::collections::HashSet::new();
for row in &rows {
let mid: String = row.try_get("movie_id").map_err(map_err)?;
let mid: String = row.try_get("movie_id").map_err(adapter_common::map_sqlx_error)?;
if seen.insert(mid.clone()) {
movie_ids.push(mid);
}
@@ -318,18 +314,18 @@ impl WrapUpStatsQuery for PostgresWrapUpStatsQuery {
// 3) Build result
let mut result = Vec::with_capacity(rows.len());
for row in &rows {
let movie_id_str: String = row.try_get("movie_id").map_err(map_err)?;
let title: String = row.try_get("title").map_err(map_err)?;
let release_year: i64 = row.try_get("release_year").map_err(map_err)?;
let director: Option<String> = row.try_get("director").map_err(map_err)?;
let poster_path: Option<String> = row.try_get("poster_path").map_err(map_err)?;
let rating: i64 = row.try_get("rating").map_err(map_err)?;
let watched_at_str: String = row.try_get("watched_at").map_err(map_err)?;
let user_id_str: String = row.try_get("user_id").map_err(map_err)?;
let runtime_minutes: Option<i32> = row.try_get("runtime_minutes").map_err(map_err)?;
let budget_usd: Option<i64> = row.try_get("budget_usd").map_err(map_err)?;
let movie_id_str: String = row.try_get("movie_id").map_err(adapter_common::map_sqlx_error)?;
let title: String = row.try_get("title").map_err(adapter_common::map_sqlx_error)?;
let release_year: i64 = row.try_get("release_year").map_err(adapter_common::map_sqlx_error)?;
let director: Option<String> = row.try_get("director").map_err(adapter_common::map_sqlx_error)?;
let poster_path: Option<String> = row.try_get("poster_path").map_err(adapter_common::map_sqlx_error)?;
let rating: i64 = row.try_get("rating").map_err(adapter_common::map_sqlx_error)?;
let watched_at_str: String = row.try_get("watched_at").map_err(adapter_common::map_sqlx_error)?;
let user_id_str: String = row.try_get("user_id").map_err(adapter_common::map_sqlx_error)?;
let runtime_minutes: Option<i32> = row.try_get("runtime_minutes").map_err(adapter_common::map_sqlx_error)?;
let budget_usd: Option<i64> = row.try_get("budget_usd").map_err(adapter_common::map_sqlx_error)?;
let original_language: Option<String> =
row.try_get("original_language").map_err(map_err)?;
row.try_get("original_language").map_err(adapter_common::map_sqlx_error)?;
let genres = genres_map.get(&movie_id_str).cloned().unwrap_or_default();
let keywords = keywords_map.get(&movie_id_str).cloned().unwrap_or_default();
@@ -383,12 +379,12 @@ async fn fetch_genres_pg(
.bind(movie_ids)
.fetch_all(pool)
.await
.map_err(map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
let mut map: HashMap<String, Vec<String>> = HashMap::new();
for row in rows {
let mid: String = row.try_get("movie_id").map_err(map_err)?;
let name: String = row.try_get("name").map_err(map_err)?;
let mid: String = row.try_get("movie_id").map_err(adapter_common::map_sqlx_error)?;
let name: String = row.try_get("name").map_err(adapter_common::map_sqlx_error)?;
map.entry(mid).or_default().push(name);
}
Ok(map)
@@ -404,12 +400,12 @@ async fn fetch_keywords_pg(
.bind(movie_ids)
.fetch_all(pool)
.await
.map_err(map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
let mut map: HashMap<String, Vec<String>> = HashMap::new();
for row in rows {
let mid: String = row.try_get("movie_id").map_err(map_err)?;
let name: String = row.try_get("name").map_err(map_err)?;
let mid: String = row.try_get("movie_id").map_err(adapter_common::map_sqlx_error)?;
let name: String = row.try_get("name").map_err(adapter_common::map_sqlx_error)?;
map.entry(mid).or_default().push(name);
}
Ok(map)
@@ -428,15 +424,15 @@ async fn fetch_cast_pg(
.bind(movie_ids)
.fetch_all(pool)
.await
.map_err(map_err)?;
.map_err(adapter_common::map_sqlx_error)?;
let mut map: HashMap<String, Vec<CastEntry>> = HashMap::new();
for row in rows {
let mid: String = row.try_get("movie_id").map_err(map_err)?;
let name: String = row.try_get("name").map_err(map_err)?;
let billing_order: i32 = row.try_get("billing_order").map_err(map_err)?;
let tmdb_person_id: i64 = row.try_get("tmdb_person_id").map_err(map_err)?;
let profile_path: Option<String> = row.try_get("profile_path").map_err(map_err)?;
let mid: String = row.try_get("movie_id").map_err(adapter_common::map_sqlx_error)?;
let name: String = row.try_get("name").map_err(adapter_common::map_sqlx_error)?;
let billing_order: i32 = row.try_get("billing_order").map_err(adapter_common::map_sqlx_error)?;
let tmdb_person_id: i64 = row.try_get("tmdb_person_id").map_err(adapter_common::map_sqlx_error)?;
let profile_path: Option<String> = row.try_get("profile_path").map_err(adapter_common::map_sqlx_error)?;
map.entry(mid).or_default().push(CastEntry {
name,
billing_order: billing_order as u32,