Refactor event payload handling across adapters
- Introduced `event-payload` crate to centralize event payload definitions. - Updated NATS and PostgreSQL adapters to use the new `EventPayload` type. - Removed redundant event payload definitions and conversion implementations from NATS and PostgreSQL adapters. - Simplified SQLite event queue to utilize the new `EventPayload`. - Refactored wiring functions for PostgreSQL and SQLite to improve database connection handling and migration. - Cleaned up presentation and worker crates by removing unused event publisher dependencies and related wiring functions.
This commit is contained in:
17
Cargo.lock
generated
17
Cargo.lock
generated
@@ -1625,6 +1625,17 @@ dependencies = [
|
|||||||
"pin-project-lite",
|
"pin-project-lite",
|
||||||
]
|
]
|
||||||
|
|
||||||
|
[[package]]
|
||||||
|
name = "event-payload"
|
||||||
|
version = "0.1.0"
|
||||||
|
dependencies = [
|
||||||
|
"anyhow",
|
||||||
|
"chrono",
|
||||||
|
"domain",
|
||||||
|
"serde",
|
||||||
|
"uuid",
|
||||||
|
]
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "event-publisher"
|
name = "event-publisher"
|
||||||
version = "0.1.0"
|
version = "0.1.0"
|
||||||
@@ -2829,6 +2840,7 @@ dependencies = [
|
|||||||
"async-trait",
|
"async-trait",
|
||||||
"chrono",
|
"chrono",
|
||||||
"domain",
|
"domain",
|
||||||
|
"event-payload",
|
||||||
"futures",
|
"futures",
|
||||||
"serde",
|
"serde",
|
||||||
"serde_json",
|
"serde_json",
|
||||||
@@ -3373,6 +3385,7 @@ dependencies = [
|
|||||||
name = "postgres"
|
name = "postgres"
|
||||||
version = "0.1.0"
|
version = "0.1.0"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
|
"anyhow",
|
||||||
"async-trait",
|
"async-trait",
|
||||||
"chrono",
|
"chrono",
|
||||||
"domain",
|
"domain",
|
||||||
@@ -3390,6 +3403,7 @@ dependencies = [
|
|||||||
"async-trait",
|
"async-trait",
|
||||||
"chrono",
|
"chrono",
|
||||||
"domain",
|
"domain",
|
||||||
|
"event-payload",
|
||||||
"futures",
|
"futures",
|
||||||
"serde",
|
"serde",
|
||||||
"serde_json",
|
"serde_json",
|
||||||
@@ -3454,7 +3468,6 @@ dependencies = [
|
|||||||
"doc",
|
"doc",
|
||||||
"domain",
|
"domain",
|
||||||
"dotenvy",
|
"dotenvy",
|
||||||
"event-publisher",
|
|
||||||
"export",
|
"export",
|
||||||
"http-body-util",
|
"http-body-util",
|
||||||
"infer",
|
"infer",
|
||||||
@@ -4532,6 +4545,7 @@ dependencies = [
|
|||||||
"async-trait",
|
"async-trait",
|
||||||
"chrono",
|
"chrono",
|
||||||
"domain",
|
"domain",
|
||||||
|
"event-payload",
|
||||||
"futures",
|
"futures",
|
||||||
"serde",
|
"serde",
|
||||||
"serde_json",
|
"serde_json",
|
||||||
@@ -6242,7 +6256,6 @@ dependencies = [
|
|||||||
"chrono",
|
"chrono",
|
||||||
"domain",
|
"domain",
|
||||||
"dotenvy",
|
"dotenvy",
|
||||||
"event-publisher",
|
|
||||||
"export",
|
"export",
|
||||||
"futures",
|
"futures",
|
||||||
"metadata",
|
"metadata",
|
||||||
|
|||||||
@@ -16,6 +16,7 @@ members = [
|
|||||||
"crates/adapters/activitypub",
|
"crates/adapters/activitypub",
|
||||||
"crates/adapters/activitypub-base",
|
"crates/adapters/activitypub-base",
|
||||||
"crates/adapters/export",
|
"crates/adapters/export",
|
||||||
|
"crates/adapters/event-payload",
|
||||||
"crates/adapters/nats",
|
"crates/adapters/nats",
|
||||||
"crates/application",
|
"crates/application",
|
||||||
"crates/domain",
|
"crates/domain",
|
||||||
@@ -67,6 +68,7 @@ template-askama = { path = "crates/adapters/template-askama" }
|
|||||||
activitypub = { path = "crates/adapters/activitypub" }
|
activitypub = { path = "crates/adapters/activitypub" }
|
||||||
activitypub-base = { path = "crates/adapters/activitypub-base" }
|
activitypub-base = { path = "crates/adapters/activitypub-base" }
|
||||||
doc = { path = "crates/doc" }
|
doc = { path = "crates/doc" }
|
||||||
|
event-payload = { path = "crates/adapters/event-payload" }
|
||||||
nats = { path = "crates/adapters/nats" }
|
nats = { path = "crates/adapters/nats" }
|
||||||
sqlite-event-queue = { path = "crates/adapters/sqlite-event-queue" }
|
sqlite-event-queue = { path = "crates/adapters/sqlite-event-queue" }
|
||||||
postgres-event-queue = { path = "crates/adapters/postgres-event-queue" }
|
postgres-event-queue = { path = "crates/adapters/postgres-event-queue" }
|
||||||
|
|||||||
@@ -10,6 +10,7 @@ COPY .sqlx ./.sqlx
|
|||||||
COPY crates/adapters/activitypub/Cargo.toml crates/adapters/activitypub/Cargo.toml
|
COPY crates/adapters/activitypub/Cargo.toml crates/adapters/activitypub/Cargo.toml
|
||||||
COPY crates/adapters/activitypub-base/Cargo.toml crates/adapters/activitypub-base/Cargo.toml
|
COPY crates/adapters/activitypub-base/Cargo.toml crates/adapters/activitypub-base/Cargo.toml
|
||||||
COPY crates/adapters/auth/Cargo.toml crates/adapters/auth/Cargo.toml
|
COPY crates/adapters/auth/Cargo.toml crates/adapters/auth/Cargo.toml
|
||||||
|
COPY crates/adapters/event-payload/Cargo.toml crates/adapters/event-payload/Cargo.toml
|
||||||
COPY crates/adapters/event-publisher/Cargo.toml crates/adapters/event-publisher/Cargo.toml
|
COPY crates/adapters/event-publisher/Cargo.toml crates/adapters/event-publisher/Cargo.toml
|
||||||
COPY crates/adapters/nats/Cargo.toml crates/adapters/nats/Cargo.toml
|
COPY crates/adapters/nats/Cargo.toml crates/adapters/nats/Cargo.toml
|
||||||
COPY crates/adapters/metadata/Cargo.toml crates/adapters/metadata/Cargo.toml
|
COPY crates/adapters/metadata/Cargo.toml crates/adapters/metadata/Cargo.toml
|
||||||
|
|||||||
11
crates/adapters/event-payload/Cargo.toml
Normal file
11
crates/adapters/event-payload/Cargo.toml
Normal file
@@ -0,0 +1,11 @@
|
|||||||
|
[package]
|
||||||
|
name = "event-payload"
|
||||||
|
version = "0.1.0"
|
||||||
|
edition = "2024"
|
||||||
|
|
||||||
|
[dependencies]
|
||||||
|
domain = { workspace = true }
|
||||||
|
serde = { workspace = true }
|
||||||
|
chrono = { workspace = true }
|
||||||
|
uuid = { workspace = true }
|
||||||
|
anyhow = { workspace = true }
|
||||||
189
crates/adapters/event-payload/src/lib.rs
Normal file
189
crates/adapters/event-payload/src/lib.rs
Normal file
@@ -0,0 +1,189 @@
|
|||||||
|
use chrono::NaiveDateTime;
|
||||||
|
use domain::{
|
||||||
|
errors::DomainError,
|
||||||
|
events::DomainEvent,
|
||||||
|
value_objects::{ExternalMetadataId, MovieId, Rating, ReviewId, UserId},
|
||||||
|
};
|
||||||
|
use serde::{Deserialize, Serialize};
|
||||||
|
use uuid::Uuid;
|
||||||
|
|
||||||
|
#[derive(Debug, Serialize, Deserialize, PartialEq)]
|
||||||
|
#[serde(tag = "type", content = "data")]
|
||||||
|
pub enum EventPayload {
|
||||||
|
ReviewLogged {
|
||||||
|
review_id: String,
|
||||||
|
movie_id: String,
|
||||||
|
user_id: String,
|
||||||
|
rating: u8,
|
||||||
|
watched_at: i64,
|
||||||
|
},
|
||||||
|
ReviewUpdated {
|
||||||
|
review_id: String,
|
||||||
|
movie_id: String,
|
||||||
|
user_id: String,
|
||||||
|
rating: u8,
|
||||||
|
watched_at: i64,
|
||||||
|
},
|
||||||
|
MovieDiscovered {
|
||||||
|
movie_id: String,
|
||||||
|
external_metadata_id: String,
|
||||||
|
},
|
||||||
|
}
|
||||||
|
|
||||||
|
impl EventPayload {
|
||||||
|
pub fn event_type(&self) -> &'static str {
|
||||||
|
match self {
|
||||||
|
EventPayload::ReviewLogged { .. } => "ReviewLogged",
|
||||||
|
EventPayload::ReviewUpdated { .. } => "ReviewUpdated",
|
||||||
|
EventPayload::MovieDiscovered { .. } => "MovieDiscovered",
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn parse_uuid(s: &str, field: &str) -> Result<Uuid, DomainError> {
|
||||||
|
Uuid::parse_str(s)
|
||||||
|
.map_err(|e| DomainError::InfrastructureError(format!("{field}: {e}")))
|
||||||
|
}
|
||||||
|
|
||||||
|
fn parse_ts(ts: i64) -> Result<NaiveDateTime, DomainError> {
|
||||||
|
chrono::DateTime::from_timestamp(ts, 0)
|
||||||
|
.map(|dt| dt.naive_utc())
|
||||||
|
.ok_or_else(|| DomainError::InfrastructureError(format!("invalid timestamp: {ts}")))
|
||||||
|
}
|
||||||
|
|
||||||
|
impl From<&DomainEvent> for EventPayload {
|
||||||
|
fn from(event: &DomainEvent) -> Self {
|
||||||
|
match event {
|
||||||
|
DomainEvent::ReviewLogged { review_id, movie_id, user_id, rating, watched_at } => {
|
||||||
|
EventPayload::ReviewLogged {
|
||||||
|
review_id: review_id.value().to_string(),
|
||||||
|
movie_id: movie_id.value().to_string(),
|
||||||
|
user_id: user_id.value().to_string(),
|
||||||
|
rating: rating.value(),
|
||||||
|
watched_at: watched_at.and_utc().timestamp(),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
DomainEvent::ReviewUpdated { review_id, movie_id, user_id, rating, watched_at } => {
|
||||||
|
EventPayload::ReviewUpdated {
|
||||||
|
review_id: review_id.value().to_string(),
|
||||||
|
movie_id: movie_id.value().to_string(),
|
||||||
|
user_id: user_id.value().to_string(),
|
||||||
|
rating: rating.value(),
|
||||||
|
watched_at: watched_at.and_utc().timestamp(),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
DomainEvent::MovieDiscovered { movie_id, external_metadata_id } => {
|
||||||
|
EventPayload::MovieDiscovered {
|
||||||
|
movie_id: movie_id.value().to_string(),
|
||||||
|
external_metadata_id: external_metadata_id.value().to_owned(),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
impl TryFrom<EventPayload> for DomainEvent {
|
||||||
|
type Error = DomainError;
|
||||||
|
fn try_from(payload: EventPayload) -> Result<Self, DomainError> {
|
||||||
|
match payload {
|
||||||
|
EventPayload::ReviewLogged { review_id, movie_id, user_id, rating, watched_at } => {
|
||||||
|
Ok(DomainEvent::ReviewLogged {
|
||||||
|
review_id: ReviewId::from_uuid(parse_uuid(&review_id, "review_id")?),
|
||||||
|
movie_id: MovieId::from_uuid(parse_uuid(&movie_id, "movie_id")?),
|
||||||
|
user_id: UserId::from_uuid(parse_uuid(&user_id, "user_id")?),
|
||||||
|
rating: Rating::new(rating)?,
|
||||||
|
watched_at: parse_ts(watched_at)?,
|
||||||
|
})
|
||||||
|
}
|
||||||
|
EventPayload::ReviewUpdated { review_id, movie_id, user_id, rating, watched_at } => {
|
||||||
|
Ok(DomainEvent::ReviewUpdated {
|
||||||
|
review_id: ReviewId::from_uuid(parse_uuid(&review_id, "review_id")?),
|
||||||
|
movie_id: MovieId::from_uuid(parse_uuid(&movie_id, "movie_id")?),
|
||||||
|
user_id: UserId::from_uuid(parse_uuid(&user_id, "user_id")?),
|
||||||
|
rating: Rating::new(rating)?,
|
||||||
|
watched_at: parse_ts(watched_at)?,
|
||||||
|
})
|
||||||
|
}
|
||||||
|
EventPayload::MovieDiscovered { movie_id, external_metadata_id } => {
|
||||||
|
Ok(DomainEvent::MovieDiscovered {
|
||||||
|
movie_id: MovieId::from_uuid(parse_uuid(&movie_id, "movie_id")?),
|
||||||
|
external_metadata_id: ExternalMetadataId::new(external_metadata_id)?,
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
mod tests {
|
||||||
|
use super::*;
|
||||||
|
|
||||||
|
fn fixed_dt() -> NaiveDateTime {
|
||||||
|
chrono::DateTime::from_timestamp(1_700_000_000, 0).unwrap().naive_utc()
|
||||||
|
}
|
||||||
|
|
||||||
|
fn review_logged() -> DomainEvent {
|
||||||
|
DomainEvent::ReviewLogged {
|
||||||
|
review_id: ReviewId::from_uuid(Uuid::new_v4()),
|
||||||
|
movie_id: MovieId::from_uuid(Uuid::new_v4()),
|
||||||
|
user_id: UserId::from_uuid(Uuid::new_v4()),
|
||||||
|
rating: Rating::new(4).unwrap(),
|
||||||
|
watched_at: fixed_dt(),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn review_updated() -> DomainEvent {
|
||||||
|
DomainEvent::ReviewUpdated {
|
||||||
|
review_id: ReviewId::from_uuid(Uuid::new_v4()),
|
||||||
|
movie_id: MovieId::from_uuid(Uuid::new_v4()),
|
||||||
|
user_id: UserId::from_uuid(Uuid::new_v4()),
|
||||||
|
rating: Rating::new(3).unwrap(),
|
||||||
|
watched_at: fixed_dt(),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn movie_discovered() -> DomainEvent {
|
||||||
|
DomainEvent::MovieDiscovered {
|
||||||
|
movie_id: MovieId::from_uuid(Uuid::new_v4()),
|
||||||
|
external_metadata_id: ExternalMetadataId::new("tt1234567".into()).unwrap(),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn round_trip(event: DomainEvent) {
|
||||||
|
let payload = EventPayload::from(&event);
|
||||||
|
let json = serde_json::to_string(&payload).expect("serialize");
|
||||||
|
let back: EventPayload = serde_json::from_str(&json).expect("deserialize");
|
||||||
|
let recovered = DomainEvent::try_from(back).expect("try_from");
|
||||||
|
assert_eq!(EventPayload::from(&event), EventPayload::from(&recovered));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn round_trip_review_logged() {
|
||||||
|
round_trip(review_logged());
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn round_trip_review_updated() {
|
||||||
|
round_trip(review_updated());
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn round_trip_movie_discovered() {
|
||||||
|
round_trip(movie_discovered());
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn serialized_format_is_tagged() {
|
||||||
|
let payload = EventPayload::from(&movie_discovered());
|
||||||
|
let json = serde_json::to_string(&payload).unwrap();
|
||||||
|
assert!(json.contains(r#""type":"MovieDiscovered""#));
|
||||||
|
assert!(json.contains(r#""data":"#));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn event_type_strings() {
|
||||||
|
assert_eq!(EventPayload::from(&review_logged()).event_type(), "ReviewLogged");
|
||||||
|
assert_eq!(EventPayload::from(&review_updated()).event_type(), "ReviewUpdated");
|
||||||
|
assert_eq!(EventPayload::from(&movie_discovered()).event_type(), "MovieDiscovered");
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -7,6 +7,7 @@ edition = "2024"
|
|||||||
async-nats = "0.48.0"
|
async-nats = "0.48.0"
|
||||||
|
|
||||||
domain = { workspace = true }
|
domain = { workspace = true }
|
||||||
|
event-payload = { workspace = true }
|
||||||
async-trait = { workspace = true }
|
async-trait = { workspace = true }
|
||||||
anyhow = { workspace = true }
|
anyhow = { workspace = true }
|
||||||
thiserror = { workspace = true }
|
thiserror = { workspace = true }
|
||||||
|
|||||||
@@ -1,172 +1 @@
|
|||||||
use chrono::NaiveDateTime;
|
pub use event_payload::EventPayload as NatsEventPayload;
|
||||||
use domain::{
|
|
||||||
errors::DomainError,
|
|
||||||
events::DomainEvent,
|
|
||||||
value_objects::{ExternalMetadataId, MovieId, Rating, ReviewId, UserId},
|
|
||||||
};
|
|
||||||
use serde::{Deserialize, Serialize};
|
|
||||||
use uuid::Uuid;
|
|
||||||
|
|
||||||
#[derive(Debug, Serialize, Deserialize, PartialEq)]
|
|
||||||
#[serde(tag = "type", content = "data")]
|
|
||||||
pub enum NatsEventPayload {
|
|
||||||
ReviewLogged {
|
|
||||||
review_id: String,
|
|
||||||
movie_id: String,
|
|
||||||
user_id: String,
|
|
||||||
rating: u8,
|
|
||||||
watched_at: i64,
|
|
||||||
},
|
|
||||||
ReviewUpdated {
|
|
||||||
review_id: String,
|
|
||||||
movie_id: String,
|
|
||||||
user_id: String,
|
|
||||||
rating: u8,
|
|
||||||
watched_at: i64,
|
|
||||||
},
|
|
||||||
MovieDiscovered {
|
|
||||||
movie_id: String,
|
|
||||||
external_metadata_id: String,
|
|
||||||
},
|
|
||||||
}
|
|
||||||
|
|
||||||
fn parse_uuid(s: &str, field: &str) -> Result<Uuid, DomainError> {
|
|
||||||
Uuid::parse_str(s)
|
|
||||||
.map_err(|e| DomainError::InfrastructureError(format!("{field}: {e}")))
|
|
||||||
}
|
|
||||||
|
|
||||||
fn parse_ts(ts: i64) -> Result<NaiveDateTime, DomainError> {
|
|
||||||
chrono::DateTime::from_timestamp(ts, 0)
|
|
||||||
.map(|dt| dt.naive_utc())
|
|
||||||
.ok_or_else(|| DomainError::InfrastructureError(format!("invalid timestamp: {ts}")))
|
|
||||||
}
|
|
||||||
|
|
||||||
impl From<&DomainEvent> for NatsEventPayload {
|
|
||||||
fn from(event: &DomainEvent) -> Self {
|
|
||||||
match event {
|
|
||||||
DomainEvent::ReviewLogged { review_id, movie_id, user_id, rating, watched_at } => {
|
|
||||||
NatsEventPayload::ReviewLogged {
|
|
||||||
review_id: review_id.value().to_string(),
|
|
||||||
movie_id: movie_id.value().to_string(),
|
|
||||||
user_id: user_id.value().to_string(),
|
|
||||||
rating: rating.value(),
|
|
||||||
watched_at: watched_at.and_utc().timestamp(),
|
|
||||||
}
|
|
||||||
}
|
|
||||||
DomainEvent::ReviewUpdated { review_id, movie_id, user_id, rating, watched_at } => {
|
|
||||||
NatsEventPayload::ReviewUpdated {
|
|
||||||
review_id: review_id.value().to_string(),
|
|
||||||
movie_id: movie_id.value().to_string(),
|
|
||||||
user_id: user_id.value().to_string(),
|
|
||||||
rating: rating.value(),
|
|
||||||
watched_at: watched_at.and_utc().timestamp(),
|
|
||||||
}
|
|
||||||
}
|
|
||||||
DomainEvent::MovieDiscovered { movie_id, external_metadata_id } => {
|
|
||||||
NatsEventPayload::MovieDiscovered {
|
|
||||||
movie_id: movie_id.value().to_string(),
|
|
||||||
external_metadata_id: external_metadata_id.value().to_owned(),
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
impl TryFrom<NatsEventPayload> for DomainEvent {
|
|
||||||
type Error = DomainError;
|
|
||||||
fn try_from(payload: NatsEventPayload) -> Result<Self, DomainError> {
|
|
||||||
match payload {
|
|
||||||
NatsEventPayload::ReviewLogged { review_id, movie_id, user_id, rating, watched_at } => {
|
|
||||||
Ok(DomainEvent::ReviewLogged {
|
|
||||||
review_id: ReviewId::from_uuid(parse_uuid(&review_id, "review_id")?),
|
|
||||||
movie_id: MovieId::from_uuid(parse_uuid(&movie_id, "movie_id")?),
|
|
||||||
user_id: UserId::from_uuid(parse_uuid(&user_id, "user_id")?),
|
|
||||||
rating: Rating::new(rating)?,
|
|
||||||
watched_at: parse_ts(watched_at)?,
|
|
||||||
})
|
|
||||||
}
|
|
||||||
NatsEventPayload::ReviewUpdated { review_id, movie_id, user_id, rating, watched_at } => {
|
|
||||||
Ok(DomainEvent::ReviewUpdated {
|
|
||||||
review_id: ReviewId::from_uuid(parse_uuid(&review_id, "review_id")?),
|
|
||||||
movie_id: MovieId::from_uuid(parse_uuid(&movie_id, "movie_id")?),
|
|
||||||
user_id: UserId::from_uuid(parse_uuid(&user_id, "user_id")?),
|
|
||||||
rating: Rating::new(rating)?,
|
|
||||||
watched_at: parse_ts(watched_at)?,
|
|
||||||
})
|
|
||||||
}
|
|
||||||
NatsEventPayload::MovieDiscovered { movie_id, external_metadata_id } => {
|
|
||||||
Ok(DomainEvent::MovieDiscovered {
|
|
||||||
movie_id: MovieId::from_uuid(parse_uuid(&movie_id, "movie_id")?),
|
|
||||||
external_metadata_id: ExternalMetadataId::new(external_metadata_id)?,
|
|
||||||
})
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
#[cfg(test)]
|
|
||||||
mod tests {
|
|
||||||
use super::*;
|
|
||||||
|
|
||||||
fn fixed_dt() -> NaiveDateTime {
|
|
||||||
chrono::DateTime::from_timestamp(1_700_000_000, 0).unwrap().naive_utc()
|
|
||||||
}
|
|
||||||
|
|
||||||
fn review_logged() -> DomainEvent {
|
|
||||||
DomainEvent::ReviewLogged {
|
|
||||||
review_id: ReviewId::from_uuid(Uuid::new_v4()),
|
|
||||||
movie_id: MovieId::from_uuid(Uuid::new_v4()),
|
|
||||||
user_id: UserId::from_uuid(Uuid::new_v4()),
|
|
||||||
rating: Rating::new(4).unwrap(),
|
|
||||||
watched_at: fixed_dt(),
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
fn review_updated() -> DomainEvent {
|
|
||||||
DomainEvent::ReviewUpdated {
|
|
||||||
review_id: ReviewId::from_uuid(Uuid::new_v4()),
|
|
||||||
movie_id: MovieId::from_uuid(Uuid::new_v4()),
|
|
||||||
user_id: UserId::from_uuid(Uuid::new_v4()),
|
|
||||||
rating: Rating::new(3).unwrap(),
|
|
||||||
watched_at: fixed_dt(),
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
fn movie_discovered() -> DomainEvent {
|
|
||||||
DomainEvent::MovieDiscovered {
|
|
||||||
movie_id: MovieId::from_uuid(Uuid::new_v4()),
|
|
||||||
external_metadata_id: ExternalMetadataId::new("tt1234567".into()).unwrap(),
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
fn round_trip(event: DomainEvent) {
|
|
||||||
let payload = NatsEventPayload::from(&event);
|
|
||||||
let json = serde_json::to_string(&payload).expect("serialize");
|
|
||||||
let back: NatsEventPayload = serde_json::from_str(&json).expect("deserialize");
|
|
||||||
let recovered = DomainEvent::try_from(back).expect("try_from");
|
|
||||||
assert_eq!(NatsEventPayload::from(&event), NatsEventPayload::from(&recovered));
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn round_trip_review_logged() {
|
|
||||||
round_trip(review_logged());
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn round_trip_review_updated() {
|
|
||||||
round_trip(review_updated());
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn round_trip_movie_discovered() {
|
|
||||||
round_trip(movie_discovered());
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn serialized_format_is_tagged() {
|
|
||||||
let payload = NatsEventPayload::from(&movie_discovered());
|
|
||||||
let json = serde_json::to_string(&payload).unwrap();
|
|
||||||
assert!(json.contains(r#""type":"MovieDiscovered""#));
|
|
||||||
assert!(json.contains(r#""data":"#));
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|||||||
@@ -6,6 +6,7 @@ edition = "2024"
|
|||||||
[dependencies]
|
[dependencies]
|
||||||
sqlx = { version = "0.8.6", features = ["runtime-tokio-rustls", "postgres", "macros", "chrono"] }
|
sqlx = { version = "0.8.6", features = ["runtime-tokio-rustls", "postgres", "macros", "chrono"] }
|
||||||
domain = { workspace = true }
|
domain = { workspace = true }
|
||||||
|
event-payload = { workspace = true }
|
||||||
anyhow = { workspace = true }
|
anyhow = { workspace = true }
|
||||||
async-trait = { workspace = true }
|
async-trait = { workspace = true }
|
||||||
serde = { workspace = true }
|
serde = { workspace = true }
|
||||||
|
|||||||
@@ -1,189 +1 @@
|
|||||||
use chrono::NaiveDateTime;
|
pub use event_payload::EventPayload as DbEventPayload;
|
||||||
use domain::{
|
|
||||||
errors::DomainError,
|
|
||||||
events::DomainEvent,
|
|
||||||
value_objects::{ExternalMetadataId, MovieId, Rating, ReviewId, UserId},
|
|
||||||
};
|
|
||||||
use serde::{Deserialize, Serialize};
|
|
||||||
use uuid::Uuid;
|
|
||||||
|
|
||||||
#[derive(Debug, Serialize, Deserialize, PartialEq)]
|
|
||||||
#[serde(tag = "type", content = "data")]
|
|
||||||
pub enum DbEventPayload {
|
|
||||||
ReviewLogged {
|
|
||||||
review_id: String,
|
|
||||||
movie_id: String,
|
|
||||||
user_id: String,
|
|
||||||
rating: u8,
|
|
||||||
watched_at: i64,
|
|
||||||
},
|
|
||||||
ReviewUpdated {
|
|
||||||
review_id: String,
|
|
||||||
movie_id: String,
|
|
||||||
user_id: String,
|
|
||||||
rating: u8,
|
|
||||||
watched_at: i64,
|
|
||||||
},
|
|
||||||
MovieDiscovered {
|
|
||||||
movie_id: String,
|
|
||||||
external_metadata_id: String,
|
|
||||||
},
|
|
||||||
}
|
|
||||||
|
|
||||||
impl DbEventPayload {
|
|
||||||
pub fn event_type(&self) -> &'static str {
|
|
||||||
match self {
|
|
||||||
DbEventPayload::ReviewLogged { .. } => "ReviewLogged",
|
|
||||||
DbEventPayload::ReviewUpdated { .. } => "ReviewUpdated",
|
|
||||||
DbEventPayload::MovieDiscovered { .. } => "MovieDiscovered",
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
fn parse_uuid(s: &str, field: &str) -> Result<Uuid, DomainError> {
|
|
||||||
Uuid::parse_str(s)
|
|
||||||
.map_err(|e| DomainError::InfrastructureError(format!("{field}: {e}")))
|
|
||||||
}
|
|
||||||
|
|
||||||
fn parse_ts(ts: i64) -> Result<NaiveDateTime, DomainError> {
|
|
||||||
chrono::DateTime::from_timestamp(ts, 0)
|
|
||||||
.map(|dt| dt.naive_utc())
|
|
||||||
.ok_or_else(|| DomainError::InfrastructureError(format!("invalid timestamp: {ts}")))
|
|
||||||
}
|
|
||||||
|
|
||||||
impl From<&DomainEvent> for DbEventPayload {
|
|
||||||
fn from(event: &DomainEvent) -> Self {
|
|
||||||
match event {
|
|
||||||
DomainEvent::ReviewLogged { review_id, movie_id, user_id, rating, watched_at } => {
|
|
||||||
DbEventPayload::ReviewLogged {
|
|
||||||
review_id: review_id.value().to_string(),
|
|
||||||
movie_id: movie_id.value().to_string(),
|
|
||||||
user_id: user_id.value().to_string(),
|
|
||||||
rating: rating.value(),
|
|
||||||
watched_at: watched_at.and_utc().timestamp(),
|
|
||||||
}
|
|
||||||
}
|
|
||||||
DomainEvent::ReviewUpdated { review_id, movie_id, user_id, rating, watched_at } => {
|
|
||||||
DbEventPayload::ReviewUpdated {
|
|
||||||
review_id: review_id.value().to_string(),
|
|
||||||
movie_id: movie_id.value().to_string(),
|
|
||||||
user_id: user_id.value().to_string(),
|
|
||||||
rating: rating.value(),
|
|
||||||
watched_at: watched_at.and_utc().timestamp(),
|
|
||||||
}
|
|
||||||
}
|
|
||||||
DomainEvent::MovieDiscovered { movie_id, external_metadata_id } => {
|
|
||||||
DbEventPayload::MovieDiscovered {
|
|
||||||
movie_id: movie_id.value().to_string(),
|
|
||||||
external_metadata_id: external_metadata_id.value().to_owned(),
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
impl TryFrom<DbEventPayload> for DomainEvent {
|
|
||||||
type Error = DomainError;
|
|
||||||
fn try_from(payload: DbEventPayload) -> Result<Self, DomainError> {
|
|
||||||
match payload {
|
|
||||||
DbEventPayload::ReviewLogged { review_id, movie_id, user_id, rating, watched_at } => {
|
|
||||||
Ok(DomainEvent::ReviewLogged {
|
|
||||||
review_id: ReviewId::from_uuid(parse_uuid(&review_id, "review_id")?),
|
|
||||||
movie_id: MovieId::from_uuid(parse_uuid(&movie_id, "movie_id")?),
|
|
||||||
user_id: UserId::from_uuid(parse_uuid(&user_id, "user_id")?),
|
|
||||||
rating: Rating::new(rating)?,
|
|
||||||
watched_at: parse_ts(watched_at)?,
|
|
||||||
})
|
|
||||||
}
|
|
||||||
DbEventPayload::ReviewUpdated { review_id, movie_id, user_id, rating, watched_at } => {
|
|
||||||
Ok(DomainEvent::ReviewUpdated {
|
|
||||||
review_id: ReviewId::from_uuid(parse_uuid(&review_id, "review_id")?),
|
|
||||||
movie_id: MovieId::from_uuid(parse_uuid(&movie_id, "movie_id")?),
|
|
||||||
user_id: UserId::from_uuid(parse_uuid(&user_id, "user_id")?),
|
|
||||||
rating: Rating::new(rating)?,
|
|
||||||
watched_at: parse_ts(watched_at)?,
|
|
||||||
})
|
|
||||||
}
|
|
||||||
DbEventPayload::MovieDiscovered { movie_id, external_metadata_id } => {
|
|
||||||
Ok(DomainEvent::MovieDiscovered {
|
|
||||||
movie_id: MovieId::from_uuid(parse_uuid(&movie_id, "movie_id")?),
|
|
||||||
external_metadata_id: ExternalMetadataId::new(external_metadata_id)?,
|
|
||||||
})
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
#[cfg(test)]
|
|
||||||
mod tests {
|
|
||||||
use super::*;
|
|
||||||
|
|
||||||
fn fixed_dt() -> NaiveDateTime {
|
|
||||||
chrono::DateTime::from_timestamp(1_700_000_000, 0).unwrap().naive_utc()
|
|
||||||
}
|
|
||||||
|
|
||||||
fn review_logged() -> DomainEvent {
|
|
||||||
DomainEvent::ReviewLogged {
|
|
||||||
review_id: ReviewId::from_uuid(Uuid::new_v4()),
|
|
||||||
movie_id: MovieId::from_uuid(Uuid::new_v4()),
|
|
||||||
user_id: UserId::from_uuid(Uuid::new_v4()),
|
|
||||||
rating: Rating::new(4).unwrap(),
|
|
||||||
watched_at: fixed_dt(),
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
fn review_updated() -> DomainEvent {
|
|
||||||
DomainEvent::ReviewUpdated {
|
|
||||||
review_id: ReviewId::from_uuid(Uuid::new_v4()),
|
|
||||||
movie_id: MovieId::from_uuid(Uuid::new_v4()),
|
|
||||||
user_id: UserId::from_uuid(Uuid::new_v4()),
|
|
||||||
rating: Rating::new(3).unwrap(),
|
|
||||||
watched_at: fixed_dt(),
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
fn movie_discovered() -> DomainEvent {
|
|
||||||
DomainEvent::MovieDiscovered {
|
|
||||||
movie_id: MovieId::from_uuid(Uuid::new_v4()),
|
|
||||||
external_metadata_id: ExternalMetadataId::new("tt1234567".into()).unwrap(),
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
fn round_trip(event: DomainEvent) {
|
|
||||||
let payload = DbEventPayload::from(&event);
|
|
||||||
let json = serde_json::to_string(&payload).expect("serialize");
|
|
||||||
let back: DbEventPayload = serde_json::from_str(&json).expect("deserialize");
|
|
||||||
let recovered = DomainEvent::try_from(back).expect("try_from");
|
|
||||||
assert_eq!(DbEventPayload::from(&event), DbEventPayload::from(&recovered));
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn round_trip_review_logged() {
|
|
||||||
round_trip(review_logged());
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn round_trip_review_updated() {
|
|
||||||
round_trip(review_updated());
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn round_trip_movie_discovered() {
|
|
||||||
round_trip(movie_discovered());
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn serialized_format_is_tagged() {
|
|
||||||
let payload = DbEventPayload::from(&movie_discovered());
|
|
||||||
let json = serde_json::to_string(&payload).unwrap();
|
|
||||||
assert!(json.contains(r#""type":"MovieDiscovered""#));
|
|
||||||
assert!(json.contains(r#""data":"#));
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn event_type_strings() {
|
|
||||||
assert_eq!(DbEventPayload::from(&review_logged()).event_type(), "ReviewLogged");
|
|
||||||
assert_eq!(DbEventPayload::from(&review_updated()).event_type(), "ReviewUpdated");
|
|
||||||
assert_eq!(DbEventPayload::from(&movie_discovered()).event_type(), "MovieDiscovered");
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|||||||
@@ -12,6 +12,7 @@ sqlx = { version = "0.8.6", features = [
|
|||||||
"chrono",
|
"chrono",
|
||||||
] }
|
] }
|
||||||
domain = { workspace = true }
|
domain = { workspace = true }
|
||||||
|
anyhow = { workspace = true }
|
||||||
uuid = { workspace = true }
|
uuid = { workspace = true }
|
||||||
chrono = { workspace = true }
|
chrono = { workspace = true }
|
||||||
tracing = { workspace = true }
|
tracing = { workspace = true }
|
||||||
|
|||||||
@@ -767,3 +767,33 @@ impl StatsRepository for PostgresRepository {
|
|||||||
})
|
})
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
pub async fn wire(database_url: &str) -> anyhow::Result<(
|
||||||
|
sqlx::PgPool,
|
||||||
|
std::sync::Arc<dyn domain::ports::MovieRepository>,
|
||||||
|
std::sync::Arc<dyn domain::ports::ReviewRepository>,
|
||||||
|
std::sync::Arc<dyn domain::ports::DiaryRepository>,
|
||||||
|
std::sync::Arc<dyn domain::ports::StatsRepository>,
|
||||||
|
std::sync::Arc<dyn domain::ports::UserRepository>,
|
||||||
|
)> {
|
||||||
|
use anyhow::Context;
|
||||||
|
|
||||||
|
let pool = sqlx::PgPool::connect(database_url)
|
||||||
|
.await
|
||||||
|
.context("Failed to connect to PostgreSQL database")?;
|
||||||
|
|
||||||
|
let repo = std::sync::Arc::new(PostgresRepository::new(pool.clone()));
|
||||||
|
repo.migrate()
|
||||||
|
.await
|
||||||
|
.map_err(|e| anyhow::anyhow!("{e}"))
|
||||||
|
.context("Database migration failed")?;
|
||||||
|
|
||||||
|
Ok((
|
||||||
|
pool.clone(),
|
||||||
|
std::sync::Arc::clone(&repo) as _,
|
||||||
|
std::sync::Arc::clone(&repo) as _,
|
||||||
|
std::sync::Arc::clone(&repo) as _,
|
||||||
|
std::sync::Arc::clone(&repo) as _,
|
||||||
|
std::sync::Arc::new(PostgresUserRepository::new(pool)) as _,
|
||||||
|
))
|
||||||
|
}
|
||||||
|
|||||||
@@ -6,6 +6,7 @@ edition = "2024"
|
|||||||
[dependencies]
|
[dependencies]
|
||||||
sqlx = { version = "0.8.6", features = ["runtime-tokio-rustls", "sqlite", "macros", "chrono"] }
|
sqlx = { version = "0.8.6", features = ["runtime-tokio-rustls", "sqlite", "macros", "chrono"] }
|
||||||
domain = { workspace = true }
|
domain = { workspace = true }
|
||||||
|
event-payload = { workspace = true }
|
||||||
anyhow = { workspace = true }
|
anyhow = { workspace = true }
|
||||||
async-trait = { workspace = true }
|
async-trait = { workspace = true }
|
||||||
serde = { workspace = true }
|
serde = { workspace = true }
|
||||||
|
|||||||
@@ -1,189 +1 @@
|
|||||||
use chrono::NaiveDateTime;
|
pub use event_payload::EventPayload as DbEventPayload;
|
||||||
use domain::{
|
|
||||||
errors::DomainError,
|
|
||||||
events::DomainEvent,
|
|
||||||
value_objects::{ExternalMetadataId, MovieId, Rating, ReviewId, UserId},
|
|
||||||
};
|
|
||||||
use serde::{Deserialize, Serialize};
|
|
||||||
use uuid::Uuid;
|
|
||||||
|
|
||||||
#[derive(Debug, Serialize, Deserialize, PartialEq)]
|
|
||||||
#[serde(tag = "type", content = "data")]
|
|
||||||
pub enum DbEventPayload {
|
|
||||||
ReviewLogged {
|
|
||||||
review_id: String,
|
|
||||||
movie_id: String,
|
|
||||||
user_id: String,
|
|
||||||
rating: u8,
|
|
||||||
watched_at: i64,
|
|
||||||
},
|
|
||||||
ReviewUpdated {
|
|
||||||
review_id: String,
|
|
||||||
movie_id: String,
|
|
||||||
user_id: String,
|
|
||||||
rating: u8,
|
|
||||||
watched_at: i64,
|
|
||||||
},
|
|
||||||
MovieDiscovered {
|
|
||||||
movie_id: String,
|
|
||||||
external_metadata_id: String,
|
|
||||||
},
|
|
||||||
}
|
|
||||||
|
|
||||||
impl DbEventPayload {
|
|
||||||
pub fn event_type(&self) -> &'static str {
|
|
||||||
match self {
|
|
||||||
DbEventPayload::ReviewLogged { .. } => "ReviewLogged",
|
|
||||||
DbEventPayload::ReviewUpdated { .. } => "ReviewUpdated",
|
|
||||||
DbEventPayload::MovieDiscovered { .. } => "MovieDiscovered",
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
fn parse_uuid(s: &str, field: &str) -> Result<Uuid, DomainError> {
|
|
||||||
Uuid::parse_str(s)
|
|
||||||
.map_err(|e| DomainError::InfrastructureError(format!("{field}: {e}")))
|
|
||||||
}
|
|
||||||
|
|
||||||
fn parse_ts(ts: i64) -> Result<NaiveDateTime, DomainError> {
|
|
||||||
chrono::DateTime::from_timestamp(ts, 0)
|
|
||||||
.map(|dt| dt.naive_utc())
|
|
||||||
.ok_or_else(|| DomainError::InfrastructureError(format!("invalid timestamp: {ts}")))
|
|
||||||
}
|
|
||||||
|
|
||||||
impl From<&DomainEvent> for DbEventPayload {
|
|
||||||
fn from(event: &DomainEvent) -> Self {
|
|
||||||
match event {
|
|
||||||
DomainEvent::ReviewLogged { review_id, movie_id, user_id, rating, watched_at } => {
|
|
||||||
DbEventPayload::ReviewLogged {
|
|
||||||
review_id: review_id.value().to_string(),
|
|
||||||
movie_id: movie_id.value().to_string(),
|
|
||||||
user_id: user_id.value().to_string(),
|
|
||||||
rating: rating.value(),
|
|
||||||
watched_at: watched_at.and_utc().timestamp(),
|
|
||||||
}
|
|
||||||
}
|
|
||||||
DomainEvent::ReviewUpdated { review_id, movie_id, user_id, rating, watched_at } => {
|
|
||||||
DbEventPayload::ReviewUpdated {
|
|
||||||
review_id: review_id.value().to_string(),
|
|
||||||
movie_id: movie_id.value().to_string(),
|
|
||||||
user_id: user_id.value().to_string(),
|
|
||||||
rating: rating.value(),
|
|
||||||
watched_at: watched_at.and_utc().timestamp(),
|
|
||||||
}
|
|
||||||
}
|
|
||||||
DomainEvent::MovieDiscovered { movie_id, external_metadata_id } => {
|
|
||||||
DbEventPayload::MovieDiscovered {
|
|
||||||
movie_id: movie_id.value().to_string(),
|
|
||||||
external_metadata_id: external_metadata_id.value().to_owned(),
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
impl TryFrom<DbEventPayload> for DomainEvent {
|
|
||||||
type Error = DomainError;
|
|
||||||
fn try_from(payload: DbEventPayload) -> Result<Self, DomainError> {
|
|
||||||
match payload {
|
|
||||||
DbEventPayload::ReviewLogged { review_id, movie_id, user_id, rating, watched_at } => {
|
|
||||||
Ok(DomainEvent::ReviewLogged {
|
|
||||||
review_id: ReviewId::from_uuid(parse_uuid(&review_id, "review_id")?),
|
|
||||||
movie_id: MovieId::from_uuid(parse_uuid(&movie_id, "movie_id")?),
|
|
||||||
user_id: UserId::from_uuid(parse_uuid(&user_id, "user_id")?),
|
|
||||||
rating: Rating::new(rating)?,
|
|
||||||
watched_at: parse_ts(watched_at)?,
|
|
||||||
})
|
|
||||||
}
|
|
||||||
DbEventPayload::ReviewUpdated { review_id, movie_id, user_id, rating, watched_at } => {
|
|
||||||
Ok(DomainEvent::ReviewUpdated {
|
|
||||||
review_id: ReviewId::from_uuid(parse_uuid(&review_id, "review_id")?),
|
|
||||||
movie_id: MovieId::from_uuid(parse_uuid(&movie_id, "movie_id")?),
|
|
||||||
user_id: UserId::from_uuid(parse_uuid(&user_id, "user_id")?),
|
|
||||||
rating: Rating::new(rating)?,
|
|
||||||
watched_at: parse_ts(watched_at)?,
|
|
||||||
})
|
|
||||||
}
|
|
||||||
DbEventPayload::MovieDiscovered { movie_id, external_metadata_id } => {
|
|
||||||
Ok(DomainEvent::MovieDiscovered {
|
|
||||||
movie_id: MovieId::from_uuid(parse_uuid(&movie_id, "movie_id")?),
|
|
||||||
external_metadata_id: ExternalMetadataId::new(external_metadata_id)?,
|
|
||||||
})
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
#[cfg(test)]
|
|
||||||
mod tests {
|
|
||||||
use super::*;
|
|
||||||
|
|
||||||
fn fixed_dt() -> NaiveDateTime {
|
|
||||||
chrono::DateTime::from_timestamp(1_700_000_000, 0).unwrap().naive_utc()
|
|
||||||
}
|
|
||||||
|
|
||||||
fn review_logged() -> DomainEvent {
|
|
||||||
DomainEvent::ReviewLogged {
|
|
||||||
review_id: ReviewId::from_uuid(Uuid::new_v4()),
|
|
||||||
movie_id: MovieId::from_uuid(Uuid::new_v4()),
|
|
||||||
user_id: UserId::from_uuid(Uuid::new_v4()),
|
|
||||||
rating: Rating::new(4).unwrap(),
|
|
||||||
watched_at: fixed_dt(),
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
fn review_updated() -> DomainEvent {
|
|
||||||
DomainEvent::ReviewUpdated {
|
|
||||||
review_id: ReviewId::from_uuid(Uuid::new_v4()),
|
|
||||||
movie_id: MovieId::from_uuid(Uuid::new_v4()),
|
|
||||||
user_id: UserId::from_uuid(Uuid::new_v4()),
|
|
||||||
rating: Rating::new(3).unwrap(),
|
|
||||||
watched_at: fixed_dt(),
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
fn movie_discovered() -> DomainEvent {
|
|
||||||
DomainEvent::MovieDiscovered {
|
|
||||||
movie_id: MovieId::from_uuid(Uuid::new_v4()),
|
|
||||||
external_metadata_id: ExternalMetadataId::new("tt1234567".into()).unwrap(),
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
fn round_trip(event: DomainEvent) {
|
|
||||||
let payload = DbEventPayload::from(&event);
|
|
||||||
let json = serde_json::to_string(&payload).expect("serialize");
|
|
||||||
let back: DbEventPayload = serde_json::from_str(&json).expect("deserialize");
|
|
||||||
let recovered = DomainEvent::try_from(back).expect("try_from");
|
|
||||||
assert_eq!(DbEventPayload::from(&event), DbEventPayload::from(&recovered));
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn round_trip_review_logged() {
|
|
||||||
round_trip(review_logged());
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn round_trip_review_updated() {
|
|
||||||
round_trip(review_updated());
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn round_trip_movie_discovered() {
|
|
||||||
round_trip(movie_discovered());
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn serialized_format_is_tagged() {
|
|
||||||
let payload = DbEventPayload::from(&movie_discovered());
|
|
||||||
let json = serde_json::to_string(&payload).unwrap();
|
|
||||||
assert!(json.contains(r#""type":"MovieDiscovered""#));
|
|
||||||
assert!(json.contains(r#""data":"#));
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn event_type_strings() {
|
|
||||||
assert_eq!(DbEventPayload::from(&review_logged()).event_type(), "ReviewLogged");
|
|
||||||
assert_eq!(DbEventPayload::from(&review_updated()).event_type(), "ReviewUpdated");
|
|
||||||
assert_eq!(DbEventPayload::from(&movie_discovered()).event_type(), "MovieDiscovered");
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|||||||
@@ -759,6 +759,43 @@ impl StatsRepository for SqliteMovieRepository {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
pub async fn wire(database_url: &str) -> anyhow::Result<(
|
||||||
|
sqlx::SqlitePool,
|
||||||
|
std::sync::Arc<dyn domain::ports::MovieRepository>,
|
||||||
|
std::sync::Arc<dyn domain::ports::ReviewRepository>,
|
||||||
|
std::sync::Arc<dyn domain::ports::DiaryRepository>,
|
||||||
|
std::sync::Arc<dyn domain::ports::StatsRepository>,
|
||||||
|
std::sync::Arc<dyn domain::ports::UserRepository>,
|
||||||
|
)> {
|
||||||
|
use std::str::FromStr;
|
||||||
|
use anyhow::Context;
|
||||||
|
use sqlx::sqlite::SqliteConnectOptions;
|
||||||
|
|
||||||
|
let opts = SqliteConnectOptions::from_str(database_url)
|
||||||
|
.context("Invalid DATABASE_URL")?
|
||||||
|
.create_if_missing(true)
|
||||||
|
.journal_mode(sqlx::sqlite::SqliteJournalMode::Wal)
|
||||||
|
.busy_timeout(std::time::Duration::from_secs(5));
|
||||||
|
let pool = sqlx::SqlitePool::connect_with(opts)
|
||||||
|
.await
|
||||||
|
.context("Failed to connect to SQLite database")?;
|
||||||
|
|
||||||
|
let repo = std::sync::Arc::new(SqliteMovieRepository::new(pool.clone()));
|
||||||
|
repo.migrate()
|
||||||
|
.await
|
||||||
|
.map_err(|e| anyhow::anyhow!("{e}"))
|
||||||
|
.context("Database migration failed")?;
|
||||||
|
|
||||||
|
Ok((
|
||||||
|
pool.clone(),
|
||||||
|
std::sync::Arc::clone(&repo) as _,
|
||||||
|
std::sync::Arc::clone(&repo) as _,
|
||||||
|
std::sync::Arc::clone(&repo) as _,
|
||||||
|
std::sync::Arc::clone(&repo) as _,
|
||||||
|
std::sync::Arc::new(SqliteUserRepository::new(pool)) as _,
|
||||||
|
))
|
||||||
|
}
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
mod feed_filter_tests {
|
mod feed_filter_tests {
|
||||||
use super::*;
|
use super::*;
|
||||||
|
|||||||
@@ -49,7 +49,6 @@ metadata = { workspace = true }
|
|||||||
poster-fetcher = { workspace = true }
|
poster-fetcher = { workspace = true }
|
||||||
poster-storage = { workspace = true }
|
poster-storage = { workspace = true }
|
||||||
template-askama = { workspace = true }
|
template-askama = { workspace = true }
|
||||||
event-publisher = { workspace = true }
|
|
||||||
nats = { workspace = true, optional = true }
|
nats = { workspace = true, optional = true }
|
||||||
rss = { workspace = true }
|
rss = { workspace = true }
|
||||||
export = { workspace = true }
|
export = { workspace = true }
|
||||||
|
|||||||
@@ -1,18 +1,12 @@
|
|||||||
use std::sync::Arc;
|
use std::sync::Arc;
|
||||||
|
|
||||||
use anyhow::Context;
|
use anyhow::Context;
|
||||||
use std::str::FromStr;
|
|
||||||
|
|
||||||
use tokio::net::TcpListener;
|
use tokio::net::TcpListener;
|
||||||
use tracing_subscriber::{layer::SubscriberExt, util::SubscriberInitExt};
|
use tracing_subscriber::{layer::SubscriberExt, util::SubscriberInitExt};
|
||||||
|
|
||||||
#[cfg(feature = "sqlite")]
|
|
||||||
use sqlite::{SqliteMovieRepository, SqliteUserRepository};
|
|
||||||
#[cfg(feature = "sqlite-federation")]
|
#[cfg(feature = "sqlite-federation")]
|
||||||
use sqlite_federation::SqliteFederationRepository;
|
use sqlite_federation::SqliteFederationRepository;
|
||||||
|
|
||||||
#[cfg(feature = "postgres")]
|
|
||||||
use postgres::{PostgresRepository, PostgresUserRepository};
|
|
||||||
#[cfg(feature = "postgres-federation")]
|
#[cfg(feature = "postgres-federation")]
|
||||||
use postgres_federation::PostgresFederationRepository;
|
use postgres_federation::PostgresFederationRepository;
|
||||||
|
|
||||||
@@ -36,9 +30,8 @@ use presentation::{openapi::ApiDoc, routes, state::AppState};
|
|||||||
use utoipa::OpenApi as _;
|
use utoipa::OpenApi as _;
|
||||||
|
|
||||||
use domain::ports::{
|
use domain::ports::{
|
||||||
AuthService, DiaryExporter, DiaryRepository, EventPublisher, MetadataClient,
|
AuthService, DiaryExporter, EventPublisher, MetadataClient,
|
||||||
MovieRepository, PasswordHasher, PosterFetcherClient, PosterStorage, ReviewRepository,
|
PasswordHasher, PosterFetcherClient, PosterStorage,
|
||||||
StatsRepository, UserRepository,
|
|
||||||
};
|
};
|
||||||
|
|
||||||
#[cfg(not(any(feature = "sqlite", feature = "postgres")))]
|
#[cfg(not(any(feature = "sqlite", feature = "postgres")))]
|
||||||
@@ -91,12 +84,12 @@ async fn wire_dependencies() -> anyhow::Result<(AppState, axum::Router)> {
|
|||||||
match backend.as_str() {
|
match backend.as_str() {
|
||||||
#[cfg(feature = "postgres")]
|
#[cfg(feature = "postgres")]
|
||||||
"postgres" => {
|
"postgres" => {
|
||||||
let (pool, m, r, d, s, u) = wire_postgres(&database_url).await?;
|
let (pool, m, r, d, s, u) = postgres::wire(&database_url).await?;
|
||||||
(m, r, d, s, u, DbPool::Postgres(pool))
|
(m, r, d, s, u, DbPool::Postgres(pool))
|
||||||
}
|
}
|
||||||
#[cfg(feature = "sqlite")]
|
#[cfg(feature = "sqlite")]
|
||||||
_ => {
|
_ => {
|
||||||
let (pool, m, r, d, s, u) = wire_sqlite(&database_url).await?;
|
let (pool, m, r, d, s, u) = sqlite::wire(&database_url).await?;
|
||||||
(m, r, d, s, u, DbPool::Sqlite(pool))
|
(m, r, d, s, u, DbPool::Sqlite(pool))
|
||||||
}
|
}
|
||||||
#[cfg(not(feature = "sqlite"))]
|
#[cfg(not(feature = "sqlite"))]
|
||||||
@@ -230,73 +223,6 @@ async fn wire_dependencies() -> anyhow::Result<(AppState, axum::Router)> {
|
|||||||
Ok((state, ap_router))
|
Ok((state, ap_router))
|
||||||
}
|
}
|
||||||
|
|
||||||
#[cfg(feature = "sqlite")]
|
|
||||||
async fn wire_sqlite(database_url: &str) -> anyhow::Result<(
|
|
||||||
sqlx::SqlitePool,
|
|
||||||
Arc<dyn MovieRepository>,
|
|
||||||
Arc<dyn ReviewRepository>,
|
|
||||||
Arc<dyn DiaryRepository>,
|
|
||||||
Arc<dyn StatsRepository>,
|
|
||||||
Arc<dyn UserRepository>,
|
|
||||||
)> {
|
|
||||||
use sqlx::sqlite::SqliteConnectOptions;
|
|
||||||
|
|
||||||
let opts = SqliteConnectOptions::from_str(database_url)
|
|
||||||
.context("Invalid DATABASE_URL")?
|
|
||||||
.create_if_missing(true)
|
|
||||||
.journal_mode(sqlx::sqlite::SqliteJournalMode::Wal)
|
|
||||||
.busy_timeout(std::time::Duration::from_secs(5));
|
|
||||||
let pool = sqlx::SqlitePool::connect_with(opts)
|
|
||||||
.await
|
|
||||||
.context("Failed to connect to SQLite database")?;
|
|
||||||
|
|
||||||
let sqlite_repo = Arc::new(SqliteMovieRepository::new(pool.clone()));
|
|
||||||
sqlite_repo
|
|
||||||
.migrate()
|
|
||||||
.await
|
|
||||||
.map_err(|e| anyhow::anyhow!("{}", e))
|
|
||||||
.context("Database migration failed")?;
|
|
||||||
|
|
||||||
let movie_repository: Arc<dyn MovieRepository> = Arc::clone(&sqlite_repo) as _;
|
|
||||||
let review_repository: Arc<dyn ReviewRepository> = Arc::clone(&sqlite_repo) as _;
|
|
||||||
let diary_repository: Arc<dyn DiaryRepository> = Arc::clone(&sqlite_repo) as _;
|
|
||||||
let stats_repository: Arc<dyn StatsRepository> = Arc::clone(&sqlite_repo) as _;
|
|
||||||
let user_repository: Arc<dyn UserRepository> =
|
|
||||||
Arc::new(SqliteUserRepository::new(pool.clone()));
|
|
||||||
|
|
||||||
Ok((pool, movie_repository, review_repository, diary_repository, stats_repository, user_repository))
|
|
||||||
}
|
|
||||||
|
|
||||||
#[cfg(feature = "postgres")]
|
|
||||||
async fn wire_postgres(database_url: &str) -> anyhow::Result<(
|
|
||||||
sqlx::PgPool,
|
|
||||||
Arc<dyn MovieRepository>,
|
|
||||||
Arc<dyn ReviewRepository>,
|
|
||||||
Arc<dyn DiaryRepository>,
|
|
||||||
Arc<dyn StatsRepository>,
|
|
||||||
Arc<dyn UserRepository>,
|
|
||||||
)> {
|
|
||||||
let pool = sqlx::PgPool::connect(database_url)
|
|
||||||
.await
|
|
||||||
.context("Failed to connect to PostgreSQL database")?;
|
|
||||||
|
|
||||||
let pg_repo = Arc::new(PostgresRepository::new(pool.clone()));
|
|
||||||
pg_repo
|
|
||||||
.migrate()
|
|
||||||
.await
|
|
||||||
.map_err(|e| anyhow::anyhow!("{}", e))
|
|
||||||
.context("Database migration failed")?;
|
|
||||||
|
|
||||||
let movie_repository: Arc<dyn MovieRepository> = Arc::clone(&pg_repo) as _;
|
|
||||||
let review_repository: Arc<dyn ReviewRepository> = Arc::clone(&pg_repo) as _;
|
|
||||||
let diary_repository: Arc<dyn DiaryRepository> = Arc::clone(&pg_repo) as _;
|
|
||||||
let stats_repository: Arc<dyn StatsRepository> = Arc::clone(&pg_repo) as _;
|
|
||||||
let user_repository: Arc<dyn UserRepository> =
|
|
||||||
Arc::new(PostgresUserRepository::new(pool.clone()));
|
|
||||||
|
|
||||||
Ok((pool, movie_repository, review_repository, diary_repository, stats_repository, user_repository))
|
|
||||||
}
|
|
||||||
|
|
||||||
enum DbPool {
|
enum DbPool {
|
||||||
#[cfg(feature = "sqlite")]
|
#[cfg(feature = "sqlite")]
|
||||||
Sqlite(sqlx::SqlitePool),
|
Sqlite(sqlx::SqlitePool),
|
||||||
|
|||||||
@@ -15,7 +15,6 @@ postgres-federation = ["postgres", "dep:postgres-federation", "dep:activitypub",
|
|||||||
[dependencies]
|
[dependencies]
|
||||||
domain = { workspace = true }
|
domain = { workspace = true }
|
||||||
application = { workspace = true }
|
application = { workspace = true }
|
||||||
event-publisher = { workspace = true }
|
|
||||||
tokio = { workspace = true }
|
tokio = { workspace = true }
|
||||||
anyhow = { workspace = true }
|
anyhow = { workspace = true }
|
||||||
thiserror = { workspace = true }
|
thiserror = { workspace = true }
|
||||||
|
|||||||
@@ -1,5 +1,4 @@
|
|||||||
use std::sync::Arc;
|
use std::sync::Arc;
|
||||||
use std::str::FromStr;
|
|
||||||
|
|
||||||
use anyhow::Context;
|
use anyhow::Context;
|
||||||
use application::{config::AppConfig, context::AppContext, event_handlers::PosterSyncHandler, worker::WorkerService};
|
use application::{config::AppConfig, context::AppContext, event_handlers::PosterSyncHandler, worker::WorkerService};
|
||||||
@@ -10,12 +9,6 @@ use poster_fetcher::{PosterFetcherConfig, ReqwestPosterFetcher};
|
|||||||
use poster_storage::{PosterStorageAdapter, StorageConfig};
|
use poster_storage::{PosterStorageAdapter, StorageConfig};
|
||||||
use tracing_subscriber::{layer::SubscriberExt, util::SubscriberInitExt};
|
use tracing_subscriber::{layer::SubscriberExt, util::SubscriberInitExt};
|
||||||
|
|
||||||
#[cfg(feature = "sqlite")]
|
|
||||||
use sqlite::{SqliteMovieRepository, SqliteUserRepository};
|
|
||||||
|
|
||||||
#[cfg(feature = "postgres")]
|
|
||||||
use postgres::{PostgresRepository, PostgresUserRepository};
|
|
||||||
|
|
||||||
#[cfg(feature = "sqlite-federation")]
|
#[cfg(feature = "sqlite-federation")]
|
||||||
use sqlite_federation::SqliteFederationRepository;
|
use sqlite_federation::SqliteFederationRepository;
|
||||||
#[cfg(feature = "postgres-federation")]
|
#[cfg(feature = "postgres-federation")]
|
||||||
@@ -27,9 +20,8 @@ use activitypub::{
|
|||||||
};
|
};
|
||||||
|
|
||||||
use domain::ports::{
|
use domain::ports::{
|
||||||
AuthService, DiaryExporter, DiaryRepository, EventHandler, MetadataClient, MovieRepository,
|
AuthService, DiaryExporter, EventHandler, MetadataClient,
|
||||||
PasswordHasher, PosterFetcherClient, PosterStorage, ReviewRepository, StatsRepository,
|
PasswordHasher, PosterFetcherClient, PosterStorage,
|
||||||
UserRepository,
|
|
||||||
};
|
};
|
||||||
|
|
||||||
#[cfg(not(any(feature = "sqlite", feature = "postgres")))]
|
#[cfg(not(any(feature = "sqlite", feature = "postgres")))]
|
||||||
@@ -65,12 +57,12 @@ async fn main() -> anyhow::Result<()> {
|
|||||||
match backend.as_str() {
|
match backend.as_str() {
|
||||||
#[cfg(feature = "postgres")]
|
#[cfg(feature = "postgres")]
|
||||||
"postgres" => {
|
"postgres" => {
|
||||||
let (pool, m, r, d, s, u) = wire_postgres(&database_url).await?;
|
let (pool, m, r, d, s, u) = postgres::wire(&database_url).await?;
|
||||||
(m, r, d, s, u, DbPool::Postgres(pool))
|
(m, r, d, s, u, DbPool::Postgres(pool))
|
||||||
}
|
}
|
||||||
#[cfg(feature = "sqlite")]
|
#[cfg(feature = "sqlite")]
|
||||||
_ => {
|
_ => {
|
||||||
let (pool, m, r, d, s, u) = wire_sqlite(&database_url).await?;
|
let (pool, m, r, d, s, u) = sqlite::wire(&database_url).await?;
|
||||||
(m, r, d, s, u, DbPool::Sqlite(pool))
|
(m, r, d, s, u, DbPool::Sqlite(pool))
|
||||||
}
|
}
|
||||||
#[cfg(not(feature = "sqlite"))]
|
#[cfg(not(feature = "sqlite"))]
|
||||||
@@ -125,50 +117,57 @@ async fn main() -> anyhow::Result<()> {
|
|||||||
config: app_config,
|
config: app_config,
|
||||||
};
|
};
|
||||||
|
|
||||||
let mut handlers: Vec<Arc<dyn EventHandler>> = vec![Arc::new(PosterSyncHandler::new(ctx, 3))];
|
let handlers: Vec<Arc<dyn EventHandler>> = {
|
||||||
|
let poster = Arc::new(PosterSyncHandler::new(ctx, 3)) as Arc<dyn EventHandler>;
|
||||||
|
|
||||||
#[cfg(feature = "federation")]
|
#[cfg(not(feature = "federation"))]
|
||||||
{
|
{ vec![poster] }
|
||||||
let (federation_repo, review_store): (
|
|
||||||
Arc<dyn activitypub::FederationRepository>,
|
|
||||||
Arc<dyn activitypub::RemoteReviewRepository>,
|
|
||||||
) = match &db_pool {
|
|
||||||
#[cfg(feature = "sqlite-federation")]
|
|
||||||
DbPool::Sqlite(pool) => {
|
|
||||||
let fed = Arc::new(SqliteFederationRepository::new(pool.clone()));
|
|
||||||
(Arc::clone(&fed) as _, fed as _)
|
|
||||||
}
|
|
||||||
#[cfg(feature = "postgres-federation")]
|
|
||||||
DbPool::Postgres(pool) => {
|
|
||||||
let fed = Arc::new(PostgresFederationRepository::new(pool.clone()));
|
|
||||||
(Arc::clone(&fed) as _, fed as _)
|
|
||||||
}
|
|
||||||
};
|
|
||||||
|
|
||||||
let ap_service = Arc::new(
|
#[cfg(feature = "federation")]
|
||||||
ActivityPubService::new(
|
{
|
||||||
federation_repo,
|
let (federation_repo, review_store): (
|
||||||
Arc::new(DomainUserRepoAdapter(fed_user_repo)),
|
Arc<dyn activitypub::FederationRepository>,
|
||||||
Arc::new(ReviewObjectHandler {
|
Arc<dyn activitypub::RemoteReviewRepository>,
|
||||||
movie_repository: Arc::clone(&fed_movie_repo),
|
) = match &db_pool {
|
||||||
diary_repository: fed_diary_repo,
|
#[cfg(feature = "sqlite-federation")]
|
||||||
review_store,
|
DbPool::Sqlite(pool) => {
|
||||||
base_url: base_url.clone(),
|
let fed = Arc::new(SqliteFederationRepository::new(pool.clone()));
|
||||||
}),
|
(Arc::clone(&fed) as _, fed as _)
|
||||||
base_url.clone(),
|
}
|
||||||
cfg!(debug_assertions),
|
#[cfg(feature = "postgres-federation")]
|
||||||
)
|
DbPool::Postgres(pool) => {
|
||||||
.await?,
|
let fed = Arc::new(PostgresFederationRepository::new(pool.clone()));
|
||||||
);
|
(Arc::clone(&fed) as _, fed as _)
|
||||||
|
}
|
||||||
|
};
|
||||||
|
|
||||||
handlers.push(Arc::new(ActivityPubEventHandler::new(
|
let ap_service = Arc::new(
|
||||||
ap_service,
|
ActivityPubService::new(
|
||||||
fed_movie_repo,
|
federation_repo,
|
||||||
fed_review_repo,
|
Arc::new(DomainUserRepoAdapter(fed_user_repo)),
|
||||||
base_url,
|
Arc::new(ReviewObjectHandler {
|
||||||
)));
|
movie_repository: Arc::clone(&fed_movie_repo),
|
||||||
tracing::info!("federation event handler registered");
|
diary_repository: fed_diary_repo,
|
||||||
}
|
review_store,
|
||||||
|
base_url: base_url.clone(),
|
||||||
|
}),
|
||||||
|
base_url.clone(),
|
||||||
|
cfg!(debug_assertions),
|
||||||
|
)
|
||||||
|
.await?,
|
||||||
|
);
|
||||||
|
|
||||||
|
let ap = Arc::new(ActivityPubEventHandler::new(
|
||||||
|
ap_service,
|
||||||
|
fed_movie_repo,
|
||||||
|
fed_review_repo,
|
||||||
|
base_url,
|
||||||
|
)) as Arc<dyn EventHandler>;
|
||||||
|
|
||||||
|
tracing::info!("federation event handler registered");
|
||||||
|
vec![poster, ap]
|
||||||
|
}
|
||||||
|
};
|
||||||
|
|
||||||
let worker = WorkerService::new(consumer_arc, handlers);
|
let worker = WorkerService::new(consumer_arc, handlers);
|
||||||
|
|
||||||
@@ -218,69 +217,3 @@ fn init_tracing() {
|
|||||||
.init();
|
.init();
|
||||||
}
|
}
|
||||||
|
|
||||||
#[cfg(feature = "sqlite")]
|
|
||||||
async fn wire_sqlite(database_url: &str) -> anyhow::Result<(
|
|
||||||
sqlx::SqlitePool,
|
|
||||||
Arc<dyn MovieRepository>,
|
|
||||||
Arc<dyn ReviewRepository>,
|
|
||||||
Arc<dyn DiaryRepository>,
|
|
||||||
Arc<dyn StatsRepository>,
|
|
||||||
Arc<dyn UserRepository>,
|
|
||||||
)> {
|
|
||||||
use sqlx::sqlite::SqliteConnectOptions;
|
|
||||||
|
|
||||||
let opts = SqliteConnectOptions::from_str(database_url)
|
|
||||||
.context("Invalid DATABASE_URL")?
|
|
||||||
.create_if_missing(true)
|
|
||||||
.journal_mode(sqlx::sqlite::SqliteJournalMode::Wal)
|
|
||||||
.busy_timeout(std::time::Duration::from_secs(5));
|
|
||||||
let pool = sqlx::SqlitePool::connect_with(opts)
|
|
||||||
.await
|
|
||||||
.context("Failed to connect to SQLite database")?;
|
|
||||||
|
|
||||||
let sqlite_repo = Arc::new(SqliteMovieRepository::new(pool.clone()));
|
|
||||||
sqlite_repo
|
|
||||||
.migrate()
|
|
||||||
.await
|
|
||||||
.map_err(|e| anyhow::anyhow!("{}", e))
|
|
||||||
.context("Database migration failed")?;
|
|
||||||
|
|
||||||
let movie_repository: Arc<dyn MovieRepository> = Arc::clone(&sqlite_repo) as _;
|
|
||||||
let review_repository: Arc<dyn ReviewRepository> = Arc::clone(&sqlite_repo) as _;
|
|
||||||
let diary_repository: Arc<dyn DiaryRepository> = Arc::clone(&sqlite_repo) as _;
|
|
||||||
let stats_repository: Arc<dyn StatsRepository> = Arc::clone(&sqlite_repo) as _;
|
|
||||||
let user_repository: Arc<dyn UserRepository> =
|
|
||||||
Arc::new(SqliteUserRepository::new(pool.clone()));
|
|
||||||
|
|
||||||
Ok((pool, movie_repository, review_repository, diary_repository, stats_repository, user_repository))
|
|
||||||
}
|
|
||||||
|
|
||||||
#[cfg(feature = "postgres")]
|
|
||||||
async fn wire_postgres(database_url: &str) -> anyhow::Result<(
|
|
||||||
sqlx::PgPool,
|
|
||||||
Arc<dyn MovieRepository>,
|
|
||||||
Arc<dyn ReviewRepository>,
|
|
||||||
Arc<dyn DiaryRepository>,
|
|
||||||
Arc<dyn StatsRepository>,
|
|
||||||
Arc<dyn UserRepository>,
|
|
||||||
)> {
|
|
||||||
let pool = sqlx::PgPool::connect(database_url)
|
|
||||||
.await
|
|
||||||
.context("Failed to connect to PostgreSQL database")?;
|
|
||||||
|
|
||||||
let pg_repo = Arc::new(PostgresRepository::new(pool.clone()));
|
|
||||||
pg_repo
|
|
||||||
.migrate()
|
|
||||||
.await
|
|
||||||
.map_err(|e| anyhow::anyhow!("{}", e))
|
|
||||||
.context("Database migration failed")?;
|
|
||||||
|
|
||||||
let movie_repository: Arc<dyn MovieRepository> = Arc::clone(&pg_repo) as _;
|
|
||||||
let review_repository: Arc<dyn ReviewRepository> = Arc::clone(&pg_repo) as _;
|
|
||||||
let diary_repository: Arc<dyn DiaryRepository> = Arc::clone(&pg_repo) as _;
|
|
||||||
let stats_repository: Arc<dyn StatsRepository> = Arc::clone(&pg_repo) as _;
|
|
||||||
let user_repository: Arc<dyn UserRepository> =
|
|
||||||
Arc::new(PostgresUserRepository::new(pool.clone()));
|
|
||||||
|
|
||||||
Ok((pool, movie_repository, review_repository, diary_repository, stats_repository, user_repository))
|
|
||||||
}
|
|
||||||
|
|||||||
Reference in New Issue
Block a user