This commit is contained in:
2026-06-29 23:54:59 +02:00
parent bb17beb80d
commit 9c06fbd33e
11 changed files with 207 additions and 195 deletions

View File

@@ -4,7 +4,6 @@ version = "0.1.0"
edition = "2024" edition = "2024"
[dependencies] [dependencies]
application = { workspace = true }
domain = { workspace = true } domain = { workspace = true }
reqwest = { workspace = true } reqwest = { workspace = true }
serde = { workspace = true } serde = { workspace = true }

View File

@@ -1,16 +1,9 @@
use async_trait::async_trait; use domain::errors::DomainError;
use chrono::Utc;
use domain::{
errors::DomainError,
models::{CastMember, CrewMember, Genre, Keyword, MovieProfile, PersonEnrichmentData},
ports::{MovieEnrichmentClient, PersonEnrichmentClient},
value_objects::MovieId,
};
use serde::Deserialize; use serde::Deserialize;
pub struct TmdbEnrichmentClient { pub struct TmdbEnrichmentClient {
api_key: String, pub(crate) api_key: String,
http: reqwest::Client, pub(crate) http: reqwest::Client,
} }
impl TmdbEnrichmentClient { impl TmdbEnrichmentClient {
@@ -49,7 +42,7 @@ impl TmdbEnrichmentClient {
.map_err(|e| DomainError::InfrastructureError(e.to_string())) .map_err(|e| DomainError::InfrastructureError(e.to_string()))
} }
async fn resolve_tmdb_id(&self, external_id: &str) -> Result<u64, DomainError> { pub(crate) async fn resolve_tmdb_id(&self, external_id: &str) -> Result<u64, DomainError> {
if let Some(numeric) = external_id.strip_prefix("tmdb:") { if let Some(numeric) = external_id.strip_prefix("tmdb:") {
return numeric.parse::<u64>().map_err(|_| { return numeric.parse::<u64>().map_err(|_| {
DomainError::InfrastructureError(format!("Invalid tmdb id: {numeric}")) DomainError::InfrastructureError(format!("Invalid tmdb id: {numeric}"))
@@ -74,174 +67,3 @@ impl TmdbEnrichmentClient {
.ok_or_else(|| DomainError::NotFound(format!("TMDb: no movie for {external_id}"))) .ok_or_else(|| DomainError::NotFound(format!("TMDb: no movie for {external_id}")))
} }
} }
#[async_trait]
impl MovieEnrichmentClient for TmdbEnrichmentClient {
async fn fetch_profile(
&self,
movie_id: MovieId,
external_metadata_id: &str,
) -> Result<MovieProfile, DomainError> {
let tmdb_id = self.resolve_tmdb_id(external_metadata_id).await?;
#[derive(Deserialize)]
struct GenreDto {
id: u32,
name: String,
}
#[derive(Deserialize)]
struct CollectionDto {
name: String,
}
#[derive(Deserialize)]
struct CastDto {
id: u64,
name: String,
character: String,
order: u32,
profile_path: Option<String>,
}
#[derive(Deserialize)]
struct CrewDto {
id: u64,
name: String,
job: String,
department: String,
profile_path: Option<String>,
}
#[derive(Deserialize)]
struct Credits {
cast: Vec<CastDto>,
crew: Vec<CrewDto>,
}
#[derive(Deserialize)]
struct KeywordDto {
id: u32,
name: String,
}
#[derive(Deserialize)]
struct Keywords {
keywords: Vec<KeywordDto>,
}
#[derive(Deserialize)]
struct Details {
imdb_id: Option<String>,
overview: Option<String>,
tagline: Option<String>,
runtime: Option<u32>,
budget: Option<i64>,
revenue: Option<i64>,
vote_average: Option<f64>,
vote_count: Option<u32>,
original_language: Option<String>,
genres: Vec<GenreDto>,
belongs_to_collection: Option<CollectionDto>,
credits: Credits,
keywords: Keywords,
}
let url = self.base(&format!("/movie/{}", tmdb_id));
let d: Details = self
.get(&url, &[("append_to_response", "credits,keywords")])
.await?;
Ok(MovieProfile {
movie_id,
tmdb_id,
imdb_id: d.imdb_id.filter(|s| !s.is_empty()),
overview: d.overview.filter(|s| !s.is_empty()),
tagline: d.tagline.filter(|s| !s.is_empty()),
runtime_minutes: d.runtime,
budget_usd: d.budget.filter(|&v| v > 0),
revenue_usd: d.revenue.filter(|&v| v > 0),
vote_average: d.vote_average,
vote_count: d.vote_count,
original_language: d.original_language,
collection_name: d.belongs_to_collection.map(|c| c.name),
genres: d
.genres
.into_iter()
.map(|g| Genre {
tmdb_id: g.id,
name: g.name,
})
.collect(),
keywords: d
.keywords
.keywords
.into_iter()
.map(|k| Keyword {
tmdb_id: k.id,
name: k.name,
})
.collect(),
cast: d
.credits
.cast
.into_iter()
.map(|c| CastMember {
tmdb_person_id: c.id,
name: c.name,
character: c.character,
billing_order: c.order,
profile_path: c.profile_path,
})
.collect(),
crew: d
.credits
.crew
.into_iter()
.map(|c| CrewMember {
tmdb_person_id: c.id,
name: c.name,
job: c.job,
department: c.department,
profile_path: c.profile_path,
})
.collect(),
enriched_at: Utc::now(),
})
}
}
#[async_trait]
impl PersonEnrichmentClient for TmdbEnrichmentClient {
async fn fetch_details(&self, external_id: &str) -> Result<PersonEnrichmentData, DomainError> {
let tmdb_id = external_id
.strip_prefix("tmdb:")
.and_then(|s| s.parse::<u64>().ok())
.ok_or_else(|| {
DomainError::InfrastructureError(format!(
"Cannot parse person external_id: {external_id}"
))
})?;
#[derive(Deserialize)]
struct PersonDetails {
biography: Option<String>,
birthday: Option<String>,
deathday: Option<String>,
place_of_birth: Option<String>,
also_known_as: Option<Vec<String>>,
homepage: Option<String>,
imdb_id: Option<String>,
}
let url = self.base(&format!("/person/{tmdb_id}"));
let d: PersonDetails = self.get(&url, &[]).await?;
Ok(PersonEnrichmentData {
biography: d.biography.filter(|s| !s.is_empty()),
birthday: d
.birthday
.and_then(|s| chrono::NaiveDate::parse_from_str(&s, "%Y-%m-%d").ok()),
deathday: d
.deathday
.and_then(|s| chrono::NaiveDate::parse_from_str(&s, "%Y-%m-%d").ok()),
place_of_birth: d.place_of_birth.filter(|s| !s.is_empty()),
also_known_as: d.also_known_as.unwrap_or_default(),
homepage: d.homepage.filter(|s| !s.is_empty()),
imdb_id: d.imdb_id.filter(|s| !s.is_empty()),
})
}
}

View File

@@ -1,7 +1,5 @@
mod client; mod client;
mod movie_handler; mod movie;
mod person_handler; mod person;
pub use client::TmdbEnrichmentClient; pub use client::TmdbEnrichmentClient;
pub use movie_handler::MovieEnrichmentHandler;
pub use person_handler::PersonEnrichmentHandler;

View File

@@ -0,0 +1,140 @@
use async_trait::async_trait;
use chrono::Utc;
use domain::{
errors::DomainError,
models::{CastMember, CrewMember, Genre, Keyword, MovieProfile},
ports::MovieEnrichmentClient,
value_objects::MovieId,
};
use serde::Deserialize;
use crate::client::TmdbEnrichmentClient;
#[async_trait]
impl MovieEnrichmentClient for TmdbEnrichmentClient {
async fn fetch_profile(
&self,
movie_id: MovieId,
external_metadata_id: &str,
) -> Result<MovieProfile, DomainError> {
let tmdb_id = self.resolve_tmdb_id(external_metadata_id).await?;
#[derive(Deserialize)]
struct GenreDto {
id: u32,
name: String,
}
#[derive(Deserialize)]
struct CollectionDto {
name: String,
}
#[derive(Deserialize)]
struct CastDto {
id: u64,
name: String,
character: String,
order: u32,
profile_path: Option<String>,
}
#[derive(Deserialize)]
struct CrewDto {
id: u64,
name: String,
job: String,
department: String,
profile_path: Option<String>,
}
#[derive(Deserialize)]
struct Credits {
cast: Vec<CastDto>,
crew: Vec<CrewDto>,
}
#[derive(Deserialize)]
struct KeywordDto {
id: u32,
name: String,
}
#[derive(Deserialize)]
struct Keywords {
keywords: Vec<KeywordDto>,
}
#[derive(Deserialize)]
struct Details {
imdb_id: Option<String>,
overview: Option<String>,
tagline: Option<String>,
runtime: Option<u32>,
budget: Option<i64>,
revenue: Option<i64>,
vote_average: Option<f64>,
vote_count: Option<u32>,
original_language: Option<String>,
genres: Vec<GenreDto>,
belongs_to_collection: Option<CollectionDto>,
credits: Credits,
keywords: Keywords,
}
let url = self.base(&format!("/movie/{}", tmdb_id));
let d: Details = self
.get(&url, &[("append_to_response", "credits,keywords")])
.await?;
Ok(MovieProfile {
movie_id,
tmdb_id,
imdb_id: d.imdb_id.filter(|s| !s.is_empty()),
overview: d.overview.filter(|s| !s.is_empty()),
tagline: d.tagline.filter(|s| !s.is_empty()),
runtime_minutes: d.runtime,
budget_usd: d.budget.filter(|&v| v > 0),
revenue_usd: d.revenue.filter(|&v| v > 0),
vote_average: d.vote_average,
vote_count: d.vote_count,
original_language: d.original_language,
collection_name: d.belongs_to_collection.map(|c| c.name),
genres: d
.genres
.into_iter()
.map(|g| Genre {
tmdb_id: g.id,
name: g.name,
})
.collect(),
keywords: d
.keywords
.keywords
.into_iter()
.map(|k| Keyword {
tmdb_id: k.id,
name: k.name,
})
.collect(),
cast: d
.credits
.cast
.into_iter()
.map(|c| CastMember {
tmdb_person_id: c.id,
name: c.name,
character: c.character,
billing_order: c.order,
profile_path: c.profile_path,
})
.collect(),
crew: d
.credits
.crew
.into_iter()
.map(|c| CrewMember {
tmdb_person_id: c.id,
name: c.name,
job: c.job,
department: c.department,
profile_path: c.profile_path,
})
.collect(),
enriched_at: Utc::now(),
})
}
}

View File

@@ -0,0 +1,47 @@
use async_trait::async_trait;
use domain::{errors::DomainError, models::PersonEnrichmentData, ports::PersonEnrichmentClient};
use serde::Deserialize;
use crate::client::TmdbEnrichmentClient;
#[async_trait]
impl PersonEnrichmentClient for TmdbEnrichmentClient {
async fn fetch_details(&self, external_id: &str) -> Result<PersonEnrichmentData, DomainError> {
let tmdb_id = external_id
.strip_prefix("tmdb:")
.and_then(|s| s.parse::<u64>().ok())
.ok_or_else(|| {
DomainError::InfrastructureError(format!(
"Cannot parse person external_id: {external_id}"
))
})?;
#[derive(Deserialize)]
struct PersonDetails {
biography: Option<String>,
birthday: Option<String>,
deathday: Option<String>,
place_of_birth: Option<String>,
also_known_as: Option<Vec<String>>,
homepage: Option<String>,
imdb_id: Option<String>,
}
let url = self.base(&format!("/person/{tmdb_id}"));
let d: PersonDetails = self.get(&url, &[]).await?;
Ok(PersonEnrichmentData {
biography: d.biography.filter(|s| !s.is_empty()),
birthday: d
.birthday
.and_then(|s| chrono::NaiveDate::parse_from_str(&s, "%Y-%m-%d").ok()),
deathday: d
.deathday
.and_then(|s| chrono::NaiveDate::parse_from_str(&s, "%Y-%m-%d").ok()),
place_of_birth: d.place_of_birth.filter(|s| !s.is_empty()),
also_known_as: d.also_known_as.unwrap_or_default(),
homepage: d.homepage.filter(|s| !s.is_empty()),
imdb_id: d.imdb_id.filter(|s| !s.is_empty()),
})
}
}

View File

@@ -6,6 +6,7 @@ edition = "2024"
[dependencies] [dependencies]
async-trait = { workspace = true } async-trait = { workspace = true }
domain = { workspace = true } domain = { workspace = true }
reqwest = { workspace = true }
uuid = { workspace = true } uuid = { workspace = true }
chrono = { workspace = true } chrono = { workspace = true }
tracing = { workspace = true } tracing = { workspace = true }

View File

@@ -1,8 +1,5 @@
use std::sync::Arc; use std::sync::Arc;
use application::movies::{
commands::EnrichMovieCommand, deps::EnrichMovieDeps, enrich_movie, request_enrichment,
};
use async_trait::async_trait; use async_trait::async_trait;
use domain::{ use domain::{
errors::DomainError, errors::DomainError,
@@ -14,6 +11,10 @@ use domain::{
}, },
}; };
use crate::movies::{
commands::EnrichMovieCommand, deps::EnrichMovieDeps, enrich_movie, request_enrichment,
};
pub struct MovieEnrichmentHandler { pub struct MovieEnrichmentHandler {
enrichment_client: Arc<dyn MovieEnrichmentClient>, enrichment_client: Arc<dyn MovieEnrichmentClient>,
movie_repository: Arc<dyn MovieRepository>, movie_repository: Arc<dyn MovieRepository>,

View File

@@ -2,6 +2,7 @@ pub mod commands;
pub mod deps; pub mod deps;
pub mod discovery_indexer; pub mod discovery_indexer;
pub mod enrich_movie; pub mod enrich_movie;
pub mod event_handler;
pub mod get_movie_profile; pub mod get_movie_profile;
pub mod get_movies; pub mod get_movies;
pub mod queries; pub mod queries;
@@ -11,5 +12,6 @@ pub mod search_cleanup;
pub mod sync_poster; pub mod sync_poster;
pub use discovery_indexer::MovieDiscoveryIndexer; pub use discovery_indexer::MovieDiscoveryIndexer;
pub use event_handler::MovieEnrichmentHandler;
pub use reindex_search::SearchReindexHandler; pub use reindex_search::SearchReindexHandler;
pub use search_cleanup::SearchCleanupHandler; pub use search_cleanup::SearchCleanupHandler;

View File

@@ -7,7 +7,7 @@ use domain::{
ports::{EventHandler, PersonCommand, PersonEnrichmentClient, PersonQuery}, ports::{EventHandler, PersonCommand, PersonEnrichmentClient, PersonQuery},
}; };
use application::person::deps::EnrichPersonDeps; use super::deps::EnrichPersonDeps;
pub struct PersonEnrichmentHandler { pub struct PersonEnrichmentHandler {
deps: EnrichPersonDeps, deps: EnrichPersonDeps,
@@ -40,7 +40,6 @@ impl EventHandler for PersonEnrichmentHandler {
_ => return Ok(()), _ => return Ok(()),
}; };
application::person::enrich::execute(&self.deps, person_id, external_person_id.value()) super::enrich::execute(&self.deps, person_id, external_person_id.value()).await
.await
} }
} }

View File

@@ -1,4 +1,7 @@
pub mod deps; pub mod deps;
pub mod enrich; pub mod enrich;
pub mod event_handler;
pub mod get; pub mod get;
pub mod get_credits; pub mod get_credits;
pub use event_handler::PersonEnrichmentHandler;

View File

@@ -92,7 +92,7 @@ async fn main() -> anyhow::Result<()> {
Ok(client) => { Ok(client) => {
tracing::info!("TMDb enrichment enabled"); tracing::info!("TMDb enrichment enabled");
let client = Arc::new(client); let client = Arc::new(client);
let handler = Arc::new(tmdb_enrichment::MovieEnrichmentHandler::new( let handler = Arc::new(application::movies::MovieEnrichmentHandler::new(
Arc::clone(&client) as Arc<dyn MovieEnrichmentClient>, Arc::clone(&client) as Arc<dyn MovieEnrichmentClient>,
Arc::clone(&movie), Arc::clone(&movie),
Arc::clone(&movie_profile), Arc::clone(&movie_profile),
@@ -101,7 +101,7 @@ async fn main() -> anyhow::Result<()> {
Arc::clone(&object_storage), Arc::clone(&object_storage),
)) as Arc<dyn EventHandler>; )) as Arc<dyn EventHandler>;
let person_enrichment_arc = Arc::clone(&client) as Arc<dyn PersonEnrichmentClient>; let person_enrichment_arc = Arc::clone(&client) as Arc<dyn PersonEnrichmentClient>;
let person_handler = Arc::new(tmdb_enrichment::PersonEnrichmentHandler::new( let person_handler = Arc::new(application::person::PersonEnrichmentHandler::new(
Arc::clone(&person_query), Arc::clone(&person_query),
Some(person_enrichment_arc), Some(person_enrichment_arc),
Arc::clone(&person_command), Arc::clone(&person_command),