making it AP complaint
Some checks failed
CI / Check / Test (push) Has been cancelled

This commit is contained in:
2026-06-30 02:25:42 +02:00
parent 943a0abe54
commit 5ae38c834a
66 changed files with 3635 additions and 2369 deletions

View File

@@ -75,6 +75,11 @@ impl MovieRepository for RepoWithExternalMovie {
{
panic!("unexpected")
}
async fn list_movies_with_external_id(
&self,
) -> Result<Vec<domain::models::Movie>, DomainError> {
Ok(vec![])
}
}
#[async_trait::async_trait]
@@ -121,6 +126,11 @@ impl MovieRepository for RepoEmpty {
{
panic!("unexpected")
}
async fn list_movies_with_external_id(
&self,
) -> Result<Vec<domain::models::Movie>, DomainError> {
Ok(vec![])
}
}
#[async_trait::async_trait]
@@ -167,6 +177,11 @@ impl MovieRepository for RepoWithTitleMatch {
{
panic!("unexpected")
}
async fn list_movies_with_external_id(
&self,
) -> Result<Vec<domain::models::Movie>, DomainError> {
Ok(vec![])
}
}
struct MetaReturnsMovie(Movie);

View File

@@ -1,11 +1,13 @@
mod enrichment_staleness;
mod import_cleanup;
mod movie_dedup;
mod refresh_session_cleanup;
mod watch_event_cleanup;
mod wrapup;
pub use enrichment_staleness::EnrichmentStalenessJob;
pub use import_cleanup::ImportSessionCleanupJob;
pub use movie_dedup::MovieDeduplicationJob;
pub use refresh_session_cleanup::RefreshSessionCleanupJob;
pub use watch_event_cleanup::WatchEventCleanupJob;
pub use wrapup::{WrapUpAutoGenerateJob, WrapUpCleanupJob};

View File

@@ -0,0 +1,49 @@
use std::sync::Arc;
use std::time::Duration;
use async_trait::async_trait;
use domain::{
errors::DomainError,
ports::{MovieDeduplicator, MovieRepository, ObjectStorage, PeriodicJob},
};
use crate::movies::merge_duplicates::{MergeDuplicatesDeps, execute};
pub struct MovieDeduplicationJob {
deps: MergeDuplicatesDeps,
}
impl MovieDeduplicationJob {
pub fn new(
movie: Arc<dyn MovieRepository>,
deduplicator: Arc<dyn MovieDeduplicator>,
object_storage: Arc<dyn ObjectStorage>,
) -> Self {
Self {
deps: MergeDuplicatesDeps {
movie,
deduplicator,
object_storage,
},
}
}
}
#[async_trait]
impl PeriodicJob for MovieDeduplicationJob {
fn interval(&self) -> Duration {
Duration::from_secs(86_400) // once per day
}
async fn run(&self) -> Result<(), DomainError> {
let report = execute(&self.deps).await?;
if report.pairs_found > 0 {
tracing::info!(
pairs_found = report.pairs_found,
rows_repointed = report.rows_repointed,
"movie dedup: merged duplicate records"
);
}
Ok(())
}
}

View File

@@ -0,0 +1,90 @@
use std::sync::Arc;
use domain::{
errors::DomainError,
ports::{MovieDeduplicator, MovieRepository, ObjectStorage},
value_objects::MovieId,
};
pub struct MergeDuplicatesDeps {
pub movie: Arc<dyn MovieRepository>,
pub deduplicator: Arc<dyn MovieDeduplicator>,
pub object_storage: Arc<dyn ObjectStorage>,
}
pub struct MergeReport {
pub pairs_found: u64,
pub rows_repointed: u64,
}
pub async fn execute(deps: &MergeDuplicatesDeps) -> Result<MergeReport, DomainError> {
let movies = deps.movie.list_movies_with_external_id().await?;
let mut pairs_found = 0u64;
let mut rows_repointed = 0u64;
for movie in movies {
let external_id = match movie.external_metadata_id() {
Some(id) => id,
None => continue,
};
let canonical_id = MovieId::from_external(external_id);
if movie.id() == &canonical_id {
continue; // already canonical
}
pairs_found += 1;
// Determine which poster will be dropped after merge
let canonical = match deps.movie.get_movie_by_id(&canonical_id).await? {
Some(existing) => existing,
None => domain::models::Movie::from_persistence(
canonical_id,
movie.external_metadata_id().cloned(),
movie.title().clone(),
movie.release_year().clone(),
movie.director().map(str::to_string),
movie.poster_path().cloned(),
),
};
// The COALESCE in merge_into_canonical keeps canonical's poster if it has one,
// otherwise takes old's. Work out which poster key will be orphaned.
let orphaned_poster = match (canonical.poster_path(), movie.poster_path()) {
(Some(_), Some(old_poster)) if canonical.poster_path() != movie.poster_path() => {
// Canonical wins — old movie's poster will be orphaned
Some(old_poster.value().to_string())
}
(None, Some(_)) => None, // old poster moves to canonical, nothing orphaned
_ => None,
};
let repointed = deps
.deduplicator
.merge_into_canonical(movie.id(), &canonical)
.await?;
// Delete the orphaned poster file from object storage
if let Some(key) = orphaned_poster
&& let Err(e) = deps.object_storage.delete(&key).await
{
tracing::warn!(key, "failed to delete orphaned poster: {e}");
}
rows_repointed += repointed;
tracing::info!(
old_id = %movie.id().value(),
canonical_id = %canonical.id().value(),
external_id = %external_id.value(),
rows_repointed = repointed,
"merged duplicate movie"
);
}
Ok(MergeReport {
pairs_found,
rows_repointed,
})
}

View File

@@ -5,6 +5,7 @@ pub mod enrich_movie;
pub mod event_handler;
pub mod get_movie_profile;
pub mod get_movies;
pub mod merge_duplicates;
pub mod queries;
pub mod reindex_search;
pub mod request_enrichment;

View File

@@ -84,6 +84,9 @@ impl EventHandler for RecordingHandler {
| DomainEvent::GoalUpdated { .. }
| DomainEvent::GoalDeleted { .. } => "goal",
DomainEvent::PersonEnrichmentRequested { .. } => "person_enrichment_requested",
DomainEvent::UserDeleted { .. } | DomainEvent::UserAccountMoved { .. } => {
"user_lifecycle"
}
};
self.calls.lock().unwrap().push(label);
Ok(())

View File

@@ -0,0 +1,19 @@
use domain::{errors::DomainError, events::DomainEvent, value_objects::UserId};
use crate::users::deps::UpdateProfileDeps;
pub async fn execute(deps: &UpdateProfileDeps, user_id: uuid::Uuid) -> Result<(), DomainError> {
let uid = UserId::from_uuid(user_id);
deps.user
.find_by_id(&uid)
.await?
.ok_or_else(|| DomainError::NotFound("User not found".into()))?;
// Notify federation peers before any data is removed so they can process the tombstone.
deps.event_publisher
.publish(&DomainEvent::UserDeleted { user_id: uid })
.await?;
Ok(())
}

View File

@@ -1,4 +1,5 @@
pub mod commands;
pub mod delete_account;
pub mod deps;
pub mod get_current_profile;
pub mod get_profile;

View File

@@ -68,6 +68,14 @@ pub async fn execute(
user.banner_path().map(|s| s.to_string())
};
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) {
Some(new_url.to_string())
} else {
None
}
});
deps.user
.update_profile(
&user_id,
@@ -83,9 +91,21 @@ pub async fn execute(
.await?;
deps.event_publisher
.publish(&DomainEvent::UserUpdated { user_id })
.publish(&DomainEvent::UserUpdated {
user_id: user_id.clone(),
})
.await?;
if let Some(new_actor_url) = moved_to {
let _ = deps
.event_publisher
.publish(&DomainEvent::UserAccountMoved {
user_id,
new_actor_url,
})
.await;
}
Ok(())
}