diff --git a/Cargo.lock b/Cargo.lock index 5857d08..fb6ce6a 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -15,6 +15,22 @@ dependencies = [ "uuid", ] +[[package]] +name = "adapter-postgres" +version = "0.1.0" +dependencies = [ + "adapter-common", + "async-trait", + "chrono", + "domain", + "infra-wiring", + "serde", + "serde_json", + "sqlx", + "tracing", + "uuid", +] + [[package]] name = "adapter-sqlite" version = "0.1.0" diff --git a/Cargo.toml b/Cargo.toml index 0c4a311..04b068d 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,5 +1,5 @@ [workspace] -members = ["crates/domain", "crates/application", "crates/api-types", "crates/infra-wiring", "crates/adapters/adapter-common", "crates/adapters/sqlite"] +members = ["crates/domain", "crates/application", "crates/api-types", "crates/infra-wiring", "crates/adapters/adapter-common", "crates/adapters/sqlite", "crates/adapters/postgres"] exclude = ["k-tv-backend", "k-tv-frontend"] resolver = "2" diff --git a/crates/adapters/postgres/Cargo.toml b/crates/adapters/postgres/Cargo.toml new file mode 100644 index 0000000..3c2aa54 --- /dev/null +++ b/crates/adapters/postgres/Cargo.toml @@ -0,0 +1,16 @@ +[package] +name = "adapter-postgres" +version = "0.1.0" +edition = "2024" + +[dependencies] +domain = { workspace = true } +adapter-common = { workspace = true } +infra-wiring = { workspace = true, features = ["postgres"] } +async-trait = { workspace = true } +sqlx = { workspace = true, features = ["postgres"] } +chrono = { workspace = true } +uuid = { workspace = true } +serde = { workspace = true } +serde_json = { workspace = true } +tracing = { workspace = true } diff --git a/crates/adapters/postgres/src/activity.rs b/crates/adapters/postgres/src/activity.rs new file mode 100644 index 0000000..e3fee6e --- /dev/null +++ b/crates/adapters/postgres/src/activity.rs @@ -0,0 +1,84 @@ +//! PostgreSQL adapter for activity log (ActivityLogCommand + ActivityLogQuery). + +use async_trait::async_trait; +use chrono::Utc; +use sqlx::PgPool; +use uuid::Uuid; + +use adapter_common::{map_sqlx_error, parse_dt, parse_uuid}; +use domain::{ + ports::activity::{ActivityLogCommand, ActivityLogQuery}, + ActivityEvent, ChannelId, DomainResult, +}; + +pub struct PgActivityLog { + pool: PgPool, +} + +impl PgActivityLog { + pub fn new(pool: PgPool) -> Self { + Self { pool } + } +} + +#[async_trait] +impl ActivityLogCommand for PgActivityLog { + async fn log( + &self, + event_type: &str, + detail: &str, + channel_id: Option, + ) -> DomainResult<()> { + let id = Uuid::new_v4().to_string(); + let timestamp = Utc::now().to_rfc3339(); + let channel_id_str = channel_id.map(|id| id.value().to_string()); + + sqlx::query( + "INSERT INTO activity_log (id, timestamp, event_type, detail, channel_id) VALUES ($1, $2, $3, $4, $5)", + ) + .bind(&id) + .bind(×tamp) + .bind(event_type) + .bind(detail) + .bind(&channel_id_str) + .execute(&self.pool) + .await + .map_err(map_sqlx_error)?; + + Ok(()) + } +} + +#[async_trait] +impl ActivityLogQuery for PgActivityLog { + async fn recent(&self, limit: u32) -> DomainResult> { + let rows: Vec<(String, String, String, String, Option)> = sqlx::query_as( + "SELECT id, timestamp, event_type, detail, channel_id FROM activity_log ORDER BY timestamp DESC LIMIT $1", + ) + .bind(limit as i64) + .fetch_all(&self.pool) + .await + .map_err(map_sqlx_error)?; + + let mut events = Vec::with_capacity(rows.len()); + for (id_str, ts_str, event_type, detail, channel_id_str) in rows { + // Silently skip rows with bad UUIDs/timestamps (mirrors old behaviour) + let Ok(id) = parse_uuid(&id_str, "activity id") else { + continue; + }; + let Ok(timestamp) = parse_dt(&ts_str) else { + continue; + }; + let channel_id = channel_id_str.and_then(|s| Uuid::parse_str(&s).ok()); + events.push(ActivityEvent::from_persistence( + id, + timestamp, + event_type, + detail, + channel_id, + )); + } + + Ok(events) + } +} diff --git a/crates/adapters/postgres/src/channel.rs b/crates/adapters/postgres/src/channel.rs new file mode 100644 index 0000000..08e4789 --- /dev/null +++ b/crates/adapters/postgres/src/channel.rs @@ -0,0 +1,353 @@ +//! PostgreSQL adapter for channel persistence (ChannelCommand + ChannelQuery). + +use async_trait::async_trait; +use chrono::{DateTime, Utc}; +use sqlx::{PgPool, Row}; +use uuid::Uuid; + +use adapter_common::{ + map_sqlx_error, parse_dt, parse_enum_or_default, parse_recycle_policy, parse_schedule_config, + parse_uuid, +}; +use domain::{ + ports::channel::{ChannelCommand, ChannelQuery}, + AccessMode, Channel, ChannelConfigSnapshot, ChannelId, DomainError, DomainResult, LogoPosition, + ScheduleConfig, UserId, +}; + +pub struct PgChannelRepository { + pool: PgPool, +} + +impl PgChannelRepository { + pub fn new(pool: PgPool) -> Self { + Self { pool } + } +} + +// -- Row type ---------------------------------------------------------------- + +const SELECT_COLS: &str = "id, owner_id, name, description, timezone, schedule_config, recycle_policy, auto_schedule, access_mode, access_password_hash, logo, logo_position, logo_opacity, webhook_url, webhook_poll_interval_secs, webhook_body_template, webhook_headers, created_at, updated_at"; + +#[derive(Debug, sqlx::FromRow)] +struct ChannelRow { + id: String, + owner_id: String, + name: String, + description: Option, + timezone: String, + schedule_config: String, + recycle_policy: String, + auto_schedule: bool, + access_mode: String, + access_password_hash: Option, + logo: Option, + logo_position: String, + logo_opacity: f32, + webhook_url: Option, + webhook_poll_interval_secs: i64, + webhook_body_template: Option, + webhook_headers: Option, + created_at: String, + updated_at: String, +} + +impl ChannelRow { + fn into_channel(self) -> DomainResult { + let id = ChannelId::from_uuid(parse_uuid(&self.id, "channel id")?); + let owner_id = UserId::from_uuid(parse_uuid(&self.owner_id, "owner id")?); + let schedule_config = parse_schedule_config(&self.schedule_config)?; + let recycle_policy = parse_recycle_policy(&self.recycle_policy)?; + let access_mode: AccessMode = parse_enum_or_default(self.access_mode); + let logo_position: LogoPosition = parse_enum_or_default(self.logo_position); + + Ok(Channel::from_persistence( + id, + owner_id, + self.name, + self.description, + self.timezone, + schedule_config, + recycle_policy, + self.auto_schedule, + access_mode, + self.access_password_hash, + self.logo, + logo_position, + self.logo_opacity, + self.webhook_url, + self.webhook_poll_interval_secs as u32, + self.webhook_body_template, + self.webhook_headers, + parse_dt(&self.created_at)?, + parse_dt(&self.updated_at)?, + )) + } +} + +// -- Helpers ------------------------------------------------------------------ + +fn serialize_enum_as_string(v: &T, fallback: &str) -> String { + serde_json::to_value(v) + .ok() + .and_then(|v| v.as_str().map(str::to_owned)) + .unwrap_or_else(|| fallback.to_owned()) +} + +fn map_snapshot_row( + row: &sqlx::postgres::PgRow, + channel_id: ChannelId, +) -> DomainResult { + let id_str: String = row.get("id"); + let id = parse_uuid(&id_str, "snapshot id")?; + let config_json: String = row.get("config_json"); + let config = parse_schedule_config(&config_json)?; + let version_num: i64 = row.get("version_num"); + let label: Option = row.get("label"); + let created_at_str: String = row.get("created_at"); + let created_at: DateTime = parse_dt(&created_at_str)?; + + Ok(ChannelConfigSnapshot::from_persistence( + id, + channel_id, + config, + version_num, + label, + created_at, + )) +} + +// -- Command ----------------------------------------------------------------- + +#[async_trait] +impl ChannelCommand for PgChannelRepository { + async fn save(&self, channel: &Channel) -> DomainResult<()> { + let schedule_config = serde_json::to_string(channel.schedule_config()) + .map_err(|e| DomainError::RepositoryError(format!("serialize schedule_config: {e}")))?; + let recycle_policy = serde_json::to_string(channel.recycle_policy()) + .map_err(|e| DomainError::RepositoryError(format!("serialize recycle_policy: {e}")))?; + let access_mode = serialize_enum_as_string(channel.access_mode(), "public"); + let logo_position = serialize_enum_as_string(channel.logo_position(), "top_right"); + + sqlx::query( + r#" + INSERT INTO channels + (id, owner_id, name, description, timezone, schedule_config, recycle_policy, + auto_schedule, access_mode, access_password_hash, logo, logo_position, + logo_opacity, webhook_url, webhook_poll_interval_secs, webhook_body_template, + webhook_headers, created_at, updated_at) + VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16, $17, $18, $19) + ON CONFLICT(id) DO UPDATE SET + name = EXCLUDED.name, + description = EXCLUDED.description, + timezone = EXCLUDED.timezone, + schedule_config = EXCLUDED.schedule_config, + recycle_policy = EXCLUDED.recycle_policy, + auto_schedule = EXCLUDED.auto_schedule, + access_mode = EXCLUDED.access_mode, + access_password_hash = EXCLUDED.access_password_hash, + logo = EXCLUDED.logo, + logo_position = EXCLUDED.logo_position, + logo_opacity = EXCLUDED.logo_opacity, + webhook_url = EXCLUDED.webhook_url, + webhook_poll_interval_secs = EXCLUDED.webhook_poll_interval_secs, + webhook_body_template = EXCLUDED.webhook_body_template, + webhook_headers = EXCLUDED.webhook_headers, + updated_at = EXCLUDED.updated_at + "#, + ) + .bind(channel.id().value().to_string()) + .bind(channel.owner_id().value().to_string()) + .bind(channel.name()) + .bind(channel.description()) + .bind(channel.timezone()) + .bind(&schedule_config) + .bind(&recycle_policy) + .bind(channel.auto_schedule()) + .bind(&access_mode) + .bind(channel.access_password_hash()) + .bind(channel.logo()) + .bind(&logo_position) + .bind(channel.logo_opacity()) + .bind(channel.webhook_url()) + .bind(channel.webhook_poll_interval_secs() as i64) + .bind(channel.webhook_body_template()) + .bind(channel.webhook_headers()) + .bind(channel.created_at().to_rfc3339()) + .bind(channel.updated_at().to_rfc3339()) + .execute(&self.pool) + .await + .map_err(map_sqlx_error)?; + + Ok(()) + } + + async fn delete(&self, id: ChannelId) -> DomainResult<()> { + sqlx::query("DELETE FROM channels WHERE id = $1") + .bind(id.value().to_string()) + .execute(&self.pool) + .await + .map_err(map_sqlx_error)?; + Ok(()) + } + + async fn save_config_snapshot( + &self, + channel_id: ChannelId, + config: &ScheduleConfig, + label: Option, + ) -> DomainResult { + let id = Uuid::new_v4(); + let now = Utc::now(); + let config_json = serde_json::to_string(config) + .map_err(|e| DomainError::RepositoryError(e.to_string()))?; + + let mut tx = self.pool.begin().await.map_err(map_sqlx_error)?; + + let version_num: i64 = sqlx::query_scalar( + "SELECT COALESCE(MAX(version_num), 0) + 1 FROM channel_config_snapshots WHERE channel_id = $1", + ) + .bind(channel_id.value().to_string()) + .fetch_one(&mut *tx) + .await + .map_err(map_sqlx_error)?; + + sqlx::query( + "INSERT INTO channel_config_snapshots (id, channel_id, config_json, version_num, label, created_at) + VALUES ($1, $2, $3, $4, $5, $6)", + ) + .bind(id.to_string()) + .bind(channel_id.value().to_string()) + .bind(&config_json) + .bind(version_num) + .bind(&label) + .bind(now.to_rfc3339()) + .execute(&mut *tx) + .await + .map_err(map_sqlx_error)?; + + tx.commit().await.map_err(map_sqlx_error)?; + + Ok(ChannelConfigSnapshot::from_persistence( + id, + channel_id, + config.clone(), + version_num, + label, + now, + )) + } + + async fn patch_config_snapshot_label( + &self, + channel_id: ChannelId, + snapshot_id: Uuid, + label: Option, + ) -> DomainResult> { + let updated = sqlx::query( + "UPDATE channel_config_snapshots SET label = $1 WHERE id = $2 AND channel_id = $3 RETURNING id", + ) + .bind(&label) + .bind(snapshot_id.to_string()) + .bind(channel_id.value().to_string()) + .fetch_optional(&self.pool) + .await + .map_err(map_sqlx_error)?; + + if updated.is_none() { + return Ok(None); + } + self.get_config_snapshot(channel_id, snapshot_id).await + } +} + +// -- Query ------------------------------------------------------------------- + +#[async_trait] +impl ChannelQuery for PgChannelRepository { + async fn find_by_id(&self, id: ChannelId) -> DomainResult> { + let sql = format!("SELECT {SELECT_COLS} FROM channels WHERE id = $1"); + let row: Option = sqlx::query_as(&sql) + .bind(id.value().to_string()) + .fetch_optional(&self.pool) + .await + .map_err(map_sqlx_error)?; + + row.map(ChannelRow::into_channel).transpose() + } + + async fn find_by_owner(&self, owner_id: UserId) -> DomainResult> { + let sql = format!( + "SELECT {SELECT_COLS} FROM channels WHERE owner_id = $1 ORDER BY created_at ASC" + ); + let rows: Vec = sqlx::query_as(&sql) + .bind(owner_id.value().to_string()) + .fetch_all(&self.pool) + .await + .map_err(map_sqlx_error)?; + + rows.into_iter().map(ChannelRow::into_channel).collect() + } + + async fn find_all(&self) -> DomainResult> { + let sql = format!("SELECT {SELECT_COLS} FROM channels ORDER BY created_at ASC"); + let rows: Vec = sqlx::query_as(&sql) + .fetch_all(&self.pool) + .await + .map_err(map_sqlx_error)?; + + rows.into_iter().map(ChannelRow::into_channel).collect() + } + + async fn find_auto_schedule_enabled(&self) -> DomainResult> { + let sql = format!( + "SELECT {SELECT_COLS} FROM channels WHERE auto_schedule = TRUE ORDER BY created_at ASC" + ); + let rows: Vec = sqlx::query_as(&sql) + .fetch_all(&self.pool) + .await + .map_err(map_sqlx_error)?; + + rows.into_iter().map(ChannelRow::into_channel).collect() + } + + async fn list_config_snapshots( + &self, + channel_id: ChannelId, + ) -> DomainResult> { + let rows = sqlx::query( + "SELECT id, config_json, version_num, label, created_at + FROM channel_config_snapshots WHERE channel_id = $1 + ORDER BY version_num DESC", + ) + .bind(channel_id.value().to_string()) + .fetch_all(&self.pool) + .await + .map_err(map_sqlx_error)?; + + rows.iter() + .map(|row| map_snapshot_row(row, channel_id)) + .collect() + } + + async fn get_config_snapshot( + &self, + channel_id: ChannelId, + snapshot_id: Uuid, + ) -> DomainResult> { + let row = sqlx::query( + "SELECT id, config_json, version_num, label, created_at + FROM channel_config_snapshots WHERE id = $1 AND channel_id = $2", + ) + .bind(snapshot_id.to_string()) + .bind(channel_id.value().to_string()) + .fetch_optional(&self.pool) + .await + .map_err(map_sqlx_error)?; + + match row { + None => Ok(None), + Some(row) => Ok(Some(map_snapshot_row(&row, channel_id)?)), + } + } +} diff --git a/crates/adapters/postgres/src/lib.rs b/crates/adapters/postgres/src/lib.rs new file mode 100644 index 0000000..36fcf0b --- /dev/null +++ b/crates/adapters/postgres/src/lib.rs @@ -0,0 +1,14 @@ +//! PostgreSQL adapter crate — implements all CQRS-split repository port traits +//! for PostgreSQL via sqlx. + +pub mod activity; +pub mod channel; +pub mod library; +pub mod provider_config; +pub mod schedule; +pub mod settings; +pub mod transcode; +pub mod user; +pub mod wire; + +pub use wire::{wire, PostgresWireOutput}; diff --git a/crates/adapters/postgres/src/library.rs b/crates/adapters/postgres/src/library.rs new file mode 100644 index 0000000..b0058d2 --- /dev/null +++ b/crates/adapters/postgres/src/library.rs @@ -0,0 +1,575 @@ +//! PostgreSQL adapter for library persistence (LibraryCommand + LibraryQuery). + +use std::collections::HashSet; + +use async_trait::async_trait; +use sqlx::PgPool; + +use domain::{ + ports::library::{LibraryCommand, LibraryQuery}, + ContentType, DomainError, DomainResult, LibraryCollection, LibraryItem, LibrarySearchFilter, + LibrarySyncLogEntry, LibrarySyncResult, SeasonSummary, ShowSummary, +}; + +pub struct PgLibraryRepository { + pool: PgPool, +} + +impl PgLibraryRepository { + pub fn new(pool: PgPool) -> Self { + Self { pool } + } +} + +// -- Helpers ----------------------------------------------------------------- + +fn content_type_str(ct: &ContentType) -> &'static str { + match ct { + ContentType::Movie => "movie", + ContentType::Episode => "episode", + ContentType::Short => "short", + } +} + +fn parse_content_type(s: &str) -> ContentType { + match s { + "episode" => ContentType::Episode, + "short" => ContentType::Short, + _ => ContentType::Movie, + } +} + +// -- Row types --------------------------------------------------------------- + +#[derive(sqlx::FromRow)] +struct LibraryItemRow { + id: String, + provider_id: String, + external_id: String, + title: String, + content_type: String, + duration_secs: i64, + series_name: Option, + season_number: Option, + episode_number: Option, + year: Option, + genres: String, + tags: String, + collection_id: Option, + collection_name: Option, + collection_type: Option, + thumbnail_url: Option, + synced_at: String, +} + +impl LibraryItemRow { + fn into_library_item(self) -> LibraryItem { + LibraryItem::from_persistence( + self.id, + self.provider_id, + self.external_id, + self.title, + parse_content_type(&self.content_type), + self.duration_secs as u32, + self.series_name, + self.season_number.map(|n| n as u32), + self.episode_number.map(|n| n as u32), + self.year.map(|n| n as u16), + serde_json::from_str(&self.genres).unwrap_or_default(), + serde_json::from_str(&self.tags).unwrap_or_default(), + self.collection_id, + self.collection_name, + self.collection_type, + self.thumbnail_url, + self.synced_at, + ) + } +} + +#[derive(sqlx::FromRow)] +struct SyncLogRow { + id: i64, + provider_id: String, + started_at: String, + finished_at: Option, + items_found: i64, + status: String, + error_msg: Option, +} + +#[derive(sqlx::FromRow)] +struct ShowSummaryRow { + series_name: String, + episode_count: i64, + season_count: i64, + thumbnail_url: Option, + genres_blob: String, +} + +#[derive(sqlx::FromRow)] +struct SeasonSummaryRow { + season_number: i64, + episode_count: i64, + thumbnail_url: Option, +} + +// -- Command ----------------------------------------------------------------- + +#[async_trait] +impl LibraryCommand for PgLibraryRepository { + async fn upsert_items(&self, _provider_id: &str, items: Vec) -> DomainResult<()> { + let mut tx = self + .pool + .begin() + .await + .map_err(|e| DomainError::InfrastructureError(e.to_string()))?; + + for item in items { + sqlx::query( + "INSERT INTO library_items + (id, provider_id, external_id, title, content_type, duration_secs, + series_name, season_number, episode_number, year, genres, tags, + collection_id, collection_name, collection_type, thumbnail_url, synced_at) + VALUES ($1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11,$12,$13,$14,$15,$16,$17) + ON CONFLICT(id) DO UPDATE SET + provider_id = EXCLUDED.provider_id, + external_id = EXCLUDED.external_id, + title = EXCLUDED.title, + content_type = EXCLUDED.content_type, + duration_secs = EXCLUDED.duration_secs, + series_name = EXCLUDED.series_name, + season_number = EXCLUDED.season_number, + episode_number = EXCLUDED.episode_number, + year = EXCLUDED.year, + genres = EXCLUDED.genres, + tags = EXCLUDED.tags, + collection_id = EXCLUDED.collection_id, + collection_name = EXCLUDED.collection_name, + collection_type = EXCLUDED.collection_type, + thumbnail_url = EXCLUDED.thumbnail_url, + synced_at = EXCLUDED.synced_at", + ) + .bind(item.id()) + .bind(item.provider_id()) + .bind(item.external_id()) + .bind(item.title()) + .bind(content_type_str(item.content_type())) + .bind(item.duration_secs() as i64) + .bind(item.series_name()) + .bind(item.season_number().map(|n| n as i64)) + .bind(item.episode_number().map(|n| n as i64)) + .bind(item.year().map(|n| n as i64)) + .bind(serde_json::to_string(item.genres()).unwrap_or_default()) + .bind(serde_json::to_string(item.tags()).unwrap_or_default()) + .bind(item.collection_id()) + .bind(item.collection_name()) + .bind(item.collection_type()) + .bind(item.thumbnail_url()) + .bind(item.synced_at()) + .execute(&mut *tx) + .await + .map_err(|e| DomainError::InfrastructureError(e.to_string()))?; + } + + tx.commit() + .await + .map_err(|e| DomainError::InfrastructureError(e.to_string())) + } + + async fn clear_provider(&self, provider_id: &str) -> DomainResult<()> { + sqlx::query("DELETE FROM library_items WHERE provider_id = $1") + .bind(provider_id) + .execute(&self.pool) + .await + .map(|_| ()) + .map_err(|e| DomainError::InfrastructureError(e.to_string())) + } + + async fn log_sync_start(&self, provider_id: &str) -> DomainResult { + let now = chrono::Utc::now().to_rfc3339(); + let id = sqlx::query_scalar::<_, i64>( + "INSERT INTO library_sync_log (provider_id, started_at, status) + VALUES ($1, $2, 'running') RETURNING id", + ) + .bind(provider_id) + .bind(&now) + .fetch_one(&self.pool) + .await + .map_err(|e| DomainError::InfrastructureError(e.to_string()))?; + Ok(id) + } + + async fn log_sync_finish(&self, log_id: i64, result: &LibrarySyncResult) -> DomainResult<()> { + let now = chrono::Utc::now().to_rfc3339(); + let status = if result.error().is_none() { + "done" + } else { + "error" + }; + sqlx::query( + "UPDATE library_sync_log + SET finished_at = $1, items_found = $2, status = $3, error_msg = $4 + WHERE id = $5", + ) + .bind(&now) + .bind(result.items_found() as i64) + .bind(status) + .bind(result.error()) + .bind(log_id) + .execute(&self.pool) + .await + .map(|_| ()) + .map_err(|e| DomainError::InfrastructureError(e.to_string())) + } +} + +// -- Query ------------------------------------------------------------------- + +#[async_trait] +impl LibraryQuery for PgLibraryRepository { + async fn search( + &self, + filter: &LibrarySearchFilter, + ) -> DomainResult<(Vec, u32)> { + let mut conditions: Vec = vec![]; + + if let Some(p) = filter.provider_id() { + conditions.push(format!("provider_id = '{}'", p.replace('\'', "''"))); + } + if let Some(ct) = filter.content_type() { + conditions.push(format!("content_type = '{}'", content_type_str(ct))); + } + if let Some(st) = filter.search_term() { + conditions.push(format!("title ILIKE '%{}%'", st.replace('\'', "''"))); + } + if let Some(cid) = filter.collection_id() { + conditions.push(format!("collection_id = '{}'", cid.replace('\'', "''"))); + } + if let Some(decade) = filter.decade() { + let end = decade + 10; + conditions.push(format!("year >= {} AND year < {}", decade, end)); + } + if let Some(min) = filter.min_duration_secs() { + conditions.push(format!("duration_secs >= {}", min)); + } + if let Some(max) = filter.max_duration_secs() { + conditions.push(format!("duration_secs <= {}", max)); + } + if !filter.series_names().is_empty() { + let quoted: Vec = filter + .series_names() + .iter() + .map(|s| format!("'{}'", s.replace('\'', "''"))) + .collect(); + conditions.push(format!("series_name IN ({})", quoted.join(","))); + } + if !filter.genres().is_empty() { + let genre_conditions: Vec = filter + .genres() + .iter() + .map(|g| { + format!( + "EXISTS (SELECT 1 FROM jsonb_array_elements_text(library_items.genres::jsonb) je WHERE je = '{}')", + g.replace('\'', "''") + ) + }) + .collect(); + conditions.push(format!("({})", genre_conditions.join(" OR "))); + } + if let Some(sn) = filter.season_number() { + conditions.push(format!("season_number = {}", sn)); + } + + let where_clause = if conditions.is_empty() { + String::new() + } else { + format!("WHERE {}", conditions.join(" AND ")) + }; + + let count_sql = format!("SELECT COUNT(*) FROM library_items {}", where_clause); + let total: i64 = sqlx::query_scalar(&count_sql) + .fetch_one(&self.pool) + .await + .map_err(|e| DomainError::InfrastructureError(e.to_string()))?; + + let items_sql = format!( + "SELECT * FROM library_items {} ORDER BY title ASC LIMIT {} OFFSET {}", + where_clause, + filter.limit(), + filter.offset() + ); + + let rows = sqlx::query_as::<_, LibraryItemRow>(&items_sql) + .fetch_all(&self.pool) + .await + .map_err(|e| DomainError::InfrastructureError(e.to_string()))?; + + Ok(( + rows.into_iter() + .map(LibraryItemRow::into_library_item) + .collect(), + total as u32, + )) + } + + async fn get_by_id(&self, id: &str) -> DomainResult> { + let row = sqlx::query_as::<_, LibraryItemRow>( + "SELECT * FROM library_items WHERE id = $1", + ) + .bind(id) + .fetch_optional(&self.pool) + .await + .map_err(|e| DomainError::InfrastructureError(e.to_string()))?; + Ok(row.map(LibraryItemRow::into_library_item)) + } + + async fn list_collections( + &self, + provider_id: Option<&str>, + ) -> DomainResult> { + let rows: Vec<(String, Option, Option)> = if let Some(p) = provider_id { + sqlx::query_as( + "SELECT DISTINCT collection_id, collection_name, collection_type + FROM library_items WHERE collection_id IS NOT NULL AND provider_id = $1 + ORDER BY collection_name ASC", + ) + .bind(p) + .fetch_all(&self.pool) + .await + } else { + sqlx::query_as( + "SELECT DISTINCT collection_id, collection_name, collection_type + FROM library_items WHERE collection_id IS NOT NULL + ORDER BY collection_name ASC", + ) + .fetch_all(&self.pool) + .await + } + .map_err(|e| DomainError::InfrastructureError(e.to_string()))?; + + Ok(rows + .into_iter() + .map(|(id, name, ct)| { + LibraryCollection::from_persistence(id, name.unwrap_or_default(), ct) + }) + .collect()) + } + + async fn list_series(&self, provider_id: Option<&str>) -> DomainResult> { + let rows: Vec<(String,)> = if let Some(p) = provider_id { + sqlx::query_as( + "SELECT DISTINCT series_name FROM library_items + WHERE series_name IS NOT NULL AND provider_id = $1 ORDER BY series_name ASC", + ) + .bind(p) + .fetch_all(&self.pool) + .await + } else { + sqlx::query_as( + "SELECT DISTINCT series_name FROM library_items + WHERE series_name IS NOT NULL ORDER BY series_name ASC", + ) + .fetch_all(&self.pool) + .await + } + .map_err(|e| DomainError::InfrastructureError(e.to_string()))?; + + Ok(rows.into_iter().map(|(s,)| s).collect()) + } + + async fn list_genres( + &self, + content_type: Option<&ContentType>, + provider_id: Option<&str>, + ) -> DomainResult> { + let sql = match (content_type, provider_id) { + (Some(ct), Some(p)) => format!( + "SELECT DISTINCT je AS value FROM library_items li, \ + LATERAL jsonb_array_elements_text(li.genres::jsonb) je \ + WHERE li.content_type = '{}' AND li.provider_id = '{}' ORDER BY value ASC", + content_type_str(ct), + p.replace('\'', "''") + ), + (Some(ct), None) => format!( + "SELECT DISTINCT je AS value FROM library_items li, \ + LATERAL jsonb_array_elements_text(li.genres::jsonb) je \ + WHERE li.content_type = '{}' ORDER BY value ASC", + content_type_str(ct) + ), + (None, Some(p)) => format!( + "SELECT DISTINCT je AS value FROM library_items li, \ + LATERAL jsonb_array_elements_text(li.genres::jsonb) je \ + WHERE li.provider_id = '{}' ORDER BY value ASC", + p.replace('\'', "''") + ), + (None, None) => { + "SELECT DISTINCT je AS value FROM library_items li, \ + LATERAL jsonb_array_elements_text(li.genres::jsonb) je \ + ORDER BY value ASC" + .to_string() + } + }; + let rows: Vec<(String,)> = sqlx::query_as(&sql) + .fetch_all(&self.pool) + .await + .map_err(|e| DomainError::InfrastructureError(e.to_string()))?; + Ok(rows.into_iter().map(|(s,)| s).collect()) + } + + async fn latest_sync_status(&self) -> DomainResult> { + let rows = sqlx::query_as::<_, SyncLogRow>( + "SELECT * FROM library_sync_log + WHERE id IN ( + SELECT MAX(id) FROM library_sync_log GROUP BY provider_id + ) + ORDER BY started_at DESC", + ) + .fetch_all(&self.pool) + .await + .map_err(|e| DomainError::InfrastructureError(e.to_string()))?; + + Ok(rows + .into_iter() + .map(|r| { + LibrarySyncLogEntry::from_persistence( + r.id, + r.provider_id, + r.started_at, + r.finished_at, + r.items_found as u32, + r.status, + r.error_msg, + ) + }) + .collect()) + } + + async fn is_sync_running(&self, provider_id: &str) -> DomainResult { + let count: i64 = sqlx::query_scalar( + "SELECT COUNT(*) FROM library_sync_log WHERE provider_id = $1 AND status = 'running'", + ) + .bind(provider_id) + .fetch_one(&self.pool) + .await + .map_err(|e| DomainError::InfrastructureError(e.to_string()))?; + Ok(count > 0) + } + + async fn list_shows( + &self, + provider_id: Option<&str>, + search_term: Option<&str>, + genres: &[String], + ) -> DomainResult> { + let mut conditions = vec![ + "content_type = 'episode'".to_string(), + "series_name IS NOT NULL".to_string(), + ]; + if let Some(p) = provider_id { + conditions.push(format!("provider_id = '{}'", p.replace('\'', "''"))); + } + if let Some(st) = search_term { + let escaped = st.replace('\'', "''"); + conditions.push(format!( + "(title ILIKE '%{escaped}%' OR series_name ILIKE '%{escaped}%')" + )); + } + if !genres.is_empty() { + let genre_conditions: Vec = genres + .iter() + .map(|g| { + format!( + "EXISTS (SELECT 1 FROM jsonb_array_elements_text(library_items.genres::jsonb) je WHERE je = '{}')", + g.replace('\'', "''") + ) + }) + .collect(); + conditions.push(format!("({})", genre_conditions.join(" OR "))); + } + + let where_clause = format!("WHERE {}", conditions.join(" AND ")); + let sql = format!( + "SELECT series_name, COUNT(*) AS episode_count, \ + COUNT(DISTINCT season_number) AS season_count, \ + MAX(thumbnail_url) AS thumbnail_url, \ + STRING_AGG(genres, ',') AS genres_blob \ + FROM library_items {} GROUP BY series_name ORDER BY series_name ASC", + where_clause + ); + + let rows = sqlx::query_as::<_, ShowSummaryRow>(&sql) + .fetch_all(&self.pool) + .await + .map_err(|e| DomainError::InfrastructureError(e.to_string()))?; + + Ok(rows + .into_iter() + .map(|r| { + let genres: Vec = r + .genres_blob + .split("],[") + .flat_map(|chunk| { + let cleaned = chunk.trim_start_matches('[').trim_end_matches(']'); + cleaned + .split(',') + .filter_map(|s| { + let s = s.trim().trim_matches('"'); + if s.is_empty() { + None + } else { + Some(s.to_string()) + } + }) + .collect::>() + }) + .collect::>() + .into_iter() + .collect(); + ShowSummary::from_persistence( + r.series_name, + r.episode_count as u32, + r.season_count as u32, + r.thumbnail_url, + genres, + ) + }) + .collect()) + } + + async fn list_seasons( + &self, + series_name: &str, + provider_id: Option<&str>, + ) -> DomainResult> { + let mut conditions = vec![ + format!("series_name = '{}'", series_name.replace('\'', "''")), + "content_type = 'episode'".to_string(), + ]; + if let Some(p) = provider_id { + conditions.push(format!("provider_id = '{}'", p.replace('\'', "''"))); + } + let where_clause = format!("WHERE {}", conditions.join(" AND ")); + let sql = format!( + "SELECT season_number, COUNT(*) AS episode_count, \ + MAX(thumbnail_url) AS thumbnail_url \ + FROM library_items {} GROUP BY season_number ORDER BY season_number ASC", + where_clause + ); + + let rows = sqlx::query_as::<_, SeasonSummaryRow>(&sql) + .fetch_all(&self.pool) + .await + .map_err(|e| DomainError::InfrastructureError(e.to_string()))?; + + Ok(rows + .into_iter() + .map(|r| { + SeasonSummary::from_persistence( + r.season_number as u32, + r.episode_count as u32, + r.thumbnail_url, + ) + }) + .collect()) + } +} diff --git a/crates/adapters/postgres/src/provider_config.rs b/crates/adapters/postgres/src/provider_config.rs new file mode 100644 index 0000000..856867c --- /dev/null +++ b/crates/adapters/postgres/src/provider_config.rs @@ -0,0 +1,92 @@ +//! PostgreSQL adapter for provider config (ProviderConfigCommand + ProviderConfigQuery). + +use async_trait::async_trait; +use sqlx::PgPool; + +use adapter_common::map_sqlx_error; +use domain::{ + ports::provider_config::{ProviderConfigCommand, ProviderConfigQuery}, + DomainResult, ProviderConfigRow, +}; + +pub struct PgProviderConfig { + pool: PgPool, +} + +impl PgProviderConfig { + pub fn new(pool: PgPool) -> Self { + Self { pool } + } +} + +#[async_trait] +impl ProviderConfigCommand for PgProviderConfig { + async fn upsert(&self, row: &ProviderConfigRow) -> DomainResult<()> { + sqlx::query( + r#"INSERT INTO provider_configs (id, provider_type, config_json, enabled, updated_at) + VALUES ($1, $2, $3, $4, $5) + ON CONFLICT(id) DO UPDATE SET + provider_type = EXCLUDED.provider_type, + config_json = EXCLUDED.config_json, + enabled = EXCLUDED.enabled, + updated_at = EXCLUDED.updated_at"#, + ) + .bind(row.id()) + .bind(row.provider_type()) + .bind(row.config_json()) + .bind(row.enabled()) + .bind(row.updated_at()) + .execute(&self.pool) + .await + .map_err(map_sqlx_error)?; + Ok(()) + } + + async fn delete(&self, id: &str) -> DomainResult<()> { + sqlx::query("DELETE FROM provider_configs WHERE id = $1") + .bind(id) + .execute(&self.pool) + .await + .map_err(map_sqlx_error)?; + Ok(()) + } +} + +#[async_trait] +impl ProviderConfigQuery for PgProviderConfig { + async fn get_all(&self) -> DomainResult> { + let rows: Vec<(String, String, String, bool, String)> = sqlx::query_as( + "SELECT id, provider_type, config_json, enabled, updated_at FROM provider_configs", + ) + .fetch_all(&self.pool) + .await + .map_err(map_sqlx_error)?; + + Ok(rows + .into_iter() + .map(|(id, provider_type, config_json, enabled, updated_at)| { + ProviderConfigRow::from_persistence( + id, + provider_type, + config_json, + enabled, + updated_at, + ) + }) + .collect()) + } + + async fn get_by_id(&self, id: &str) -> DomainResult> { + let row: Option<(String, String, String, bool, String)> = sqlx::query_as( + "SELECT id, provider_type, config_json, enabled, updated_at FROM provider_configs WHERE id = $1", + ) + .bind(id) + .fetch_optional(&self.pool) + .await + .map_err(map_sqlx_error)?; + + Ok(row.map(|(id, provider_type, config_json, enabled, updated_at)| { + ProviderConfigRow::from_persistence(id, provider_type, config_json, enabled, updated_at) + })) + } +} diff --git a/crates/adapters/postgres/src/schedule.rs b/crates/adapters/postgres/src/schedule.rs new file mode 100644 index 0000000..e1b016c --- /dev/null +++ b/crates/adapters/postgres/src/schedule.rs @@ -0,0 +1,362 @@ +//! PostgreSQL adapter for schedule persistence (ScheduleCommand + ScheduleQuery). + +use std::collections::HashMap; + +use async_trait::async_trait; +use sqlx::PgPool; + +use adapter_common::{map_sqlx_error, parse_dt, parse_json, parse_uuid}; +use domain::{ + ports::schedule::{ScheduleCommand, ScheduleQuery}, + BlockId, ChannelId, DomainError, DomainResult, GeneratedSchedule, MediaItem, MediaItemId, + PlaybackRecord, ScheduleId, ScheduledSlot, SlotId, +}; + +pub struct PgScheduleRepository { + pool: PgPool, +} + +impl PgScheduleRepository { + pub fn new(pool: PgPool) -> Self { + Self { pool } + } +} + +// -- Row types --------------------------------------------------------------- + +#[derive(Debug, sqlx::FromRow)] +struct ScheduleRow { + id: String, + channel_id: String, + valid_from: String, + valid_until: String, + generation: i64, +} + +#[derive(Debug, sqlx::FromRow)] +struct SlotRow { + id: String, + #[allow(dead_code)] + schedule_id: String, + start_at: String, + end_at: String, + item: String, + source_block_id: String, +} + +#[derive(Debug, sqlx::FromRow)] +struct LastSlotRow { + source_block_id: String, + item: String, +} + +#[derive(Debug, sqlx::FromRow)] +struct PlaybackRecordRow { + id: String, + channel_id: String, + item_id: String, + played_at: String, + generation: i64, +} + +// -- Mapping ----------------------------------------------------------------- + +fn map_slot_row(row: SlotRow) -> DomainResult { + let id = SlotId::from_uuid(parse_uuid(&row.id, "slot id")?); + let source_block_id = BlockId::from_uuid(parse_uuid(&row.source_block_id, "block id")?); + let item: MediaItem = parse_json(&row.item, "slot item")?; + + Ok(ScheduledSlot::from_persistence( + id, + parse_dt(&row.start_at)?, + parse_dt(&row.end_at)?, + item, + source_block_id, + )) +} + +fn map_schedule(row: ScheduleRow, slot_rows: Vec) -> DomainResult { + let id = ScheduleId::from_uuid(parse_uuid(&row.id, "schedule id")?); + let channel_id = ChannelId::from_uuid(parse_uuid(&row.channel_id, "channel id")?); + let slots: Result, _> = slot_rows.into_iter().map(map_slot_row).collect(); + + Ok(GeneratedSchedule::from_persistence( + id, + channel_id, + parse_dt(&row.valid_from)?, + parse_dt(&row.valid_until)?, + row.generation as u32, + slots?, + )) +} + +fn map_playback_row(row: PlaybackRecordRow) -> DomainResult { + let id = parse_uuid(&row.id, "playback record id")?; + let channel_id = ChannelId::from_uuid(parse_uuid(&row.channel_id, "channel id")?); + + Ok(PlaybackRecord::from_persistence( + id, + channel_id, + MediaItemId::new(row.item_id), + parse_dt(&row.played_at)?, + row.generation as u32, + )) +} + +// -- Internal helpers -------------------------------------------------------- + +impl PgScheduleRepository { + async fn fetch_slots(&self, schedule_id: &str) -> DomainResult> { + sqlx::query_as( + "SELECT id, schedule_id, start_at, end_at, item, source_block_id \ + FROM scheduled_slots WHERE schedule_id = $1 ORDER BY start_at", + ) + .bind(schedule_id) + .fetch_all(&self.pool) + .await + .map_err(map_sqlx_error) + } +} + +// -- Command ----------------------------------------------------------------- + +#[async_trait] +impl ScheduleCommand for PgScheduleRepository { + async fn save(&self, schedule: &GeneratedSchedule) -> DomainResult<()> { + sqlx::query( + r#" + INSERT INTO generated_schedules (id, channel_id, valid_from, valid_until, generation) + VALUES ($1, $2, $3, $4, $5) + ON CONFLICT(id) DO UPDATE SET + valid_from = EXCLUDED.valid_from, + valid_until = EXCLUDED.valid_until, + generation = EXCLUDED.generation + "#, + ) + .bind(schedule.id().value().to_string()) + .bind(schedule.channel_id().value().to_string()) + .bind(schedule.valid_from().to_rfc3339()) + .bind(schedule.valid_until().to_rfc3339()) + .bind(schedule.generation() as i64) + .execute(&self.pool) + .await + .map_err(map_sqlx_error)?; + + // Delete-then-insert all slots + sqlx::query("DELETE FROM scheduled_slots WHERE schedule_id = $1") + .bind(schedule.id().value().to_string()) + .execute(&self.pool) + .await + .map_err(map_sqlx_error)?; + + for slot in schedule.slots() { + let item_json = serde_json::to_string(slot.item()) + .map_err(|e| DomainError::RepositoryError(format!("serialize slot item: {e}")))?; + + sqlx::query( + "INSERT INTO scheduled_slots (id, schedule_id, start_at, end_at, item, source_block_id) + VALUES ($1, $2, $3, $4, $5, $6)", + ) + .bind(slot.id().value().to_string()) + .bind(schedule.id().value().to_string()) + .bind(slot.start_at().to_rfc3339()) + .bind(slot.end_at().to_rfc3339()) + .bind(&item_json) + .bind(slot.source_block_id().value().to_string()) + .execute(&self.pool) + .await + .map_err(map_sqlx_error)?; + } + + Ok(()) + } + + async fn save_playback_record(&self, record: &PlaybackRecord) -> DomainResult<()> { + sqlx::query( + r#" + INSERT INTO playback_records (id, channel_id, item_id, played_at, generation) + VALUES ($1, $2, $3, $4, $5) + ON CONFLICT(id) DO NOTHING + "#, + ) + .bind(record.id().to_string()) + .bind(record.channel_id().value().to_string()) + .bind(record.item_id().value()) + .bind(record.played_at().to_rfc3339()) + .bind(record.generation() as i64) + .execute(&self.pool) + .await + .map_err(map_sqlx_error)?; + Ok(()) + } + + async fn delete_schedules_after( + &self, + channel_id: ChannelId, + target_generation: u32, + ) -> DomainResult<()> { + let ch = channel_id.value().to_string(); + let target_gen = target_generation as i64; + + sqlx::query("DELETE FROM playback_records WHERE channel_id = $1 AND generation > $2") + .bind(&ch) + .bind(target_gen) + .execute(&self.pool) + .await + .map_err(map_sqlx_error)?; + + sqlx::query("DELETE FROM generated_schedules WHERE channel_id = $1 AND generation > $2") + .bind(&ch) + .bind(target_gen) + .execute(&self.pool) + .await + .map_err(map_sqlx_error)?; + + Ok(()) + } +} + +// -- Query ------------------------------------------------------------------- + +#[async_trait] +impl ScheduleQuery for PgScheduleRepository { + async fn find_active( + &self, + channel_id: ChannelId, + at: chrono::DateTime, + ) -> DomainResult> { + let at_str = at.to_rfc3339(); + let row: Option = sqlx::query_as( + "SELECT id, channel_id, valid_from, valid_until, generation \ + FROM generated_schedules \ + WHERE channel_id = $1 AND valid_from <= $2 AND valid_until > $3 \ + LIMIT 1", + ) + .bind(channel_id.value().to_string()) + .bind(&at_str) + .bind(&at_str) + .fetch_optional(&self.pool) + .await + .map_err(map_sqlx_error)?; + + match row { + None => Ok(None), + Some(r) => { + let slots = self.fetch_slots(&r.id).await?; + Some(map_schedule(r, slots)).transpose() + } + } + } + + async fn find_latest(&self, channel_id: ChannelId) -> DomainResult> { + let row: Option = sqlx::query_as( + "SELECT id, channel_id, valid_from, valid_until, generation \ + FROM generated_schedules \ + WHERE channel_id = $1 ORDER BY valid_from DESC LIMIT 1", + ) + .bind(channel_id.value().to_string()) + .fetch_optional(&self.pool) + .await + .map_err(map_sqlx_error)?; + + match row { + None => Ok(None), + Some(r) => { + let slots = self.fetch_slots(&r.id).await?; + Some(map_schedule(r, slots)).transpose() + } + } + } + + async fn find_playback_history( + &self, + channel_id: ChannelId, + ) -> DomainResult> { + let rows: Vec = sqlx::query_as( + "SELECT id, channel_id, item_id, played_at, generation \ + FROM playback_records WHERE channel_id = $1 ORDER BY played_at DESC", + ) + .bind(channel_id.value().to_string()) + .fetch_all(&self.pool) + .await + .map_err(map_sqlx_error)?; + + rows.into_iter().map(map_playback_row).collect() + } + + async fn find_last_slot_per_block( + &self, + channel_id: ChannelId, + ) -> DomainResult> { + let channel_id_str = channel_id.value().to_string(); + let rows: Vec = sqlx::query_as( + "SELECT ss.source_block_id, ss.item \ + FROM scheduled_slots ss \ + INNER JOIN generated_schedules gs ON gs.id = ss.schedule_id \ + WHERE gs.channel_id = $1 \ + AND ss.start_at = ( \ + SELECT MAX(ss2.start_at) \ + FROM scheduled_slots ss2 \ + INNER JOIN generated_schedules gs2 ON gs2.id = ss2.schedule_id \ + WHERE ss2.source_block_id = ss.source_block_id \ + AND gs2.channel_id = $2 \ + )", + ) + .bind(&channel_id_str) + .bind(&channel_id_str) + .fetch_all(&self.pool) + .await + .map_err(map_sqlx_error)?; + + let mut map = HashMap::new(); + for row in rows { + let block_id = + BlockId::from_uuid(parse_uuid(&row.source_block_id, "block id")?); + let item: MediaItem = parse_json(&row.item, "slot item")?; + map.insert(block_id, item.id().clone()); + } + Ok(map) + } + + async fn list_schedule_history( + &self, + channel_id: ChannelId, + ) -> DomainResult> { + let rows: Vec = sqlx::query_as( + "SELECT id, channel_id, valid_from, valid_until, generation \ + FROM generated_schedules WHERE channel_id = $1 ORDER BY generation DESC", + ) + .bind(channel_id.value().to_string()) + .fetch_all(&self.pool) + .await + .map_err(map_sqlx_error)?; + + rows.into_iter() + .map(|r| map_schedule(r, vec![])) + .collect() + } + + async fn get_schedule_by_id( + &self, + channel_id: ChannelId, + schedule_id: ScheduleId, + ) -> DomainResult> { + let row: Option = sqlx::query_as( + "SELECT id, channel_id, valid_from, valid_until, generation \ + FROM generated_schedules WHERE id = $1 AND channel_id = $2", + ) + .bind(schedule_id.value().to_string()) + .bind(channel_id.value().to_string()) + .fetch_optional(&self.pool) + .await + .map_err(map_sqlx_error)?; + + match row { + None => Ok(None), + Some(r) => { + let slots = self.fetch_slots(&r.id).await?; + Some(map_schedule(r, slots)).transpose() + } + } + } +} diff --git a/crates/adapters/postgres/src/settings.rs b/crates/adapters/postgres/src/settings.rs new file mode 100644 index 0000000..259c050 --- /dev/null +++ b/crates/adapters/postgres/src/settings.rs @@ -0,0 +1,50 @@ +//! PostgreSQL adapter for app settings (AppSettingsRepository). + +use async_trait::async_trait; +use sqlx::PgPool; + +use domain::{ + ports::settings::AppSettingsRepository, + DomainError, DomainResult, +}; + +pub struct PgAppSettings { + pool: PgPool, +} + +impl PgAppSettings { + pub fn new(pool: PgPool) -> Self { + Self { pool } + } +} + +#[async_trait] +impl AppSettingsRepository for PgAppSettings { + async fn get(&self, key: &str) -> DomainResult> { + sqlx::query_scalar::<_, String>("SELECT value FROM app_settings WHERE key = $1") + .bind(key) + .fetch_optional(&self.pool) + .await + .map_err(|e| DomainError::InfrastructureError(e.to_string())) + } + + async fn set(&self, key: &str, value: &str) -> DomainResult<()> { + sqlx::query( + "INSERT INTO app_settings (key, value) VALUES ($1, $2) + ON CONFLICT(key) DO UPDATE SET value = EXCLUDED.value", + ) + .bind(key) + .bind(value) + .execute(&self.pool) + .await + .map(|_| ()) + .map_err(|e| DomainError::InfrastructureError(e.to_string())) + } + + async fn get_all(&self) -> DomainResult> { + sqlx::query_as::<_, (String, String)>("SELECT key, value FROM app_settings ORDER BY key") + .fetch_all(&self.pool) + .await + .map_err(|e| DomainError::InfrastructureError(e.to_string())) + } +} diff --git a/crates/adapters/postgres/src/transcode.rs b/crates/adapters/postgres/src/transcode.rs new file mode 100644 index 0000000..60863a8 --- /dev/null +++ b/crates/adapters/postgres/src/transcode.rs @@ -0,0 +1,40 @@ +//! PostgreSQL adapter for transcode settings (TranscodeSettingsRepository). + +use async_trait::async_trait; +use sqlx::PgPool; + +use domain::{ + ports::transcode::TranscodeSettingsRepository, + DomainError, DomainResult, +}; + +pub struct PgTranscodeSettings { + pool: PgPool, +} + +impl PgTranscodeSettings { + pub fn new(pool: PgPool) -> Self { + Self { pool } + } +} + +#[async_trait] +impl TranscodeSettingsRepository for PgTranscodeSettings { + async fn load_cleanup_ttl(&self) -> DomainResult> { + let row: Option<(i64,)> = + sqlx::query_as("SELECT cleanup_ttl_hours FROM transcode_settings WHERE id = 1") + .fetch_optional(&self.pool) + .await + .map_err(|e| DomainError::InfrastructureError(e.to_string()))?; + Ok(row.map(|(h,)| h as u32)) + } + + async fn save_cleanup_ttl(&self, hours: u32) -> DomainResult<()> { + sqlx::query("UPDATE transcode_settings SET cleanup_ttl_hours = $1 WHERE id = 1") + .bind(hours as i64) + .execute(&self.pool) + .await + .map_err(|e| DomainError::InfrastructureError(e.to_string()))?; + Ok(()) + } +} diff --git a/crates/adapters/postgres/src/user.rs b/crates/adapters/postgres/src/user.rs new file mode 100644 index 0000000..e312110 --- /dev/null +++ b/crates/adapters/postgres/src/user.rs @@ -0,0 +1,148 @@ +//! PostgreSQL adapter for user persistence (UserCommand + UserQuery). + +use async_trait::async_trait; +use sqlx::PgPool; + +use adapter_common::{map_sqlx_error, parse_dt, parse_uuid}; +use domain::{ + ports::user::{UserCommand, UserQuery}, + DomainError, DomainResult, Email, User, UserId, +}; + +pub struct PgUserRepository { + pool: PgPool, +} + +impl PgUserRepository { + pub fn new(pool: PgPool) -> Self { + Self { pool } + } +} + +// -- Row type for query_as -------------------------------------------------- + +#[derive(Debug, sqlx::FromRow)] +struct UserRow { + id: String, + subject: String, + email: String, + password_hash: Option, + is_admin: bool, + created_at: String, +} + +impl UserRow { + fn into_user(self) -> DomainResult { + let id = UserId::from_uuid(parse_uuid(&self.id, "user id")?); + let email = Email::new(&self.email) + .map_err(|e| DomainError::RepositoryError(format!("Invalid email: {e}")))?; + let created_at = parse_dt(&self.created_at)?; + + Ok(User::from_persistence( + id, + self.subject, + email, + self.password_hash, + self.is_admin, + created_at, + )) + } +} + +// -- Command ----------------------------------------------------------------- + +#[async_trait] +impl UserCommand for PgUserRepository { + async fn save(&self, user: &User) -> DomainResult<()> { + let id = user.id().value().to_string(); + let created_at = user.created_at().to_rfc3339(); + + sqlx::query( + r#" + INSERT INTO users (id, subject, email, password_hash, is_admin, created_at) + VALUES ($1, $2, $3, $4, $5, $6) + ON CONFLICT(id) DO UPDATE SET + subject = EXCLUDED.subject, + email = EXCLUDED.email, + password_hash = EXCLUDED.password_hash, + is_admin = EXCLUDED.is_admin + "#, + ) + .bind(&id) + .bind(user.subject()) + .bind(user.email().as_ref()) + .bind(user.password_hash()) + .bind(user.is_admin()) + .bind(&created_at) + .execute(&self.pool) + .await + .map_err(|e| { + let msg = e.to_string(); + if msg.contains("unique constraint") || msg.contains("duplicate key") { + DomainError::UserAlreadyExists(user.email().as_ref().to_string()) + } else { + map_sqlx_error(e) + } + })?; + + Ok(()) + } + + async fn delete(&self, id: UserId) -> DomainResult<()> { + sqlx::query("DELETE FROM users WHERE id = $1") + .bind(id.value().to_string()) + .execute(&self.pool) + .await + .map_err(map_sqlx_error)?; + Ok(()) + } +} + +// -- Query ------------------------------------------------------------------- + +#[async_trait] +impl UserQuery for PgUserRepository { + async fn find_by_id(&self, id: UserId) -> DomainResult> { + let row: Option = sqlx::query_as( + "SELECT id, subject, email, password_hash, is_admin, created_at FROM users WHERE id = $1", + ) + .bind(id.value().to_string()) + .fetch_optional(&self.pool) + .await + .map_err(map_sqlx_error)?; + + row.map(UserRow::into_user).transpose() + } + + async fn find_by_subject(&self, subject: &str) -> DomainResult> { + let row: Option = sqlx::query_as( + "SELECT id, subject, email, password_hash, is_admin, created_at FROM users WHERE subject = $1", + ) + .bind(subject) + .fetch_optional(&self.pool) + .await + .map_err(map_sqlx_error)?; + + row.map(UserRow::into_user).transpose() + } + + async fn find_by_email(&self, email: &str) -> DomainResult> { + let row: Option = sqlx::query_as( + "SELECT id, subject, email, password_hash, is_admin, created_at FROM users WHERE email = $1", + ) + .bind(email) + .fetch_optional(&self.pool) + .await + .map_err(map_sqlx_error)?; + + row.map(UserRow::into_user).transpose() + } + + async fn count_users(&self) -> DomainResult { + let (count,): (i64,) = sqlx::query_as("SELECT COUNT(*) FROM users") + .fetch_one(&self.pool) + .await + .map_err(map_sqlx_error)?; + Ok(count as u64) + } +} diff --git a/crates/adapters/postgres/src/wire.rs b/crates/adapters/postgres/src/wire.rs new file mode 100644 index 0000000..e891aa0 --- /dev/null +++ b/crates/adapters/postgres/src/wire.rs @@ -0,0 +1,78 @@ +//! Wiring function that instantiates all PostgreSQL repositories and returns them +//! as trait-object Arcs. + +use std::sync::Arc; + +use sqlx::PgPool; + +use domain::ports::{ + activity::{ActivityLogCommand, ActivityLogQuery}, + channel::{ChannelCommand, ChannelQuery}, + library::{LibraryCommand, LibraryQuery}, + provider_config::{ProviderConfigCommand, ProviderConfigQuery}, + schedule::{ScheduleCommand, ScheduleQuery}, + settings::AppSettingsRepository, + transcode::TranscodeSettingsRepository, + user::{UserCommand, UserQuery}, +}; + +use crate::{ + activity::PgActivityLog, + channel::PgChannelRepository, + library::PgLibraryRepository, + provider_config::PgProviderConfig, + schedule::PgScheduleRepository, + settings::PgAppSettings, + transcode::PgTranscodeSettings, + user::PgUserRepository, +}; + +/// All PostgreSQL adapter outputs, ready to be injected into the application layer. +pub struct PostgresWireOutput { + pub user_command: Arc, + pub user_query: Arc, + pub channel_command: Arc, + pub channel_query: Arc, + pub schedule_command: Arc, + pub schedule_query: Arc, + pub library_command: Arc, + pub library_query: Arc, + pub activity_command: Arc, + pub activity_query: Arc, + pub settings: Arc, + pub provider_config_command: Arc, + pub provider_config_query: Arc, + pub transcode_settings: Arc, +} + +/// Create all PostgreSQL repository implementations from a single pool. +/// +/// Each struct wraps a clone of the same pool. Repositories that implement +/// both Command and Query traits share a single `Arc` via `.clone()`. +pub fn wire(pool: PgPool) -> PostgresWireOutput { + let user = Arc::new(PgUserRepository::new(pool.clone())); + let channel = Arc::new(PgChannelRepository::new(pool.clone())); + let schedule = Arc::new(PgScheduleRepository::new(pool.clone())); + let library = Arc::new(PgLibraryRepository::new(pool.clone())); + let activity = Arc::new(PgActivityLog::new(pool.clone())); + let settings = Arc::new(PgAppSettings::new(pool.clone())); + let provider_config = Arc::new(PgProviderConfig::new(pool.clone())); + let transcode_settings = Arc::new(PgTranscodeSettings::new(pool)); + + PostgresWireOutput { + user_command: user.clone(), + user_query: user, + channel_command: channel.clone(), + channel_query: channel, + schedule_command: schedule.clone(), + schedule_query: schedule, + library_command: library.clone(), + library_query: library, + activity_command: activity.clone(), + activity_query: activity, + settings, + provider_config_command: provider_config.clone(), + provider_config_query: provider_config, + transcode_settings, + } +}