refactor: remaining MEDIUM — CQRS splits, DI Deps, profile dedup, event Value, response enum

M1: MovieRepository→MovieCommand/MovieQuery, WatchEventRepository→
WatchEventCommand/WatchEventQuery
M2: goals/ and import/ use Deps structs
M7: extract upload_image helper in update_profile
M8: FederationDeliveryRequested activity_json String→serde_json::Value
M11: UserProfileResponse uses ProfileViewData enum
This commit is contained in:
2026-07-10 03:50:43 +02:00
parent 12da356a40
commit dee013c7eb
99 changed files with 1262 additions and 896 deletions

View File

@@ -36,7 +36,7 @@ pub async fn execute(deps: &DeleteReviewDeps, cmd: DeleteReviewCommand) -> Resul
let history = deps.diary.get_review_history(&movie_id).await?;
if history.viewings().is_empty() {
let poster_path = history.movie().poster_path().cloned();
deps.movie.delete_movie(&movie_id).await?;
deps.movie_command.delete_movie(&movie_id).await?;
// best-effort: movie is already deleted, so publish failure is non-fatal
if let Err(e) = deps
.event_publisher

View File

@@ -1,8 +1,8 @@
use std::sync::Arc;
use domain::ports::{
DiaryRepository, EventPublisher, MovieProfileRepository, MovieRepository, ReviewRepository,
SocialQueryPort,
DiaryRepository, EventPublisher, MovieCommand, MovieProfileRepository, MovieQuery,
ReviewRepository, SocialQueryPort,
};
use crate::config::AppConfig;
@@ -10,7 +10,7 @@ use crate::config::AppConfig;
pub struct DeleteReviewDeps {
pub review: Arc<dyn ReviewRepository>,
pub diary: Arc<dyn DiaryRepository>,
pub movie: Arc<dyn MovieRepository>,
pub movie_command: Arc<dyn MovieCommand>,
pub event_publisher: Arc<dyn EventPublisher>,
}
@@ -20,7 +20,7 @@ pub struct EditReviewDeps {
}
pub struct GetMovieSocialPageDeps {
pub movie: Arc<dyn MovieRepository>,
pub movie_query: Arc<dyn MovieQuery>,
pub diary: Arc<dyn DiaryRepository>,
pub movie_profile: Arc<dyn MovieProfileRepository>,
}

View File

@@ -24,7 +24,7 @@ pub async fn execute(
let page = PageParams::new(Some(query.limit), Some(query.offset))?;
let movie = deps
.movie
.movie_query
.get_movie_by_id(&movie_id)
.await?
.ok_or_else(|| DomainError::NotFound(format!("Movie {}", query.movie_id)))?;

View File

@@ -2,14 +2,14 @@ use async_trait::async_trait;
use domain::{
errors::DomainError,
models::{MetadataSearchCriteria, Movie},
ports::{MetadataClient, MovieRepository},
ports::{MetadataClient, MovieQuery},
value_objects::{ExternalMetadataId, MovieTitle, ReleaseYear},
};
use crate::diary::commands::MovieInput;
pub struct MovieResolverDeps<'a> {
pub repository: &'a dyn MovieRepository,
pub repository: &'a dyn MovieQuery,
pub metadata_client: &'a dyn MetadataClient,
}

View File

@@ -6,7 +6,8 @@ use domain::{
events::DomainEvent,
models::Review,
ports::{
EventPublisher, MetadataClient, MovieRepository, ReviewRepository, WatchlistRepository,
EventPublisher, MetadataClient, MovieCommand, MovieQuery, ReviewRepository,
WatchlistRepository,
},
value_objects::{Comment, Rating, UserId},
};
@@ -16,7 +17,8 @@ use crate::movies::resolve::resolve_and_persist_movie;
use crate::ports::ReviewLogger;
pub struct DefaultReviewLogger {
movie_repo: Arc<dyn MovieRepository>,
movie_command: Arc<dyn MovieCommand>,
movie_query: Arc<dyn MovieQuery>,
review_repo: Arc<dyn ReviewRepository>,
watchlist_repo: Arc<dyn WatchlistRepository>,
metadata_client: Arc<dyn MetadataClient>,
@@ -25,14 +27,16 @@ pub struct DefaultReviewLogger {
impl DefaultReviewLogger {
pub fn new(
movie_repo: Arc<dyn MovieRepository>,
movie_command: Arc<dyn MovieCommand>,
movie_query: Arc<dyn MovieQuery>,
review_repo: Arc<dyn ReviewRepository>,
watchlist_repo: Arc<dyn WatchlistRepository>,
metadata_client: Arc<dyn MetadataClient>,
event_publisher: Arc<dyn EventPublisher>,
) -> Self {
Self {
movie_repo,
movie_command,
movie_query,
review_repo,
watchlist_repo,
metadata_client,
@@ -50,7 +54,8 @@ impl ReviewLogger for DefaultReviewLogger {
let (movie, is_new_movie) = resolve_and_persist_movie(
&cmd.input,
self.movie_repo.as_ref(),
self.movie_command.as_ref(),
self.movie_query.as_ref(),
self.metadata_client.as_ref(),
self.event_publisher.as_ref(),
)
@@ -58,7 +63,7 @@ impl ReviewLogger for DefaultReviewLogger {
// Always upsert: even existing movies may have updated metadata
if !is_new_movie {
self.movie_repo.upsert_movie(&movie).await?;
self.movie_command.upsert_movie(&movie).await?;
}
let review = Review::new(

View File

@@ -4,7 +4,7 @@ use chrono::Utc;
use domain::{
models::{Movie, Review},
ports::{MovieRepository, ReviewRepository},
ports::{MovieCommand, MovieQuery, ReviewRepository},
testing::{
FakeDiaryRepository, InMemoryMovieRepository, InMemoryReviewRepository, NoopEventPublisher,
},
@@ -55,7 +55,7 @@ async fn test_delete_review_removes_it() {
let deps = DeleteReviewDeps {
review: Arc::clone(&reviews) as _,
diary: diary.clone() as _,
movie: Arc::clone(&movies) as _,
movie_command: Arc::clone(&movies) as _,
event_publisher: Arc::clone(&events) as _,
};
@@ -93,7 +93,7 @@ async fn test_delete_review_wrong_user_is_unauthorized() {
let deps = DeleteReviewDeps {
review: Arc::clone(&reviews) as _,
diary: diary as _,
movie: movies as _,
movie_command: movies as _,
event_publisher: Arc::clone(&events) as _,
};

View File

@@ -4,7 +4,7 @@ use uuid::Uuid;
use domain::{
models::Movie,
ports::MovieRepository,
ports::MovieCommand,
testing::{FakeDiaryRepository, InMemoryMovieProfileRepository, InMemoryMovieRepository},
value_objects::{MovieTitle, ReleaseYear},
};
@@ -17,7 +17,7 @@ use crate::{
#[tokio::test]
async fn fails_when_movie_not_found() {
let deps = GetMovieSocialPageDeps {
movie: InMemoryMovieRepository::new(),
movie_query: InMemoryMovieRepository::new(),
diary: FakeDiaryRepository::new() as _,
movie_profile: InMemoryMovieProfileRepository::new(),
};
@@ -50,7 +50,7 @@ async fn returns_movie_social_page() {
movies.upsert_movie(&movie).await.unwrap();
let deps = GetMovieSocialPageDeps {
movie: Arc::clone(&movies) as _,
movie_query: Arc::clone(&movies) as _,
diary: FakeDiaryRepository::new() as _,
movie_profile: InMemoryMovieProfileRepository::new(),
};

View File

@@ -4,10 +4,10 @@ use chrono::Utc;
use domain::{
models::Movie,
ports::MovieCommand,
value_objects::{MovieTitle, ReleaseYear},
};
use domain::ports::MovieRepository;
use domain::testing::{InMemoryMovieRepository, InMemoryReviewRepository, NoopEventPublisher};
use crate::{
@@ -23,6 +23,7 @@ fn build_logger(
events: &Arc<NoopEventPublisher>,
) -> Arc<dyn crate::ports::ReviewLogger> {
Arc::new(DefaultReviewLogger::new(
Arc::clone(movies) as _,
Arc::clone(movies) as _,
Arc::clone(reviews) as _,
TestContextBuilder::new().watchlist_repo,

View File

@@ -3,7 +3,7 @@ use crate::diary::commands::MovieInput;
use domain::{
errors::DomainError,
models::{MetadataSearchCriteria, Movie},
ports::MovieRepository,
ports::MovieQuery,
value_objects::{ExternalMetadataId, MovieId, MovieTitle, PosterUrl, ReleaseYear},
};
@@ -32,7 +32,7 @@ struct RepoEmpty;
struct RepoWithTitleMatch(Movie);
#[async_trait::async_trait]
impl MovieRepository for RepoWithExternalMovie {
impl MovieQuery for RepoWithExternalMovie {
async fn get_movie_by_external_id(
&self,
_: &ExternalMetadataId,
@@ -49,12 +49,6 @@ impl MovieRepository for RepoWithExternalMovie {
) -> Result<Vec<Movie>, DomainError> {
panic!("unexpected")
}
async fn upsert_movie(&self, _: &Movie) -> Result<(), DomainError> {
panic!("unexpected")
}
async fn delete_movie(&self, _: &MovieId) -> Result<(), DomainError> {
panic!("unexpected")
}
async fn existing_external_ids(
&self,
_: &[ExternalMetadataId],
@@ -83,7 +77,7 @@ impl MovieRepository for RepoWithExternalMovie {
}
#[async_trait::async_trait]
impl MovieRepository for RepoEmpty {
impl MovieQuery for RepoEmpty {
async fn get_movie_by_external_id(
&self,
_: &ExternalMetadataId,
@@ -100,12 +94,6 @@ impl MovieRepository for RepoEmpty {
) -> Result<Vec<Movie>, DomainError> {
Ok(vec![])
}
async fn upsert_movie(&self, _: &Movie) -> Result<(), DomainError> {
panic!("unexpected")
}
async fn delete_movie(&self, _: &MovieId) -> Result<(), DomainError> {
panic!("unexpected")
}
async fn existing_external_ids(
&self,
_: &[ExternalMetadataId],
@@ -134,7 +122,7 @@ impl MovieRepository for RepoEmpty {
}
#[async_trait::async_trait]
impl MovieRepository for RepoWithTitleMatch {
impl MovieQuery for RepoWithTitleMatch {
async fn get_movie_by_external_id(
&self,
_: &ExternalMetadataId,
@@ -151,12 +139,6 @@ impl MovieRepository for RepoWithTitleMatch {
) -> Result<Vec<Movie>, DomainError> {
Ok(vec![self.0.clone()])
}
async fn upsert_movie(&self, _: &Movie) -> Result<(), DomainError> {
panic!("unexpected")
}
async fn delete_movie(&self, _: &MovieId) -> Result<(), DomainError> {
panic!("unexpected")
}
async fn existing_external_ids(
&self,
_: &[ExternalMetadataId],

View File

@@ -6,7 +6,7 @@ use domain::{
errors::DomainError,
models::WatchlistEntry,
models::{MetadataSearchCriteria, Movie},
ports::{MetadataClient, MovieRepository, WatchlistRepository},
ports::{MetadataClient, MovieCommand, WatchlistRepository},
testing::{
FakeMetadataClient, InMemoryMovieRepository, InMemoryReviewRepository,
InMemoryWatchlistRepository, NoopEventPublisher,
@@ -26,6 +26,7 @@ fn make_logger(
events: &Arc<NoopEventPublisher>,
) -> DefaultReviewLogger {
DefaultReviewLogger::new(
Arc::clone(movies) as _,
Arc::clone(movies) as _,
Arc::clone(reviews) as _,
Arc::clone(watchlist) as _,
@@ -276,6 +277,7 @@ async fn publishes_movie_discovered_for_new_movie_with_external_id() {
let events = NoopEventPublisher::new();
let logger = DefaultReviewLogger::new(
Arc::clone(&movies) as _,
Arc::clone(&movies) as _,
Arc::clone(&reviews) as _,
Arc::clone(&watchlist) as _,

View File

@@ -1,24 +1,19 @@
use std::sync::Arc;
use domain::{
errors::DomainError,
events::DomainEvent,
models::{Goal, GoalType, GoalWithProgress},
ports::{EventPublisher, GoalRepository, StatsRepository},
value_objects::UserId,
};
use super::commands::CreateGoalCommand;
use super::{commands::CreateGoalCommand, deps::GoalCommandDeps};
pub async fn execute(
goal: Arc<dyn GoalRepository>,
stats: Arc<dyn StatsRepository>,
event_publisher: Arc<dyn EventPublisher>,
deps: &GoalCommandDeps,
cmd: CreateGoalCommand,
) -> Result<GoalWithProgress, DomainError> {
let user_id = UserId::from_uuid(cmd.user_id);
let existing = goal.find_by_user_and_year(&user_id, cmd.year).await?;
let existing = deps.goal.find_by_user_and_year(&user_id, cmd.year).await?;
if existing.is_some() {
return Err(DomainError::ValidationError(
"Goal already exists for this year".into(),
@@ -31,11 +26,11 @@ pub async fn execute(
cmd.target_count,
GoalType::Movies,
)?;
goal.save(&g).await?;
deps.goal.save(&g).await?;
let current_count = stats.count_reviews_in_year(&user_id, cmd.year).await?;
let current_count = deps.stats.count_reviews_in_year(&user_id, cmd.year).await?;
event_publisher
deps.event_publisher
.publish(&DomainEvent::GoalCreated {
goal_id: g.id().clone(),
user_id,

View File

@@ -1,29 +1,26 @@
use std::sync::Arc;
use domain::{
errors::DomainError,
events::DomainEvent,
ports::{EventPublisher, GoalRepository},
value_objects::UserId,
};
use super::commands::DeleteGoalCommand;
use super::{commands::DeleteGoalCommand, deps::GoalCommandDeps};
pub async fn execute(
goal: Arc<dyn GoalRepository>,
event_publisher: Arc<dyn EventPublisher>,
deps: &GoalCommandDeps,
cmd: DeleteGoalCommand,
) -> Result<(), DomainError> {
let user_id = UserId::from_uuid(cmd.user_id);
let g = goal
let g = deps
.goal
.find_by_user_and_year(&user_id, cmd.year)
.await?
.ok_or_else(|| DomainError::NotFound(format!("Goal for year {}", cmd.year)))?;
goal.delete(g.id(), &user_id).await?;
deps.goal.delete(g.id(), &user_id).await?;
event_publisher
deps.event_publisher
.publish(&DomainEvent::GoalDeleted {
goal_id: g.id().clone(),
user_id,

View File

@@ -0,0 +1,14 @@
use std::sync::Arc;
use domain::ports::{EventPublisher, GoalRepository, StatsRepository};
pub struct GoalCommandDeps {
pub goal: Arc<dyn GoalRepository>,
pub stats: Arc<dyn StatsRepository>,
pub event_publisher: Arc<dyn EventPublisher>,
}
pub struct GoalQueryDeps {
pub goal: Arc<dyn GoalRepository>,
pub stats: Arc<dyn StatsRepository>,
}

View File

@@ -1,26 +1,22 @@
use std::sync::Arc;
use domain::{
errors::DomainError,
models::GoalWithProgress,
ports::{GoalRepository, StatsRepository},
value_objects::UserId,
};
use super::queries::GetGoalQuery;
use super::{deps::GoalQueryDeps, queries::GetGoalQuery};
pub async fn execute(
goal: Arc<dyn GoalRepository>,
stats: Arc<dyn StatsRepository>,
deps: &GoalQueryDeps,
query: GetGoalQuery,
) -> Result<Option<GoalWithProgress>, DomainError> {
let user_id = UserId::from_uuid(query.user_id);
let found = goal.find_by_user_and_year(&user_id, query.year).await?;
let found = deps.goal.find_by_user_and_year(&user_id, query.year).await?;
let Some(g) = found else { return Ok(None) };
let current_count = stats.count_reviews_in_year(&user_id, query.year).await?;
let current_count = deps.stats.count_reviews_in_year(&user_id, query.year).await?;
Ok(Some(GoalWithProgress {
goal: g,

View File

@@ -1,25 +1,21 @@
use std::sync::Arc;
use domain::{
errors::DomainError,
models::GoalWithProgress,
ports::{GoalRepository, StatsRepository},
value_objects::UserId,
};
use super::queries::ListGoalsQuery;
use super::{deps::GoalQueryDeps, queries::ListGoalsQuery};
pub async fn execute(
goal: Arc<dyn GoalRepository>,
stats: Arc<dyn StatsRepository>,
deps: &GoalQueryDeps,
query: ListGoalsQuery,
) -> Result<Vec<GoalWithProgress>, DomainError> {
let user_id = UserId::from_uuid(query.user_id);
let goals = goal.list_for_user(&user_id).await?;
let goals = deps.goal.list_for_user(&user_id).await?;
let mut result = Vec::with_capacity(goals.len());
for g in goals {
let current_count = stats.count_reviews_in_year(&user_id, g.year()).await?;
let current_count = deps.stats.count_reviews_in_year(&user_id, g.year()).await?;
result.push(GoalWithProgress {
goal: g,
current_count,

View File

@@ -1,6 +1,7 @@
pub mod commands;
pub mod create;
pub mod delete;
pub mod deps;
pub mod get;
pub mod list;
pub mod queries;

View File

@@ -4,6 +4,7 @@ use domain::events::DomainEvent;
use domain::testing::{FakeStatsRepository, InMemoryGoalRepository, NoopEventPublisher};
use uuid::Uuid;
use crate::goals::deps::GoalCommandDeps;
use crate::goals::{commands::CreateGoalCommand, create};
use crate::test_helpers::TestContextBuilder;
@@ -12,11 +13,14 @@ async fn creates_goal_and_returns_progress() {
let goals = InMemoryGoalRepository::new();
let stats = FakeStatsRepository::new();
let events = NoopEventPublisher::new();
let deps = GoalCommandDeps {
goal: Arc::clone(&goals) as _,
stats: Arc::clone(&stats) as _,
event_publisher: Arc::clone(&events) as _,
};
let result = create::execute(
Arc::clone(&goals) as _,
Arc::clone(&stats) as _,
Arc::clone(&events) as _,
&deps,
CreateGoalCommand {
user_id: Uuid::nil(),
year: 2025,
@@ -38,11 +42,14 @@ async fn creates_goal_with_review_count() {
let stats = FakeStatsRepository::new();
stats.set_review_count(Uuid::nil(), 2025, 5);
let events = NoopEventPublisher::new();
let deps = GoalCommandDeps {
goal: Arc::clone(&goals) as _,
stats: Arc::clone(&stats) as _,
event_publisher: Arc::clone(&events) as _,
};
let result = create::execute(
Arc::clone(&goals) as _,
Arc::clone(&stats) as _,
Arc::clone(&events) as _,
&deps,
CreateGoalCommand {
user_id: Uuid::nil(),
year: 2025,
@@ -60,11 +67,14 @@ async fn creates_goal_with_review_count() {
async fn emits_goal_created_event() {
let b = TestContextBuilder::new();
let events = NoopEventPublisher::new();
let deps = GoalCommandDeps {
goal: b.goal_repo.clone(),
stats: b.stats_repo.clone(),
event_publisher: Arc::clone(&events) as _,
};
create::execute(
b.goal_repo.clone(),
b.stats_repo.clone(),
Arc::clone(&events) as _,
&deps,
CreateGoalCommand {
user_id: Uuid::nil(),
year: 2025,
@@ -85,25 +95,21 @@ async fn emits_goal_created_event() {
#[tokio::test]
async fn rejects_duplicate_year() {
let b = TestContextBuilder::new();
let deps = GoalCommandDeps {
goal: b.goal_repo.clone(),
stats: b.stats_repo.clone(),
event_publisher: b.event_publisher.clone(),
};
let cmd = CreateGoalCommand {
user_id: Uuid::nil(),
year: 2025,
target_count: 10,
};
create::execute(
b.goal_repo.clone(),
b.stats_repo.clone(),
b.event_publisher.clone(),
cmd,
)
.await
.unwrap();
create::execute(&deps, cmd).await.unwrap();
let result = create::execute(
b.goal_repo.clone(),
b.stats_repo.clone(),
b.event_publisher.clone(),
&deps,
CreateGoalCommand {
user_id: Uuid::nil(),
year: 2025,
@@ -118,10 +124,13 @@ async fn rejects_duplicate_year() {
#[tokio::test]
async fn rejects_year_before_2020() {
let b = TestContextBuilder::new();
let deps = GoalCommandDeps {
goal: b.goal_repo.clone(),
stats: b.stats_repo.clone(),
event_publisher: b.event_publisher.clone(),
};
let result = create::execute(
b.goal_repo.clone(),
b.stats_repo.clone(),
b.event_publisher.clone(),
&deps,
CreateGoalCommand {
user_id: Uuid::nil(),
year: 2019,
@@ -136,10 +145,13 @@ async fn rejects_year_before_2020() {
#[tokio::test]
async fn rejects_zero_target() {
let b = TestContextBuilder::new();
let deps = GoalCommandDeps {
goal: b.goal_repo.clone(),
stats: b.stats_repo.clone(),
event_publisher: b.event_publisher.clone(),
};
let result = create::execute(
b.goal_repo.clone(),
b.stats_repo.clone(),
b.event_publisher.clone(),
&deps,
CreateGoalCommand {
user_id: Uuid::nil(),
year: 2025,

View File

@@ -3,6 +3,7 @@ use std::sync::Arc;
use domain::testing::{FakeStatsRepository, InMemoryGoalRepository, NoopEventPublisher};
use uuid::Uuid;
use crate::goals::deps::GoalCommandDeps;
use crate::goals::{
commands::{CreateGoalCommand, DeleteGoalCommand},
create, delete,
@@ -14,11 +15,14 @@ async fn deletes_existing_goal() {
let goals = InMemoryGoalRepository::new();
let stats = FakeStatsRepository::new();
let events = NoopEventPublisher::new();
let deps = GoalCommandDeps {
goal: Arc::clone(&goals) as _,
stats: Arc::clone(&stats) as _,
event_publisher: Arc::clone(&events) as _,
};
create::execute(
Arc::clone(&goals) as _,
Arc::clone(&stats) as _,
Arc::clone(&events) as _,
&deps,
CreateGoalCommand {
user_id: Uuid::nil(),
year: 2025,
@@ -30,8 +34,7 @@ async fn deletes_existing_goal() {
assert_eq!(goals.count(), 1);
delete::execute(
Arc::clone(&goals) as _,
Arc::clone(&events) as _,
&deps,
DeleteGoalCommand {
user_id: Uuid::nil(),
year: 2025,
@@ -46,9 +49,13 @@ async fn deletes_existing_goal() {
#[tokio::test]
async fn fails_when_not_found() {
let b = TestContextBuilder::new();
let deps = GoalCommandDeps {
goal: b.goal_repo.clone(),
stats: b.stats_repo.clone(),
event_publisher: b.event_publisher.clone(),
};
let result = delete::execute(
b.goal_repo.clone(),
b.event_publisher.clone(),
&deps,
DeleteGoalCommand {
user_id: Uuid::nil(),
year: 2025,

View File

@@ -1,15 +1,24 @@
use uuid::Uuid;
use crate::goals::deps::{GoalCommandDeps, GoalQueryDeps};
use crate::goals::{commands::CreateGoalCommand, create, get, queries::GetGoalQuery};
use crate::test_helpers::TestContextBuilder;
#[tokio::test]
async fn returns_goal_when_exists() {
let b = TestContextBuilder::new();
let cmd_deps = GoalCommandDeps {
goal: b.goal_repo.clone(),
stats: b.stats_repo.clone(),
event_publisher: b.event_publisher.clone(),
};
let query_deps = GoalQueryDeps {
goal: b.goal_repo.clone(),
stats: b.stats_repo.clone(),
};
create::execute(
b.goal_repo.clone(),
b.stats_repo.clone(),
b.event_publisher.clone(),
&cmd_deps,
CreateGoalCommand {
user_id: Uuid::nil(),
year: 2025,
@@ -20,8 +29,7 @@ async fn returns_goal_when_exists() {
.unwrap();
let result = get::execute(
b.goal_repo.clone(),
b.stats_repo.clone(),
&query_deps,
GetGoalQuery {
user_id: Uuid::nil(),
year: 2025,
@@ -37,9 +45,12 @@ async fn returns_goal_when_exists() {
#[tokio::test]
async fn returns_none_when_missing() {
let b = TestContextBuilder::new();
let query_deps = GoalQueryDeps {
goal: b.goal_repo.clone(),
stats: b.stats_repo.clone(),
};
let result = get::execute(
b.goal_repo.clone(),
b.stats_repo.clone(),
&query_deps,
GetGoalQuery {
user_id: Uuid::nil(),
year: 2025,

View File

@@ -1,14 +1,18 @@
use uuid::Uuid;
use crate::goals::deps::{GoalCommandDeps, GoalQueryDeps};
use crate::goals::{commands::CreateGoalCommand, create, list, queries::ListGoalsQuery};
use crate::test_helpers::TestContextBuilder;
#[tokio::test]
async fn returns_empty_when_no_goals() {
let b = TestContextBuilder::new();
let query_deps = GoalQueryDeps {
goal: b.goal_repo.clone(),
stats: b.stats_repo.clone(),
};
let result = list::execute(
b.goal_repo.clone(),
b.stats_repo.clone(),
&query_deps,
ListGoalsQuery {
user_id: Uuid::nil(),
},
@@ -22,11 +26,19 @@ async fn returns_empty_when_no_goals() {
#[tokio::test]
async fn returns_all_goals_for_user() {
let b = TestContextBuilder::new();
let cmd_deps = GoalCommandDeps {
goal: b.goal_repo.clone(),
stats: b.stats_repo.clone(),
event_publisher: b.event_publisher.clone(),
};
let query_deps = GoalQueryDeps {
goal: b.goal_repo.clone(),
stats: b.stats_repo.clone(),
};
for year in [2023, 2024, 2025] {
create::execute(
b.goal_repo.clone(),
b.stats_repo.clone(),
b.event_publisher.clone(),
&cmd_deps,
CreateGoalCommand {
user_id: Uuid::nil(),
year,
@@ -38,8 +50,7 @@ async fn returns_all_goals_for_user() {
}
let result = list::execute(
b.goal_repo.clone(),
b.stats_repo.clone(),
&query_deps,
ListGoalsQuery {
user_id: Uuid::nil(),
},

View File

@@ -1,5 +1,6 @@
use uuid::Uuid;
use crate::goals::deps::GoalCommandDeps;
use crate::goals::{
commands::{CreateGoalCommand, UpdateGoalCommand},
create, update,
@@ -9,10 +10,14 @@ use crate::test_helpers::TestContextBuilder;
#[tokio::test]
async fn updates_target_count() {
let b = TestContextBuilder::new();
let deps = GoalCommandDeps {
goal: b.goal_repo.clone(),
stats: b.stats_repo.clone(),
event_publisher: b.event_publisher.clone(),
};
create::execute(
b.goal_repo.clone(),
b.stats_repo.clone(),
b.event_publisher.clone(),
&deps,
CreateGoalCommand {
user_id: Uuid::nil(),
year: 2025,
@@ -23,9 +28,7 @@ async fn updates_target_count() {
.unwrap();
let result = update::execute(
b.goal_repo.clone(),
b.stats_repo.clone(),
b.event_publisher.clone(),
&deps,
UpdateGoalCommand {
user_id: Uuid::nil(),
year: 2025,
@@ -41,10 +44,13 @@ async fn updates_target_count() {
#[tokio::test]
async fn fails_when_goal_not_found() {
let b = TestContextBuilder::new();
let deps = GoalCommandDeps {
goal: b.goal_repo.clone(),
stats: b.stats_repo.clone(),
event_publisher: b.event_publisher.clone(),
};
let result = update::execute(
b.goal_repo.clone(),
b.stats_repo.clone(),
b.event_publisher.clone(),
&deps,
UpdateGoalCommand {
user_id: Uuid::nil(),
year: 2025,
@@ -59,10 +65,14 @@ async fn fails_when_goal_not_found() {
#[tokio::test]
async fn rejects_zero_target() {
let b = TestContextBuilder::new();
let deps = GoalCommandDeps {
goal: b.goal_repo.clone(),
stats: b.stats_repo.clone(),
event_publisher: b.event_publisher.clone(),
};
create::execute(
b.goal_repo.clone(),
b.stats_repo.clone(),
b.event_publisher.clone(),
&deps,
CreateGoalCommand {
user_id: Uuid::nil(),
year: 2025,
@@ -73,9 +83,7 @@ async fn rejects_zero_target() {
.unwrap();
let result = update::execute(
b.goal_repo.clone(),
b.stats_repo.clone(),
b.event_publisher.clone(),
&deps,
UpdateGoalCommand {
user_id: Uuid::nil(),
year: 2025,

View File

@@ -1,34 +1,30 @@
use std::sync::Arc;
use domain::{
errors::DomainError,
events::DomainEvent,
models::GoalWithProgress,
ports::{EventPublisher, GoalRepository, StatsRepository},
value_objects::UserId,
};
use super::commands::UpdateGoalCommand;
use super::{commands::UpdateGoalCommand, deps::GoalCommandDeps};
pub async fn execute(
goal: Arc<dyn GoalRepository>,
stats: Arc<dyn StatsRepository>,
event_publisher: Arc<dyn EventPublisher>,
deps: &GoalCommandDeps,
cmd: UpdateGoalCommand,
) -> Result<GoalWithProgress, DomainError> {
let user_id = UserId::from_uuid(cmd.user_id);
let mut g = goal
let mut g = deps
.goal
.find_by_user_and_year(&user_id, cmd.year)
.await?
.ok_or_else(|| DomainError::NotFound(format!("Goal for year {}", cmd.year)))?;
g.update_target(cmd.target_count)?;
goal.update(&g).await?;
deps.goal.update(&g).await?;
let current_count = stats.count_reviews_in_year(&user_id, cmd.year).await?;
let current_count = deps.stats.count_reviews_in_year(&user_id, cmd.year).await?;
event_publisher
deps.event_publisher
.publish(&DomainEvent::GoalUpdated {
goal_id: g.id().clone(),
user_id,

View File

@@ -3,22 +3,21 @@ use std::sync::Arc;
use domain::{
errors::DomainError,
models::{AnnotatedRow, import::RowResult},
ports::{DocumentParser, ImportSessionRepository, MovieRepository},
ports::MovieQuery,
value_objects::{ExternalMetadataId, ImportSessionId, MovieTitle, ReleaseYear, UserId},
};
use crate::import::commands::ApplyImportMappingCommand;
use super::{commands::ApplyImportMappingCommand, deps::ApplyMappingDeps};
pub async fn execute(
import_session: Arc<dyn ImportSessionRepository>,
document_parser: Arc<dyn DocumentParser>,
movie: Arc<dyn MovieRepository>,
deps: &ApplyMappingDeps,
cmd: ApplyImportMappingCommand,
) -> Result<Vec<AnnotatedRow>, DomainError> {
let user_id = UserId::from_uuid(cmd.user_id);
let session_id = ImportSessionId::from_uuid(cmd.session_id);
let mappings = cmd.mappings;
let mut session = import_session
let mut session = deps
.import_session
.get(&session_id, &user_id)
.await?
.ok_or_else(|| DomainError::NotFound("import session".into()))?;
@@ -28,20 +27,20 @@ pub async fn execute(
.clone()
.ok_or_else(|| DomainError::ValidationError("session has no parsed file".into()))?;
let mut annotated = document_parser.apply_mapping(&parsed, &mappings);
let mut annotated = deps.document_parser.apply_mapping(&parsed, &mappings);
mark_duplicates(movie, &mut annotated).await?;
mark_duplicates(deps.movie_query.clone(), &mut annotated).await?;
session.field_mappings = Some(mappings);
session.row_results = Some(annotated.clone());
import_session.update(&session).await?;
deps.import_session.update(&session).await?;
Ok(annotated)
}
async fn mark_duplicates(
movie: Arc<dyn MovieRepository>,
movie: Arc<dyn MovieQuery>,
rows: &mut [AnnotatedRow],
) -> Result<(), DomainError> {
let mut ext_ids = Vec::new();

View File

@@ -1,34 +1,33 @@
use std::sync::Arc;
use crate::import::commands::ApplyImportProfileCommand;
use domain::{
errors::DomainError,
ports::{ImportProfileRepository, ImportSessionRepository},
value_objects::{ImportProfileId, ImportSessionId, UserId},
};
use super::{commands::ApplyImportProfileCommand, deps::ApplyProfileDeps};
/// Copies the profile's field_mappings onto the session. Caller must then invoke
/// apply_import_mapping to regenerate row_results with the new mappings.
pub async fn execute(
import_profile: Arc<dyn ImportProfileRepository>,
import_session: Arc<dyn ImportSessionRepository>,
deps: &ApplyProfileDeps,
cmd: ApplyImportProfileCommand,
) -> Result<(), DomainError> {
let user_id = UserId::from_uuid(cmd.user_id);
let session_id = ImportSessionId::from_uuid(cmd.session_id);
let profile_id = ImportProfileId::from_uuid(cmd.profile_id);
let profile = import_profile
let profile = deps
.import_profile
.get(&profile_id, &user_id)
.await?
.ok_or_else(|| DomainError::NotFound("import profile".into()))?;
let mut session = import_session
let mut session = deps
.import_session
.get(&session_id, &user_id)
.await?
.ok_or_else(|| DomainError::NotFound("import session".into()))?;
session.field_mappings = Some(profile.field_mappings);
session.row_results = None;
import_session.update(&session).await
deps.import_session.update(&session).await
}
#[cfg(test)]

View File

@@ -1,13 +1,10 @@
use std::sync::Arc;
use domain::{
errors::DomainError,
models::ImportSession,
ports::{DocumentParser, ImportSessionRepository},
value_objects::{ImportSessionId, UserId},
};
use crate::import::commands::CreateImportSessionCommand;
use super::{commands::CreateImportSessionCommand, deps::CreateSessionDeps};
pub struct CreateSessionResult {
pub session_id: ImportSessionId,
@@ -16,14 +13,14 @@ pub struct CreateSessionResult {
}
pub async fn execute(
import_session: Arc<dyn ImportSessionRepository>,
document_parser: Arc<dyn DocumentParser>,
deps: &CreateSessionDeps,
cmd: CreateImportSessionCommand,
) -> Result<CreateSessionResult, DomainError> {
let user_id = UserId::from_uuid(cmd.user_id);
import_session.delete_expired_for_user(&user_id).await?;
deps.import_session.delete_expired_for_user(&user_id).await?;
let parsed = document_parser
let parsed = deps
.document_parser
.parse(&cmd.bytes, cmd.format)
.map_err(|e| DomainError::ValidationError(e.to_string()))?;
@@ -34,7 +31,7 @@ pub async fn execute(
let session_id = session.id.clone();
session.parsed_file = Some(parsed);
import_session.create(&session).await?;
deps.import_session.create(&session).await?;
Ok(CreateSessionResult {
session_id,

View File

@@ -0,0 +1,31 @@
use std::sync::Arc;
use domain::ports::{DocumentParser, ImportProfileRepository, ImportSessionRepository, MovieQuery};
use crate::ports::ReviewLogger;
pub struct CreateSessionDeps {
pub import_session: Arc<dyn ImportSessionRepository>,
pub document_parser: Arc<dyn DocumentParser>,
}
pub struct ApplyMappingDeps {
pub import_session: Arc<dyn ImportSessionRepository>,
pub document_parser: Arc<dyn DocumentParser>,
pub movie_query: Arc<dyn MovieQuery>,
}
pub struct ApplyProfileDeps {
pub import_profile: Arc<dyn ImportProfileRepository>,
pub import_session: Arc<dyn ImportSessionRepository>,
}
pub struct ExecuteImportDeps {
pub import_session: Arc<dyn ImportSessionRepository>,
pub review_logger: Arc<dyn ReviewLogger>,
}
pub struct SaveProfileDeps {
pub import_session: Arc<dyn ImportSessionRepository>,
pub import_profile: Arc<dyn ImportProfileRepository>,
}

View File

@@ -4,7 +4,6 @@ use chrono::NaiveDateTime;
use domain::{
errors::DomainError,
models::{ImportRow, import::RowResult},
ports::ImportSessionRepository,
value_objects::{ImportSessionId, UserId},
};
use uuid::Uuid;
@@ -12,9 +11,10 @@ use uuid::Uuid;
use crate::{
diary::commands::{LogReviewCommand, MovieInput},
import::commands::ExecuteImportCommand,
ports::ReviewLogger,
};
use super::deps::ExecuteImportDeps;
const CONCURRENCY_LIMIT: usize = 10;
pub struct ImportSummary {
@@ -24,14 +24,14 @@ pub struct ImportSummary {
}
pub async fn execute(
import_session: Arc<dyn ImportSessionRepository>,
review_logger: Arc<dyn ReviewLogger>,
deps: &ExecuteImportDeps,
cmd: ExecuteImportCommand,
) -> Result<ImportSummary, DomainError> {
let user_id = UserId::from_uuid(cmd.user_id);
let session_id = ImportSessionId::from_uuid(cmd.session_id);
let confirmed_indices = cmd.confirmed_indices;
let session = import_session
let session = deps
.import_session
.get(&session_id, &user_id)
.await?
.ok_or_else(|| DomainError::NotFound("import session".into()))?;
@@ -59,7 +59,7 @@ pub async fn execute(
Err(e) => failed.push((idx, e)),
Ok(log_cmd) => {
let permit = Arc::clone(&semaphore).acquire_owned().await.unwrap();
let logger = Arc::clone(&review_logger);
let logger = deps.review_logger.clone();
tasks.spawn(async move {
let result = logger.log_review(log_cmd).await.map_err(|e| e.to_string());
drop(permit);
@@ -78,7 +78,7 @@ pub async fn execute(
}
}
import_session.delete(&session_id).await?;
deps.import_session.delete(&session_id).await?;
Ok(ImportSummary {
imported,

View File

@@ -4,6 +4,7 @@ pub mod cleanup;
pub mod commands;
pub mod create_session;
pub mod delete_profile;
pub mod deps;
pub mod execute;
pub mod list_profiles;
pub mod save_profile;

View File

@@ -1,23 +1,21 @@
use std::sync::Arc;
use crate::import::commands::SaveImportProfileCommand;
use chrono::Utc;
use domain::{
errors::DomainError,
models::ImportProfile,
ports::{ImportProfileRepository, ImportSessionRepository},
value_objects::{ImportProfileId, ImportSessionId, UserId},
};
use super::{commands::SaveImportProfileCommand, deps::SaveProfileDeps};
pub async fn execute(
import_session: Arc<dyn ImportSessionRepository>,
import_profile: Arc<dyn ImportProfileRepository>,
deps: &SaveProfileDeps,
cmd: SaveImportProfileCommand,
) -> Result<ImportProfileId, DomainError> {
let user_id = UserId::from_uuid(cmd.user_id);
let session_id = ImportSessionId::from_uuid(cmd.session_id);
let session = import_session
let session = deps
.import_session
.get(&session_id, &user_id)
.await?
.ok_or_else(|| DomainError::NotFound("import session".into()))?;
@@ -32,7 +30,7 @@ pub async fn execute(
Utc::now().naive_utc(),
);
let id = profile.id.clone();
import_profile.save(&profile).await?;
deps.import_profile.save(&profile).await?;
Ok(id)
}

View File

@@ -7,11 +7,12 @@ use domain::{
AnnotatedRow, Movie,
import::{ImportRow, ParsedFile, RowResult},
},
ports::{DocumentParser, MovieRepository},
ports::{DocumentParser, MovieCommand},
testing::{InMemoryImportSessionRepository, InMemoryMovieRepository},
value_objects::{ExternalMetadataId, MovieTitle, ReleaseYear},
};
use crate::import::deps::{ApplyMappingDeps, CreateSessionDeps};
use crate::import::{
apply_mapping,
commands::{ApplyImportMappingCommand, CreateImportSessionCommand},
@@ -25,9 +26,13 @@ async fn applies_mapping_to_session() {
let b = TestContextBuilder::new();
let user_id = Uuid::new_v4();
let create_deps = CreateSessionDeps {
import_session: Arc::clone(&sessions) as _,
document_parser: b.document_parser.clone(),
};
let session = create_session::execute(
Arc::clone(&sessions) as _,
b.document_parser.clone(),
&create_deps,
CreateImportSessionCommand {
user_id,
bytes: b"title\nTest".to_vec(),
@@ -37,10 +42,14 @@ async fn applies_mapping_to_session() {
.await
.unwrap();
let mapping_deps = ApplyMappingDeps {
import_session: Arc::clone(&sessions) as _,
document_parser: b.document_parser.clone(),
movie_query: b.movie_query.clone(),
};
let rows = apply_mapping::execute(
Arc::clone(&sessions) as _,
b.document_parser.clone(),
b.movie_repo.clone(),
&mapping_deps,
ApplyImportMappingCommand {
user_id,
session_id: session.session_id.value(),
@@ -58,10 +67,14 @@ async fn fails_when_session_not_found() {
let sessions = InMemoryImportSessionRepository::new();
let b = TestContextBuilder::new();
let deps = ApplyMappingDeps {
import_session: Arc::clone(&sessions) as _,
document_parser: b.document_parser.clone(),
movie_query: b.movie_query.clone(),
};
let result = apply_mapping::execute(
Arc::clone(&sessions) as _,
b.document_parser.clone(),
b.movie_repo.clone(),
&deps,
ApplyImportMappingCommand {
user_id: Uuid::new_v4(),
session_id: Uuid::new_v4(),
@@ -132,9 +145,13 @@ async fn marks_duplicate_by_external_id() {
let user_id = Uuid::new_v4();
let create_deps = CreateSessionDeps {
import_session: Arc::clone(&sessions) as _,
document_parser: Arc::clone(&parser) as _,
};
let session = create_session::execute(
Arc::clone(&sessions) as _,
Arc::clone(&parser) as _,
&create_deps,
CreateImportSessionCommand {
user_id,
bytes: b"title\nKnown Movie".to_vec(),
@@ -144,10 +161,14 @@ async fn marks_duplicate_by_external_id() {
.await
.unwrap();
let mapping_deps = ApplyMappingDeps {
import_session: Arc::clone(&sessions) as _,
document_parser: Arc::clone(&parser) as _,
movie_query: Arc::clone(&movies) as _,
};
let rows = apply_mapping::execute(
Arc::clone(&sessions) as _,
Arc::clone(&parser) as _,
Arc::clone(&movies) as _,
&mapping_deps,
ApplyImportMappingCommand {
user_id,
session_id: session.session_id.value(),
@@ -185,9 +206,13 @@ async fn marks_duplicate_by_title_and_year() {
let user_id = Uuid::new_v4();
let create_deps = CreateSessionDeps {
import_session: Arc::clone(&sessions) as _,
document_parser: Arc::clone(&parser) as _,
};
let session = create_session::execute(
Arc::clone(&sessions) as _,
Arc::clone(&parser) as _,
&create_deps,
CreateImportSessionCommand {
user_id,
bytes: b"title\nDuplicate Film".to_vec(),
@@ -197,10 +222,14 @@ async fn marks_duplicate_by_title_and_year() {
.await
.unwrap();
let mapping_deps = ApplyMappingDeps {
import_session: Arc::clone(&sessions) as _,
document_parser: Arc::clone(&parser) as _,
movie_query: Arc::clone(&movies) as _,
};
let rows = apply_mapping::execute(
Arc::clone(&sessions) as _,
Arc::clone(&parser) as _,
Arc::clone(&movies) as _,
&mapping_deps,
ApplyImportMappingCommand {
user_id,
session_id: session.session_id.value(),

View File

@@ -7,6 +7,7 @@ use domain::testing::{InMemoryImportProfileRepository, InMemoryImportSessionRepo
use domain::value_objects::{ImportProfileId, UserId};
use uuid::Uuid;
use crate::import::deps::ApplyProfileDeps;
use crate::import::{apply_profile, commands::ApplyImportProfileCommand};
#[tokio::test]
@@ -14,9 +15,13 @@ async fn fails_when_profile_not_found() {
let profiles = InMemoryImportProfileRepository::new();
let sessions = InMemoryImportSessionRepository::new();
let deps = ApplyProfileDeps {
import_profile: Arc::clone(&profiles) as _,
import_session: Arc::clone(&sessions) as _,
};
let result = apply_profile::execute(
Arc::clone(&profiles) as _,
Arc::clone(&sessions) as _,
&deps,
ApplyImportProfileCommand {
user_id: Uuid::new_v4(),
session_id: Uuid::new_v4(),
@@ -44,9 +49,13 @@ async fn fails_when_session_not_found() {
let profile_id = profile.id.clone();
profiles.save(&profile).await.unwrap();
let deps = ApplyProfileDeps {
import_profile: Arc::clone(&profiles) as _,
import_session: Arc::clone(&sessions) as _,
};
let result = apply_profile::execute(
Arc::clone(&profiles) as _,
Arc::clone(&sessions) as _,
&deps,
ApplyImportProfileCommand {
user_id,
session_id: Uuid::new_v4(),
@@ -82,9 +91,13 @@ async fn applies_profile_mappings_to_session() {
let session_id = session.id.clone();
sessions.create(&session).await.unwrap();
let deps = ApplyProfileDeps {
import_profile: Arc::clone(&profiles) as _,
import_session: Arc::clone(&sessions) as _,
};
apply_profile::execute(
Arc::clone(&profiles) as _,
Arc::clone(&sessions) as _,
&deps,
ApplyImportProfileCommand {
user_id,
session_id: session_id.value(),

View File

@@ -4,6 +4,7 @@ use uuid::Uuid;
use domain::testing::InMemoryImportSessionRepository;
use crate::import::deps::CreateSessionDeps;
use crate::import::{commands::CreateImportSessionCommand, create_session};
use crate::test_helpers::TestContextBuilder;
@@ -12,9 +13,13 @@ async fn creates_session_with_parsed_file() {
let sessions = InMemoryImportSessionRepository::new();
let b = TestContextBuilder::new();
let deps = CreateSessionDeps {
import_session: Arc::clone(&sessions) as _,
document_parser: b.document_parser.clone(),
};
let result = create_session::execute(
Arc::clone(&sessions) as _,
b.document_parser.clone(),
&deps,
CreateImportSessionCommand {
user_id: Uuid::new_v4(),
bytes: b"col1\nval1".to_vec(),

View File

@@ -7,6 +7,7 @@ use domain::value_objects::UserId;
use uuid::Uuid;
use crate::import::commands::ExecuteImportCommand;
use crate::import::deps::ExecuteImportDeps;
use crate::import::execute;
use crate::test_helpers::NoopReviewLogger;
@@ -50,9 +51,13 @@ async fn imports_confirmed_rows() {
let sid = session.id.clone();
sessions.create(&session).await.unwrap();
let deps = ExecuteImportDeps {
import_session: Arc::clone(&sessions) as _,
review_logger: Arc::new(NoopReviewLogger),
};
let result = execute::execute(
Arc::clone(&sessions) as _,
Arc::new(NoopReviewLogger),
&deps,
ExecuteImportCommand {
user_id: uid,
session_id: sid.value(),
@@ -76,9 +81,13 @@ async fn skips_unconfirmed_rows() {
let sid = session.id.clone();
sessions.create(&session).await.unwrap();
let deps = ExecuteImportDeps {
import_session: Arc::clone(&sessions) as _,
review_logger: Arc::new(NoopReviewLogger),
};
let result = execute::execute(
Arc::clone(&sessions) as _,
Arc::new(NoopReviewLogger),
&deps,
ExecuteImportCommand {
user_id: uid,
session_id: sid.value(),
@@ -96,9 +105,13 @@ async fn skips_unconfirmed_rows() {
async fn fails_when_session_not_found() {
let sessions = InMemoryImportSessionRepository::new();
let deps = ExecuteImportDeps {
import_session: Arc::clone(&sessions) as _,
review_logger: Arc::new(NoopReviewLogger),
};
let result = execute::execute(
Arc::clone(&sessions) as _,
Arc::new(NoopReviewLogger),
&deps,
ExecuteImportCommand {
user_id: Uuid::new_v4(),
session_id: Uuid::new_v4(),
@@ -131,9 +144,13 @@ async fn handles_datetime_format() {
}]);
sessions.create(&session).await.unwrap();
let deps = ExecuteImportDeps {
import_session: Arc::clone(&sessions) as _,
review_logger: Arc::new(NoopReviewLogger),
};
let result = execute::execute(
Arc::clone(&sessions) as _,
Arc::new(NoopReviewLogger),
&deps,
ExecuteImportCommand {
user_id: uid,
session_id: sid.value(),
@@ -168,9 +185,13 @@ async fn fails_on_invalid_rating() {
}]);
sessions.create(&session).await.unwrap();
let deps = ExecuteImportDeps {
import_session: Arc::clone(&sessions) as _,
review_logger: Arc::new(NoopReviewLogger),
};
let result = execute::execute(
Arc::clone(&sessions) as _,
Arc::new(NoopReviewLogger),
&deps,
ExecuteImportCommand {
user_id: uid,
session_id: sid.value(),
@@ -205,9 +226,13 @@ async fn fails_on_missing_watched_at() {
}]);
sessions.create(&session).await.unwrap();
let deps = ExecuteImportDeps {
import_session: Arc::clone(&sessions) as _,
review_logger: Arc::new(NoopReviewLogger),
};
let result = execute::execute(
Arc::clone(&sessions) as _,
Arc::new(NoopReviewLogger),
&deps,
ExecuteImportCommand {
user_id: uid,
session_id: sid.value(),
@@ -242,9 +267,13 @@ async fn imports_row_with_external_metadata_id() {
}]);
sessions.create(&session).await.unwrap();
let deps = ExecuteImportDeps {
import_session: Arc::clone(&sessions) as _,
review_logger: Arc::new(NoopReviewLogger),
};
let result = execute::execute(
Arc::clone(&sessions) as _,
Arc::new(NoopReviewLogger),
&deps,
ExecuteImportCommand {
user_id: uid,
session_id: sid.value(),
@@ -279,9 +308,13 @@ async fn imports_row_with_director_and_comment() {
}]);
sessions.create(&session).await.unwrap();
let deps = ExecuteImportDeps {
import_session: Arc::clone(&sessions) as _,
review_logger: Arc::new(NoopReviewLogger),
};
let result = execute::execute(
Arc::clone(&sessions) as _,
Arc::new(NoopReviewLogger),
&deps,
ExecuteImportCommand {
user_id: uid,
session_id: sid.value(),
@@ -316,9 +349,13 @@ async fn handles_space_separated_datetime_format() {
}]);
sessions.create(&session).await.unwrap();
let deps = ExecuteImportDeps {
import_session: Arc::clone(&sessions) as _,
review_logger: Arc::new(NoopReviewLogger),
};
let result = execute::execute(
Arc::clone(&sessions) as _,
Arc::new(NoopReviewLogger),
&deps,
ExecuteImportCommand {
user_id: uid,
session_id: sid.value(),
@@ -348,9 +385,13 @@ async fn reports_invalid_row_result_errors() {
}]);
sessions.create(&session).await.unwrap();
let deps = ExecuteImportDeps {
import_session: Arc::clone(&sessions) as _,
review_logger: Arc::new(NoopReviewLogger),
};
let result = execute::execute(
Arc::clone(&sessions) as _,
Arc::new(NoopReviewLogger),
&deps,
ExecuteImportCommand {
user_id: uid,
session_id: sid.value(),
@@ -387,9 +428,13 @@ async fn fails_on_missing_rating() {
}]);
sessions.create(&session).await.unwrap();
let deps = ExecuteImportDeps {
import_session: Arc::clone(&sessions) as _,
review_logger: Arc::new(NoopReviewLogger),
};
let result = execute::execute(
Arc::clone(&sessions) as _,
Arc::new(NoopReviewLogger),
&deps,
ExecuteImportCommand {
user_id: uid,
session_id: sid.value(),
@@ -425,9 +470,13 @@ async fn fails_on_unparseable_date() {
}]);
sessions.create(&session).await.unwrap();
let deps = ExecuteImportDeps {
import_session: Arc::clone(&sessions) as _,
review_logger: Arc::new(NoopReviewLogger),
};
let result = execute::execute(
Arc::clone(&sessions) as _,
Arc::new(NoopReviewLogger),
&deps,
ExecuteImportCommand {
user_id: uid,
session_id: sid.value(),
@@ -463,9 +512,13 @@ async fn imports_row_without_release_year() {
}]);
sessions.create(&session).await.unwrap();
let deps = ExecuteImportDeps {
import_session: Arc::clone(&sessions) as _,
review_logger: Arc::new(NoopReviewLogger),
};
let result = execute::execute(
Arc::clone(&sessions) as _,
Arc::new(NoopReviewLogger),
&deps,
ExecuteImportCommand {
user_id: uid,
session_id: sid.value(),
@@ -489,9 +542,13 @@ async fn deletes_session_after_import() {
sessions.create(&session).await.unwrap();
assert_eq!(sessions.count(), 1);
let deps = ExecuteImportDeps {
import_session: Arc::clone(&sessions) as _,
review_logger: Arc::new(NoopReviewLogger),
};
execute::execute(
Arc::clone(&sessions) as _,
Arc::new(NoopReviewLogger),
&deps,
ExecuteImportCommand {
user_id: uid,
session_id: sid.value(),
@@ -533,10 +590,14 @@ async fn imports_more_rows_than_concurrency_limit() {
session.row_results = Some(rows);
sessions.create(&session).await.unwrap();
let deps = ExecuteImportDeps {
import_session: Arc::clone(&sessions) as _,
review_logger: Arc::new(NoopReviewLogger),
};
let confirmed_indices: Vec<usize> = (0..15).collect();
let result = execute::execute(
Arc::clone(&sessions) as _,
Arc::new(NoopReviewLogger),
&deps,
ExecuteImportCommand {
user_id: uid,
session_id: sid.value(),

View File

@@ -6,6 +6,7 @@ use domain::testing::{InMemoryImportProfileRepository, InMemoryImportSessionRepo
use domain::value_objects::UserId;
use uuid::Uuid;
use crate::import::deps::SaveProfileDeps;
use crate::import::{commands::SaveImportProfileCommand, save_profile};
#[tokio::test]
@@ -13,9 +14,13 @@ async fn fails_when_session_not_found() {
let sessions = InMemoryImportSessionRepository::new();
let profiles = InMemoryImportProfileRepository::new();
let deps = SaveProfileDeps {
import_session: Arc::clone(&sessions) as _,
import_profile: Arc::clone(&profiles) as _,
};
let result = save_profile::execute(
Arc::clone(&sessions) as _,
Arc::clone(&profiles) as _,
&deps,
SaveImportProfileCommand {
user_id: Uuid::new_v4(),
session_id: Uuid::new_v4(),
@@ -38,9 +43,13 @@ async fn saves_profile_from_session() {
session.field_mappings = Some(vec![]);
sessions.create(&session).await.unwrap();
let deps = SaveProfileDeps {
import_session: Arc::clone(&sessions) as _,
import_profile: Arc::clone(&profiles) as _,
};
let result = save_profile::execute(
Arc::clone(&sessions) as _,
Arc::clone(&profiles) as _,
&deps,
SaveImportProfileCommand {
user_id,
session_id: sid.value(),

View File

@@ -1,11 +1,11 @@
use std::sync::Arc;
use chrono::Duration;
use domain::{errors::DomainError, ports::WatchEventRepository};
use domain::{errors::DomainError, ports::WatchEventCommand};
pub async fn execute(watch_event: Arc<dyn WatchEventRepository>) -> Result<u64, DomainError> {
pub async fn execute(watch_event_command: Arc<dyn WatchEventCommand>) -> Result<u64, DomainError> {
let cutoff = chrono::Utc::now().naive_utc() - Duration::days(30);
watch_event.delete_non_pending_older_than(cutoff).await
watch_event_command.delete_non_pending_older_than(cutoff).await
}
#[cfg(test)]

View File

@@ -3,7 +3,7 @@ use std::sync::Arc;
use domain::{
errors::DomainError,
models::WatchEventStatus,
ports::WatchEventRepository,
ports::{WatchEventCommand, WatchEventQuery},
value_objects::{UserId, WatchEventId},
};
@@ -14,7 +14,8 @@ use crate::{
};
pub async fn execute(
watch_event: Arc<dyn WatchEventRepository>,
watch_event_command: Arc<dyn WatchEventCommand>,
watch_event_query: Arc<dyn WatchEventQuery>,
review_logger: Arc<dyn ReviewLogger>,
cmd: ConfirmWatchEventsCommand,
) -> Result<u32, DomainError> {
@@ -23,7 +24,7 @@ pub async fn execute(
for c in cmd.confirmations {
let event_id = WatchEventId::from_uuid(c.watch_event_id);
let event = watch_event
let event = watch_event_query
.get_by_id(&event_id)
.await?
.ok_or_else(|| DomainError::NotFound(format!("WatchEvent {}", c.watch_event_id)))?;
@@ -61,7 +62,7 @@ pub async fn execute(
review_logger.log_review(review_cmd).await?;
watch_event
watch_event_command
.update_status(&event_id, WatchEventStatus::Confirmed)
.await?;

View File

@@ -1,9 +1,10 @@
use std::sync::Arc;
use domain::ports::{EventPublisher, WatchEventRepository, WebhookTokenRepository};
use domain::ports::{EventPublisher, WatchEventCommand, WatchEventQuery, WebhookTokenRepository};
pub struct IngestWatchEventDeps {
pub webhook_token: Arc<dyn WebhookTokenRepository>,
pub watch_event: Arc<dyn WatchEventRepository>,
pub watch_event_command: Arc<dyn WatchEventCommand>,
pub watch_event_query: Arc<dyn WatchEventQuery>,
pub event_publisher: Arc<dyn EventPublisher>,
}

View File

@@ -3,14 +3,15 @@ use std::sync::Arc;
use domain::{
errors::DomainError,
models::WatchEventStatus,
ports::WatchEventRepository,
ports::{WatchEventCommand, WatchEventQuery},
value_objects::{UserId, WatchEventId},
};
use crate::integrations::commands::DismissWatchEventsCommand;
pub async fn execute(
watch_event: Arc<dyn WatchEventRepository>,
watch_event_command: Arc<dyn WatchEventCommand>,
watch_event_query: Arc<dyn WatchEventQuery>,
cmd: DismissWatchEventsCommand,
) -> Result<u32, DomainError> {
let user_id = UserId::from_uuid(cmd.user_id);
@@ -24,7 +25,7 @@ pub async fn execute(
.map(|id| WatchEventId::from_uuid(*id))
.collect();
let events = watch_event.get_by_ids(&ids).await?;
let events = watch_event_query.get_by_ids(&ids).await?;
if events.len() != ids.len() {
return Err(DomainError::NotFound(
@@ -37,7 +38,7 @@ pub async fn execute(
}
}
let count = watch_event
let count = watch_event_command
.update_status_batch(&ids, WatchEventStatus::Dismissed)
.await?;

View File

@@ -1,17 +1,17 @@
use std::sync::Arc;
use domain::{
errors::DomainError, models::WatchEvent, ports::WatchEventRepository, value_objects::UserId,
errors::DomainError, models::WatchEvent, ports::WatchEventQuery, value_objects::UserId,
};
use crate::integrations::queries::GetWatchQueueQuery;
pub async fn execute(
watch_event: Arc<dyn WatchEventRepository>,
watch_event_query: Arc<dyn WatchEventQuery>,
query: GetWatchQueueQuery,
) -> Result<Vec<WatchEvent>, DomainError> {
let user_id = UserId::from_uuid(query.user_id);
watch_event.list_pending(&user_id).await
watch_event_query.list_pending(&user_id).await
}
#[cfg(test)]

View File

@@ -30,7 +30,7 @@ pub async fn execute(
if let Some(ref ext_id) = external_metadata_id {
let one_hour_ago = chrono::Utc::now().naive_utc() - Duration::hours(1);
if deps
.watch_event
.watch_event_query
.find_duplicate(&user_id, ext_id, one_hour_ago)
.await?
{
@@ -49,7 +49,7 @@ pub async fn execute(
None,
);
deps.watch_event.save(&event).await?;
deps.watch_event_command.save(&event).await?;
let _ = deps
.event_publisher

View File

@@ -1,13 +1,10 @@
use std::sync::Arc;
use domain::ports::WatchEventRepository;
use domain::testing::InMemoryWatchEventRepository;
use crate::integrations::cleanup;
#[tokio::test]
async fn returns_zero_when_nothing_to_clean() {
let watch_events: Arc<dyn WatchEventRepository> = InMemoryWatchEventRepository::new();
let watch_events = InMemoryWatchEventRepository::new();
let count = cleanup::execute(watch_events).await.unwrap();

View File

@@ -1,7 +1,7 @@
use std::sync::Arc;
use domain::models::{WatchEvent, WatchEventSource};
use domain::ports::{MovieRepository, WatchEventRepository};
use domain::ports::{MovieCommand, WatchEventCommand};
use domain::testing::{InMemoryWatchEventRepository, NoopEventPublisher};
use domain::value_objects::UserId;
use uuid::Uuid;
@@ -16,7 +16,7 @@ fn noop_logger() -> Arc<dyn crate::ports::ReviewLogger> {
#[tokio::test]
async fn confirms_watch_event_via_review_logger() {
let watch_events: Arc<dyn WatchEventRepository> = InMemoryWatchEventRepository::new();
let watch_events = InMemoryWatchEventRepository::new();
let uid = Uuid::new_v4();
let event = WatchEvent::new(
@@ -32,7 +32,8 @@ async fn confirms_watch_event_via_review_logger() {
watch_events.save(&event).await.unwrap();
let result = confirm::execute(
Arc::clone(&watch_events),
Arc::clone(&watch_events) as _,
Arc::clone(&watch_events) as _,
noop_logger(),
ConfirmWatchEventsCommand {
user_id: uid,
@@ -51,10 +52,11 @@ async fn confirms_watch_event_via_review_logger() {
#[tokio::test]
async fn empty_confirmations_returns_zero() {
let watch_events: Arc<dyn WatchEventRepository> = InMemoryWatchEventRepository::new();
let watch_events = InMemoryWatchEventRepository::new();
let result = confirm::execute(
Arc::clone(&watch_events),
Arc::clone(&watch_events) as _,
Arc::clone(&watch_events) as _,
noop_logger(),
ConfirmWatchEventsCommand {
user_id: Uuid::new_v4(),
@@ -69,7 +71,7 @@ async fn empty_confirmations_returns_zero() {
#[tokio::test]
async fn confirms_event_with_external_metadata_id_and_no_movie_id() {
let watch_events: Arc<dyn WatchEventRepository> = InMemoryWatchEventRepository::new();
let watch_events = InMemoryWatchEventRepository::new();
let uid = Uuid::new_v4();
let event = WatchEvent::new(
@@ -85,7 +87,8 @@ async fn confirms_event_with_external_metadata_id_and_no_movie_id() {
watch_events.save(&event).await.unwrap();
let result = confirm::execute(
Arc::clone(&watch_events),
Arc::clone(&watch_events) as _,
Arc::clone(&watch_events) as _,
noop_logger(),
ConfirmWatchEventsCommand {
user_id: uid,
@@ -104,7 +107,7 @@ async fn confirms_event_with_external_metadata_id_and_no_movie_id() {
#[tokio::test]
async fn rejects_other_users_event() {
let watch_events: Arc<dyn WatchEventRepository> = InMemoryWatchEventRepository::new();
let watch_events = InMemoryWatchEventRepository::new();
let owner = Uuid::new_v4();
let intruder = Uuid::new_v4();
@@ -121,7 +124,8 @@ async fn rejects_other_users_event() {
watch_events.save(&event).await.unwrap();
let result = confirm::execute(
Arc::clone(&watch_events),
Arc::clone(&watch_events) as _,
Arc::clone(&watch_events) as _,
noop_logger(),
ConfirmWatchEventsCommand {
user_id: intruder,
@@ -139,10 +143,11 @@ async fn rejects_other_users_event() {
#[tokio::test]
async fn fails_when_event_not_found() {
let watch_events: Arc<dyn WatchEventRepository> = InMemoryWatchEventRepository::new();
let watch_events = InMemoryWatchEventRepository::new();
let result = confirm::execute(
Arc::clone(&watch_events),
Arc::clone(&watch_events) as _,
Arc::clone(&watch_events) as _,
noop_logger(),
ConfirmWatchEventsCommand {
user_id: Uuid::new_v4(),
@@ -160,7 +165,7 @@ async fn fails_when_event_not_found() {
#[tokio::test]
async fn confirms_event_with_movie_id() {
let watch_events: Arc<dyn WatchEventRepository> = InMemoryWatchEventRepository::new();
let watch_events = InMemoryWatchEventRepository::new();
let events = NoopEventPublisher::new();
let uid = Uuid::new_v4();
let movie_uuid = Uuid::new_v4();
@@ -194,6 +199,7 @@ async fn confirms_event_with_movie_id() {
let watchlist = domain::testing::InMemoryWatchlistRepository::new();
let review_logger: Arc<dyn crate::ports::ReviewLogger> =
Arc::new(crate::diary::review_logger::DefaultReviewLogger::new(
Arc::clone(&movies) as _,
Arc::clone(&movies) as _,
Arc::clone(&reviews) as _,
Arc::clone(&watchlist) as _,
@@ -202,7 +208,8 @@ async fn confirms_event_with_movie_id() {
));
let result = confirm::execute(
Arc::clone(&watch_events),
Arc::clone(&watch_events) as _,
Arc::clone(&watch_events) as _,
review_logger,
ConfirmWatchEventsCommand {
user_id: uid,
@@ -221,7 +228,7 @@ async fn confirms_event_with_movie_id() {
#[tokio::test]
async fn confirms_event_without_movie_id_and_without_external_metadata_id() {
let watch_events: Arc<dyn WatchEventRepository> = InMemoryWatchEventRepository::new();
let watch_events = InMemoryWatchEventRepository::new();
let uid = Uuid::new_v4();
let event = WatchEvent::new(
@@ -237,7 +244,8 @@ async fn confirms_event_without_movie_id_and_without_external_metadata_id() {
watch_events.save(&event).await.unwrap();
let result = confirm::execute(
Arc::clone(&watch_events),
Arc::clone(&watch_events) as _,
Arc::clone(&watch_events) as _,
noop_logger(),
ConfirmWatchEventsCommand {
user_id: uid,
@@ -256,7 +264,7 @@ async fn confirms_event_without_movie_id_and_without_external_metadata_id() {
#[tokio::test]
async fn confirms_multiple_events() {
let watch_events: Arc<dyn WatchEventRepository> = InMemoryWatchEventRepository::new();
let watch_events = InMemoryWatchEventRepository::new();
let uid = Uuid::new_v4();
let event1 = WatchEvent::new(
@@ -285,7 +293,8 @@ async fn confirms_multiple_events() {
watch_events.save(&event2).await.unwrap();
let result = confirm::execute(
Arc::clone(&watch_events),
Arc::clone(&watch_events) as _,
Arc::clone(&watch_events) as _,
noop_logger(),
ConfirmWatchEventsCommand {
user_id: uid,
@@ -311,7 +320,7 @@ async fn confirms_multiple_events() {
#[tokio::test]
async fn confirms_event_without_year() {
let watch_events: Arc<dyn WatchEventRepository> = InMemoryWatchEventRepository::new();
let watch_events = InMemoryWatchEventRepository::new();
let uid = Uuid::new_v4();
let event = WatchEvent::new(
@@ -327,7 +336,8 @@ async fn confirms_event_without_year() {
watch_events.save(&event).await.unwrap();
let result = confirm::execute(
Arc::clone(&watch_events),
Arc::clone(&watch_events) as _,
Arc::clone(&watch_events) as _,
noop_logger(),
ConfirmWatchEventsCommand {
user_id: uid,

View File

@@ -1,7 +1,7 @@
use std::sync::Arc;
use domain::models::{WatchEvent, WatchEventSource};
use domain::ports::WatchEventRepository;
use domain::ports::WatchEventCommand;
use domain::testing::InMemoryWatchEventRepository;
use domain::value_objects::UserId;
use uuid::Uuid;
@@ -10,10 +10,11 @@ use crate::integrations::{commands::DismissWatchEventsCommand, dismiss};
#[tokio::test]
async fn dismisses_empty_list_returns_zero() {
let events: Arc<dyn WatchEventRepository> = InMemoryWatchEventRepository::new();
let events = InMemoryWatchEventRepository::new();
let result = dismiss::execute(
Arc::clone(&events),
Arc::clone(&events) as _,
Arc::clone(&events) as _,
DismissWatchEventsCommand {
user_id: Uuid::new_v4(),
event_ids: vec![],
@@ -27,10 +28,11 @@ async fn dismisses_empty_list_returns_zero() {
#[tokio::test]
async fn fails_when_event_not_found() {
let events: Arc<dyn WatchEventRepository> = InMemoryWatchEventRepository::new();
let events = InMemoryWatchEventRepository::new();
let result = dismiss::execute(
Arc::clone(&events),
Arc::clone(&events) as _,
Arc::clone(&events) as _,
DismissWatchEventsCommand {
user_id: Uuid::new_v4(),
event_ids: vec![Uuid::new_v4()],
@@ -43,7 +45,7 @@ async fn fails_when_event_not_found() {
#[tokio::test]
async fn dismisses_existing_events() {
let watch_events: Arc<dyn WatchEventRepository> = InMemoryWatchEventRepository::new();
let watch_events = InMemoryWatchEventRepository::new();
let uid = Uuid::new_v4();
let user_id = UserId::from_uuid(uid);
@@ -71,7 +73,8 @@ async fn dismisses_existing_events() {
watch_events.save(&e2).await.unwrap();
let result = dismiss::execute(
Arc::clone(&watch_events),
Arc::clone(&watch_events) as _,
Arc::clone(&watch_events) as _,
DismissWatchEventsCommand {
user_id: uid,
event_ids: vec![id1, id2],

View File

@@ -2,7 +2,7 @@ use std::sync::Arc;
use chrono::Utc;
use domain::models::{WatchEvent, WatchEventSource};
use domain::ports::WatchEventRepository;
use domain::ports::WatchEventCommand;
use domain::testing::InMemoryWatchEventRepository;
use domain::value_objects::UserId;
use uuid::Uuid;
@@ -11,10 +11,10 @@ use crate::integrations::{get_queue, queries::GetWatchQueueQuery};
#[tokio::test]
async fn returns_empty_when_no_events() {
let events: Arc<dyn WatchEventRepository> = InMemoryWatchEventRepository::new();
let events = InMemoryWatchEventRepository::new();
let result = get_queue::execute(
Arc::clone(&events),
Arc::clone(&events) as _,
GetWatchQueueQuery {
user_id: Uuid::new_v4(),
},
@@ -27,7 +27,7 @@ async fn returns_empty_when_no_events() {
#[tokio::test]
async fn returns_pending_events() {
let events: Arc<dyn WatchEventRepository> = InMemoryWatchEventRepository::new();
let events = InMemoryWatchEventRepository::new();
let user_id = Uuid::new_v4();
let event = WatchEvent::new(
@@ -41,7 +41,7 @@ async fn returns_pending_events() {
);
events.save(&event).await.unwrap();
let result = get_queue::execute(Arc::clone(&events), GetWatchQueueQuery { user_id })
let result = get_queue::execute(Arc::clone(&events) as _, GetWatchQueueQuery { user_id })
.await
.unwrap();

View File

@@ -1,7 +1,7 @@
use std::sync::Arc;
use domain::models::WatchEventSource;
use domain::ports::{EventPublisher, WatchEventRepository, WebhookTokenRepository};
use domain::ports::{EventPublisher, WebhookTokenRepository};
use domain::testing::{
InMemoryWatchEventRepository, InMemoryWebhookTokenRepository, NoopEventPublisher,
};
@@ -30,7 +30,7 @@ impl domain::ports::MediaServerParser for FakeParser {
#[tokio::test]
async fn ingests_watch_event() {
let tokens: Arc<dyn WebhookTokenRepository> = InMemoryWebhookTokenRepository::new();
let watch_events: Arc<dyn WatchEventRepository> = InMemoryWatchEventRepository::new();
let watch_events = InMemoryWatchEventRepository::new();
let event_publisher: Arc<dyn EventPublisher> = NoopEventPublisher::new();
let user_id = Uuid::new_v4();
@@ -47,7 +47,8 @@ async fn ingests_watch_event() {
let deps = IngestWatchEventDeps {
webhook_token: Arc::clone(&tokens),
watch_event: Arc::clone(&watch_events),
watch_event_command: Arc::clone(&watch_events) as _,
watch_event_query: Arc::clone(&watch_events) as _,
event_publisher: Arc::clone(&event_publisher),
};
@@ -68,12 +69,13 @@ async fn ingests_watch_event() {
#[tokio::test]
async fn rejects_invalid_token() {
let tokens: Arc<dyn WebhookTokenRepository> = InMemoryWebhookTokenRepository::new();
let watch_events: Arc<dyn WatchEventRepository> = InMemoryWatchEventRepository::new();
let watch_events = InMemoryWatchEventRepository::new();
let event_publisher: Arc<dyn EventPublisher> = NoopEventPublisher::new();
let deps = IngestWatchEventDeps {
webhook_token: Arc::clone(&tokens),
watch_event: Arc::clone(&watch_events),
watch_event_command: Arc::clone(&watch_events) as _,
watch_event_query: Arc::clone(&watch_events) as _,
event_publisher: Arc::clone(&event_publisher),
};

View File

@@ -4,7 +4,7 @@ use std::time::Duration;
use async_trait::async_trait;
use domain::{
errors::DomainError,
ports::{MovieDeduplicator, MovieRepository, ObjectStorage, PeriodicJob},
ports::{MovieDeduplicator, MovieQuery, ObjectStorage, PeriodicJob},
};
use crate::movies::merge_duplicates::{MergeDuplicatesDeps, execute};
@@ -15,13 +15,13 @@ pub struct MovieDeduplicationJob {
impl MovieDeduplicationJob {
pub fn new(
movie: Arc<dyn MovieRepository>,
movie: Arc<dyn MovieQuery>,
deduplicator: Arc<dyn MovieDeduplicator>,
object_storage: Arc<dyn ObjectStorage>,
) -> Self {
Self {
deps: MergeDuplicatesDeps {
movie,
movie_query: movie,
deduplicator,
object_storage,
},

View File

@@ -4,15 +4,15 @@ use std::time::Duration;
use async_trait::async_trait;
use domain::{
errors::DomainError,
ports::{PeriodicJob, WatchEventRepository},
ports::{PeriodicJob, WatchEventCommand},
};
pub struct WatchEventCleanupJob {
watch_event: Arc<dyn WatchEventRepository>,
watch_event: Arc<dyn WatchEventCommand>,
}
impl WatchEventCleanupJob {
pub fn new(watch_event: Arc<dyn WatchEventRepository>) -> Self {
pub fn new(watch_event: Arc<dyn WatchEventCommand>) -> Self {
Self { watch_event }
}
}

View File

@@ -1,12 +1,13 @@
use std::sync::Arc;
use domain::ports::{
EventPublisher, MetadataClient, MovieProfileRepository, MovieRepository, ObjectStorage,
PersonCommand, PersonQuery, PosterFetcherClient, SearchCommand,
EventPublisher, MetadataClient, MovieCommand, MovieProfileRepository, MovieQuery,
ObjectStorage, PersonCommand, PersonQuery, PosterFetcherClient, SearchCommand,
};
pub struct SyncPosterDeps {
pub movie: Arc<dyn MovieRepository>,
pub movie_command: Arc<dyn MovieCommand>,
pub movie_query: Arc<dyn MovieQuery>,
pub movie_profile: Arc<dyn MovieProfileRepository>,
pub metadata: Arc<dyn MetadataClient>,
pub poster_fetcher: Arc<dyn PosterFetcherClient>,
@@ -16,14 +17,14 @@ pub struct SyncPosterDeps {
}
pub struct EnrichMovieDeps {
pub movie: Arc<dyn MovieRepository>,
pub movie_query: Arc<dyn MovieQuery>,
pub movie_profile: Arc<dyn MovieProfileRepository>,
pub person_command: Arc<dyn PersonCommand>,
pub search_command: Arc<dyn SearchCommand>,
}
pub struct ReindexSearchDeps {
pub movie: Arc<dyn MovieRepository>,
pub movie_query: Arc<dyn MovieQuery>,
pub movie_profile: Arc<dyn MovieProfileRepository>,
pub search_command: Arc<dyn SearchCommand>,
pub person_command: Arc<dyn PersonCommand>,

View File

@@ -5,20 +5,20 @@ use domain::{
errors::DomainError,
events::DomainEvent,
models::IndexableDocument,
ports::{EventHandler, MovieRepository, SearchCommand},
ports::{EventHandler, MovieQuery, SearchCommand},
};
/// Reacts to `MovieDiscovered` and inserts a bare search index entry immediately,
/// so movies are findable before TMDb enrichment runs.
/// Enrichment will later overwrite this with the full document (cast, genres, etc.).
pub struct MovieDiscoveryIndexer {
movie_repository: Arc<dyn MovieRepository>,
movie_repository: Arc<dyn MovieQuery>,
search_command: Arc<dyn SearchCommand>,
}
impl MovieDiscoveryIndexer {
pub fn new(
movie_repository: Arc<dyn MovieRepository>,
movie_repository: Arc<dyn MovieQuery>,
search_command: Arc<dyn SearchCommand>,
) -> Self {
Self {

View File

@@ -18,7 +18,7 @@ pub async fn execute(deps: &EnrichMovieDeps, cmd: EnrichMovieCommand) -> Result<
}
// 3. Fetch the movie for the search index document
let Some(movie) = deps.movie.get_movie_by_id(&cmd.movie_id).await? else {
let Some(movie) = deps.movie_query.get_movie_by_id(&cmd.movie_id).await? else {
tracing::warn!(movie_id = %cmd.movie_id.value(), "enrich_movie: movie not found after profile upsert");
return Ok(());
};

View File

@@ -6,7 +6,7 @@ use domain::{
events::DomainEvent,
models::MovieProfile,
ports::{
EventHandler, ImageFetcher, MovieEnrichmentClient, MovieProfileRepository, MovieRepository,
EventHandler, ImageFetcher, MovieEnrichmentClient, MovieProfileRepository, MovieQuery,
ObjectStorage, PersonCommand, SearchCommand,
},
};
@@ -17,7 +17,7 @@ use crate::movies::{
pub struct MovieEnrichmentHandler {
enrichment_client: Arc<dyn MovieEnrichmentClient>,
movie_repository: Arc<dyn MovieRepository>,
movie_repository: Arc<dyn MovieQuery>,
profile_repo: Arc<dyn MovieProfileRepository>,
person_command: Arc<dyn PersonCommand>,
search_command: Arc<dyn SearchCommand>,
@@ -28,7 +28,7 @@ pub struct MovieEnrichmentHandler {
impl MovieEnrichmentHandler {
pub fn new(
enrichment_client: Arc<dyn MovieEnrichmentClient>,
movie_repository: Arc<dyn MovieRepository>,
movie_repository: Arc<dyn MovieQuery>,
profile_repo: Arc<dyn MovieProfileRepository>,
person_command: Arc<dyn PersonCommand>,
search_command: Arc<dyn SearchCommand>,
@@ -92,7 +92,7 @@ impl EventHandler for MovieEnrichmentHandler {
self.download_cast_photos(&profile).await;
let enrich_deps = EnrichMovieDeps {
movie: self.movie_repository.clone(),
movie_query: self.movie_repository.clone(),
movie_profile: self.profile_repo.clone(),
person_command: self.person_command.clone(),
search_command: self.search_command.clone(),

View File

@@ -4,13 +4,13 @@ use domain::{
errors::DomainError,
models::collections::{PageParams, Paginated},
models::{MovieFilter, MovieSummary},
ports::MovieRepository,
ports::MovieQuery,
};
use crate::movies::queries::GetMoviesQuery;
pub async fn execute(
movie: Arc<dyn MovieRepository>,
movie: Arc<dyn MovieQuery>,
query: GetMoviesQuery,
) -> Result<Paginated<MovieSummary>, DomainError> {
let page = PageParams::new(query.limit, query.offset)?;

View File

@@ -2,12 +2,12 @@ use std::sync::Arc;
use domain::{
errors::DomainError,
ports::{MovieDeduplicator, MovieRepository, ObjectStorage},
ports::{MovieDeduplicator, MovieQuery, ObjectStorage},
value_objects::MovieId,
};
pub struct MergeDuplicatesDeps {
pub movie: Arc<dyn MovieRepository>,
pub movie_query: Arc<dyn MovieQuery>,
pub deduplicator: Arc<dyn MovieDeduplicator>,
pub object_storage: Arc<dyn ObjectStorage>,
}
@@ -18,7 +18,7 @@ pub struct MergeReport {
}
pub async fn execute(deps: &MergeDuplicatesDeps) -> Result<MergeReport, DomainError> {
let movies = deps.movie.list_movies_with_external_id().await?;
let movies = deps.movie_query.list_movies_with_external_id().await?;
let mut pairs_found = 0u64;
let mut rows_repointed = 0u64;
@@ -37,7 +37,7 @@ pub async fn execute(deps: &MergeDuplicatesDeps) -> Result<MergeReport, DomainEr
pairs_found += 1;
// Determine which poster will be dropped after merge
let canonical = match deps.movie.get_movie_by_id(&canonical_id).await? {
let canonical = match deps.movie_query.get_movie_by_id(&canonical_id).await? {
Some(existing) => existing,
None => domain::models::Movie::from_persistence(
canonical_id,

View File

@@ -34,7 +34,7 @@ async fn reindex_movies(deps: &ReindexSearchDeps) -> Result<u64, DomainError> {
let mut offset: u32 = 0;
loop {
let page = deps
.movie
.movie_query
.list_movies(
&PageParams {
limit: BATCH_SIZE,

View File

@@ -2,7 +2,7 @@ use domain::{
errors::DomainError,
events::DomainEvent,
models::Movie,
ports::{EventPublisher, MetadataClient, MovieRepository},
ports::{EventPublisher, MetadataClient, MovieCommand, MovieQuery},
value_objects::MovieId,
};
@@ -14,20 +14,21 @@ use crate::diary::movie_resolver::{MovieResolver, MovieResolverDeps};
/// Returns `(movie, is_new_movie)`.
pub async fn resolve_and_persist_movie(
input: &MovieInput,
movie_repo: &dyn MovieRepository,
movie_command: &dyn MovieCommand,
movie_query: &dyn MovieQuery,
metadata_client: &dyn MetadataClient,
event_publisher: &dyn EventPublisher,
) -> Result<(Movie, bool), DomainError> {
let (movie, is_new) = if let Some(id) = input.movie_id {
let movie_id = MovieId::from_uuid(id);
let movie = movie_repo
let movie = movie_query
.get_movie_by_id(&movie_id)
.await?
.ok_or_else(|| DomainError::NotFound(format!("Movie {id}")))?;
(movie, false)
} else {
let deps = MovieResolverDeps {
repository: movie_repo,
repository: movie_query,
metadata_client,
};
MovieResolver::default_pipeline()
@@ -36,7 +37,7 @@ pub async fn resolve_and_persist_movie(
};
if is_new {
movie_repo.upsert_movie(&movie).await?;
movie_command.upsert_movie(&movie).await?;
if let Some(ext_id) = movie.external_metadata_id() {
let _ = event_publisher
.publish(&DomainEvent::MovieDiscovered {

View File

@@ -10,7 +10,7 @@ use crate::{diary::commands::SyncPosterCommand, movies::deps::SyncPosterDeps};
pub async fn execute(deps: &SyncPosterDeps, cmd: SyncPosterCommand) -> Result<(), DomainError> {
let movie_id = MovieId::from_uuid(cmd.movie_id);
let mut movie = match deps.movie.get_movie_by_id(&movie_id).await? {
let mut movie = match deps.movie_query.get_movie_by_id(&movie_id).await? {
Some(m) => m,
None => {
tracing::warn!(
@@ -59,7 +59,7 @@ pub async fn execute(deps: &SyncPosterDeps, cmd: SyncPosterCommand) -> Result<()
let poster_path = PosterPath::new(stored_path)?;
movie.update_poster(poster_path);
deps.movie.upsert_movie(&movie).await?;
deps.movie_command.upsert_movie(&movie).await?;
// Refresh search index so the new poster_path is reflected immediately.
// Fetch existing profile if available for a complete index document.

View File

@@ -3,7 +3,7 @@ use std::sync::Arc;
use chrono::Utc;
use domain::{
models::{Movie, MovieProfile},
ports::MovieRepository,
ports::MovieCommand,
testing::{
FakeSearchCommand, InMemoryMovieProfileRepository, InMemoryMovieRepository,
PanicPersonCommand,
@@ -50,7 +50,7 @@ async fn stores_profile_and_indexes() {
};
let deps = EnrichMovieDeps {
movie: movie_repo as Arc<_>,
movie_query: movie_repo as Arc<_>,
movie_profile: Arc::clone(&profile_repo) as Arc<_>,
person_command: Arc::new(PanicPersonCommand),
search_command: Arc::new(FakeSearchCommand),
@@ -142,7 +142,7 @@ async fn extracts_and_indexes_persons() {
};
let deps = EnrichMovieDeps {
movie: movie_repo as Arc<_>,
movie_query: movie_repo as Arc<_>,
movie_profile: Arc::clone(&profile_repo) as Arc<_>,
person_command: Arc::new(NoopPersonCommand),
search_command: Arc::new(FakeSearchCommand),

View File

@@ -5,7 +5,7 @@ use uuid::Uuid;
use domain::{
errors::DomainError,
models::Movie,
ports::{MetadataClient, MovieRepository},
ports::{MetadataClient, MovieCommand, MovieQuery},
testing::{
FakeSearchCommand, InMemoryMovieProfileRepository, InMemoryMovieRepository,
NoopEventPublisher, NoopObjectStorage,
@@ -20,7 +20,8 @@ use crate::{
fn default_deps() -> SyncPosterDeps {
SyncPosterDeps {
movie: InMemoryMovieRepository::new(),
movie_command: InMemoryMovieRepository::new(),
movie_query: InMemoryMovieRepository::new(),
movie_profile: InMemoryMovieProfileRepository::new(),
metadata: Arc::new(domain::testing::FakeMetadataClient),
poster_fetcher: Arc::new(domain::testing::FakePosterFetcher),
@@ -59,7 +60,8 @@ async fn fails_when_no_external_id() {
movies.upsert_movie(&movie).await.unwrap();
let deps = SyncPosterDeps {
movie: Arc::clone(&movies) as _,
movie_command: Arc::clone(&movies) as _,
movie_query: Arc::clone(&movies) as _,
..default_deps()
};
@@ -103,7 +105,8 @@ async fn syncs_poster_for_movie_with_external_id() {
movies.upsert_movie(&movie).await.unwrap();
let deps = SyncPosterDeps {
movie: Arc::clone(&movies) as _,
movie_command: Arc::clone(&movies) as _,
movie_query: Arc::clone(&movies) as _,
metadata: Arc::new(FakeMetaWithPoster) as _,
..default_deps()
};

View File

@@ -7,10 +7,11 @@ use domain::{
ports::{
AuthService, DiaryExporter, DiaryRepository, DocumentParser, EventPublisher,
GoalRepository, ImportProfileRepository, ImportSessionRepository, MetadataClient,
MovieProfileRepository, MovieRepository, ObjectStorage, PasswordHasher, PersonCommand,
PersonQuery, PosterFetcherClient, RefreshSessionRepository, ReviewRepository,
SearchCommand, SearchPort, StatsRepository, UserProfileFieldsRepository, UserRepository,
UserSettingsRepository, WatchEventRepository, WatchlistRepository, WebhookTokenRepository,
MovieCommand, MovieProfileRepository, MovieQuery, ObjectStorage, PasswordHasher,
PersonCommand, PersonQuery, PosterFetcherClient, RefreshSessionRepository,
ReviewRepository, SearchCommand, SearchPort, StatsRepository,
UserProfileFieldsRepository, UserRepository, UserSettingsRepository, WatchEventCommand,
WatchEventQuery, WatchlistRepository, WebhookTokenRepository,
WrapUpRepository, WrapUpStatsQuery,
},
testing::{
@@ -40,7 +41,8 @@ impl ReviewLogger for NoopReviewLogger {
}
pub struct TestContextBuilder {
pub movie_repo: Arc<dyn MovieRepository>,
pub movie_command: Arc<dyn MovieCommand>,
pub movie_query: Arc<dyn MovieQuery>,
pub review_repo: Arc<dyn ReviewRepository>,
pub diary_repo: Arc<dyn DiaryRepository>,
pub diary_exporter: Arc<dyn DiaryExporter>,
@@ -57,7 +59,8 @@ pub struct TestContextBuilder {
pub import_profile_repo: Arc<dyn ImportProfileRepository>,
pub movie_profile_repo: Arc<dyn MovieProfileRepository>,
pub watchlist_repo: Arc<dyn WatchlistRepository>,
pub watch_event_repo: Arc<dyn WatchEventRepository>,
pub watch_event_command: Arc<dyn WatchEventCommand>,
pub watch_event_query: Arc<dyn WatchEventQuery>,
pub webhook_token_repo: Arc<dyn WebhookTokenRepository>,
pub profile_fields_repo: Arc<dyn UserProfileFieldsRepository>,
pub person_command: Arc<dyn PersonCommand>,
@@ -83,7 +86,8 @@ impl Default for TestContextBuilder {
impl TestContextBuilder {
pub fn new() -> Self {
Self {
movie_repo: InMemoryMovieRepository::new(),
movie_command: InMemoryMovieRepository::new(),
movie_query: InMemoryMovieRepository::new(),
review_repo: InMemoryReviewRepository::new(),
diary_repo: FakeDiaryRepository::new(),
diary_exporter: Arc::new(PanicDiaryExporter),
@@ -100,7 +104,8 @@ impl TestContextBuilder {
import_profile_repo: InMemoryImportProfileRepository::new(),
movie_profile_repo: InMemoryMovieProfileRepository::new(),
watchlist_repo: InMemoryWatchlistRepository::new(),
watch_event_repo: InMemoryWatchEventRepository::new(),
watch_event_command: InMemoryWatchEventRepository::new(),
watch_event_query: InMemoryWatchEventRepository::new(),
webhook_token_repo: InMemoryWebhookTokenRepository::new(),
profile_fields_repo: InMemoryProfileFieldsRepo::new(),
person_command: Arc::new(PanicPersonCommand),
@@ -128,8 +133,13 @@ impl TestContextBuilder {
}
}
pub fn with_movies(mut self, r: Arc<dyn MovieRepository>) -> Self {
self.movie_repo = r;
pub fn with_movie_command(mut self, r: Arc<dyn MovieCommand>) -> Self {
self.movie_command = r;
self
}
pub fn with_movie_query(mut self, r: Arc<dyn MovieQuery>) -> Self {
self.movie_query = r;
self
}
@@ -173,8 +183,13 @@ impl TestContextBuilder {
self
}
pub fn with_watch_events(mut self, r: Arc<dyn WatchEventRepository>) -> Self {
self.watch_event_repo = r;
pub fn with_watch_event_command(mut self, r: Arc<dyn WatchEventCommand>) -> Self {
self.watch_event_command = r;
self
}
pub fn with_watch_event_query(mut self, r: Arc<dyn WatchEventQuery>) -> Self {
self.watch_event_query = r;
self
}

View File

@@ -1,7 +1,51 @@
use domain::{errors::DomainError, events::DomainEvent, value_objects::UserId};
use domain::{
errors::DomainError,
events::DomainEvent,
ports::{EventPublisher, ObjectStorage},
value_objects::UserId,
};
use crate::users::{commands::UpdateProfileCommand, deps::UpdateProfileDeps};
async fn upload_image(
storage: &dyn ObjectStorage,
event_publisher: &dyn EventPublisher,
user_id: &UserId,
kind: &str,
old_path: Option<&str>,
new_bytes: Option<Vec<u8>>,
content_type: Option<&str>,
) -> Result<Option<String>, DomainError> {
let Some(bytes) = new_bytes else {
return Ok(old_path.map(|s| s.to_string()));
};
let ct = content_type.unwrap_or("");
if !["image/jpeg", "image/png", "image/webp"].contains(&ct) {
return Err(DomainError::ValidationError(
format!("{kind} must be jpeg, png, or webp"),
));
}
if let Some(old) = old_path {
let _ = storage.delete(old).await;
}
let key = format!("{kind}/{}", user_id.value());
let stored = storage.store(&key, &bytes).await?;
if let Err(e) = event_publisher
.publish(&DomainEvent::ImageStored {
key: stored.clone(),
})
.await
{
tracing::warn!("failed to emit ImageStored for {kind} {stored}: {e}");
}
Ok(Some(stored))
}
pub async fn execute(
deps: &UpdateProfileDeps,
cmd: UpdateProfileCommand,
@@ -14,59 +58,30 @@ pub async fn execute(
.await?
.ok_or_else(|| DomainError::NotFound("User not found".into()))?;
// Handle avatar
let new_avatar_path = if let Some(bytes) = cmd.avatar_bytes {
let content_type = cmd.avatar_content_type.as_deref().unwrap_or("");
if !["image/jpeg", "image/png", "image/webp"].contains(&content_type) {
return Err(DomainError::ValidationError(
"Avatar must be jpeg, png, or webp".into(),
));
}
if let Some(old_path) = user.avatar_path() {
let _ = deps.object_storage.delete(old_path).await;
}
let key = format!("avatars/{}", user_id.value());
let stored = deps.object_storage.store(&key, &bytes).await?;
if let Err(e) = deps
.event_publisher
.publish(&DomainEvent::ImageStored {
key: stored.clone(),
})
.await
{
tracing::warn!("failed to emit ImageStored for avatar {stored}: {e}");
}
Some(stored)
} else {
user.avatar_path().map(|s| s.to_string())
};
let storage = deps.object_storage.as_ref();
let events = deps.event_publisher.as_ref();
// Handle banner
let new_banner_path = if let Some(bytes) = cmd.banner_bytes {
let content_type = cmd.banner_content_type.as_deref().unwrap_or("");
if !["image/jpeg", "image/png", "image/webp"].contains(&content_type) {
return Err(DomainError::ValidationError(
"Banner must be jpeg, png, or webp".into(),
));
}
if let Some(old_path) = user.banner_path() {
let _ = deps.object_storage.delete(old_path).await;
}
let key = format!("banners/{}", user_id.value());
let stored = deps.object_storage.store(&key, &bytes).await?;
if let Err(e) = deps
.event_publisher
.publish(&DomainEvent::ImageStored {
key: stored.clone(),
})
.await
{
tracing::warn!("failed to emit ImageStored for banner {stored}: {e}");
}
Some(stored)
} else {
user.banner_path().map(|s| s.to_string())
};
let new_avatar_path = upload_image(
storage,
events,
&user_id,
"avatars",
user.avatar_path(),
cmd.avatar_bytes,
cmd.avatar_content_type.as_deref(),
)
.await?;
let new_banner_path = upload_image(
storage,
events,
&user_id,
"banners",
user.banner_path(),
cmd.banner_bytes,
cmd.banner_content_type.as_deref(),
)
.await?;
let moved_to = cmd.also_known_as.as_deref().and_then(|new_url| {
if user.also_known_as().map(|s| s != new_url).unwrap_or(true) {

View File

@@ -15,7 +15,8 @@ pub async fn execute(
let (movie, _is_new) = resolve_and_persist_movie(
&cmd.input,
deps.movie.as_ref(),
deps.movie_command.as_ref(),
deps.movie_query.as_ref(),
deps.metadata.as_ref(),
deps.event_publisher.as_ref(),
)

View File

@@ -1,9 +1,10 @@
use std::sync::Arc;
use domain::ports::{EventPublisher, MetadataClient, MovieRepository, WatchlistRepository};
use domain::ports::{EventPublisher, MetadataClient, MovieCommand, MovieQuery, WatchlistRepository};
pub struct WatchlistAddDeps {
pub movie: Arc<dyn MovieRepository>,
pub movie_command: Arc<dyn MovieCommand>,
pub movie_query: Arc<dyn MovieQuery>,
pub metadata: Arc<dyn MetadataClient>,
pub watchlist: Arc<dyn WatchlistRepository>,
pub event_publisher: Arc<dyn EventPublisher>,

View File

@@ -2,7 +2,7 @@ use std::sync::Arc;
use domain::{
models::Movie,
ports::MovieRepository,
ports::MovieCommand,
testing::{InMemoryMovieRepository, InMemoryWatchlistRepository, NoopEventPublisher},
value_objects::{MovieTitle, ReleaseYear},
};
@@ -17,7 +17,8 @@ fn make_deps(
watchlist: Arc<InMemoryWatchlistRepository>,
) -> WatchlistAddDeps {
WatchlistAddDeps {
movie: movies,
movie_command: Arc::clone(&movies) as _,
movie_query: movies,
metadata: Arc::new(domain::testing::FakeMetadataClient),
watchlist,
event_publisher: NoopEventPublisher::new(),