diff --git a/crates/adapters/jellyfin/src/mapping.rs b/crates/adapters/jellyfin/src/mapping.rs index f0ad08b..5a9c089 100644 --- a/crates/adapters/jellyfin/src/mapping.rs +++ b/crates/adapters/jellyfin/src/mapping.rs @@ -1,4 +1,4 @@ -use domain::{ContentType, MediaItem, MediaItemId, MediaItemRow}; +use domain::{ContentType, MediaItem, MediaItemId, MediaItemRow, MediaRole}; use crate::models::JellyfinItem; @@ -17,7 +17,7 @@ pub(crate) fn map_jellyfin_item(item: JellyfinItem) -> Option { .unwrap_or(0); Some(MediaItem::from_persistence(MediaItemRow { - id: MediaItemId::new(item.id), + id: MediaItemId::new(&item.id), title: item.name, content_type, duration_secs, @@ -30,5 +30,11 @@ pub(crate) fn map_jellyfin_item(item: JellyfinItem) -> Option { episode_number: item.index_number, thumbnail_url: None, collection_id: None, + provider_id: String::new(), + external_id: item.id, + collection_name: None, + collection_type: None, + synced_at: None, + role: MediaRole::default(), })) } diff --git a/crates/adapters/local-files/src/provider.rs b/crates/adapters/local-files/src/provider.rs index 5666e79..8ef9564 100644 --- a/crates/adapters/local-files/src/provider.rs +++ b/crates/adapters/local-files/src/provider.rs @@ -4,7 +4,7 @@ use async_trait::async_trait; use domain::ports::{ Collection, IMediaProvider, ProviderCapabilities, StreamQuality, StreamingProtocol, }; -use domain::{ContentType, DomainError, DomainResult, MediaFilter, MediaItem, MediaItemId, MediaItemRow}; +use domain::{ContentType, DomainError, DomainResult, MediaFilter, MediaItem, MediaItemId, MediaItemRow, MediaRole}; use crate::config::LocalFilesConfig; use crate::index::{decode_id, LocalIndex}; @@ -54,6 +54,12 @@ fn to_media_item(id: MediaItemId, item: &LocalFileItem) -> MediaItem { episode_number: None, thumbnail_url: None, collection_id: None, + provider_id: String::new(), + external_id: String::new(), + collection_name: None, + collection_type: None, + synced_at: None, + role: MediaRole::default(), }) } diff --git a/crates/adapters/sqlite/src/library.rs b/crates/adapters/sqlite/src/library.rs index 1743859..ca204e6 100644 --- a/crates/adapters/sqlite/src/library.rs +++ b/crates/adapters/sqlite/src/library.rs @@ -4,9 +4,10 @@ use sqlx::SqlitePool; use adapter_common::{content_type_str, parse_content_type, parse_genres_blob}; use domain::{ ports::library::{LibraryCommand, LibraryQuery}, - ContentType, DomainError, DomainResult, LibraryCollection, LibraryItem, - LibraryItemRow as DomainLibraryItemRow, LibrarySearchFilter, LibrarySyncLogEntry, - LibrarySyncResult, SeasonSummary, ShowSummary, + ContentType, DomainError, DomainResult, LibraryCollection, + LibrarySearchFilter, LibrarySyncLogEntry, + LibrarySyncResult, MediaItem, MediaItemRow as DomainMediaItemRow, + MediaRole, SeasonSummary, ShowSummary, }; pub struct SqliteLibraryRepository { @@ -41,14 +42,15 @@ struct LibraryItemRow { } impl LibraryItemRow { - fn into_library_item(self) -> LibraryItem { - LibraryItem::from_persistence(DomainLibraryItemRow { - id: self.id, + fn into_media_item(self) -> MediaItem { + MediaItem::from_persistence(DomainMediaItemRow { + id: domain::MediaItemId::new(&self.id), provider_id: self.provider_id, external_id: self.external_id, title: self.title, content_type: parse_content_type(&self.content_type), duration_secs: self.duration_secs as u32, + description: None, series_name: self.series_name, season_number: self.season_number.map(|n| n as u32), episode_number: self.episode_number.map(|n| n as u32), @@ -59,7 +61,8 @@ impl LibraryItemRow { collection_name: self.collection_name, collection_type: self.collection_type, thumbnail_url: self.thumbnail_url, - synced_at: self.synced_at, + synced_at: Some(self.synced_at), + role: MediaRole::default(), }) } } @@ -93,7 +96,7 @@ struct SeasonSummaryRow { #[async_trait] impl LibraryCommand for SqliteLibraryRepository { - async fn upsert_items(&self, _provider_id: &str, items: Vec) -> DomainResult<()> { + async fn upsert_items(&self, _provider_id: &str, items: Vec) -> DomainResult<()> { let mut tx = self .pool .begin() @@ -108,7 +111,7 @@ impl LibraryCommand for SqliteLibraryRepository { collection_id, collection_name, collection_type, thumbnail_url, synced_at) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)", ) - .bind(item.id()) + .bind(item.id().value()) .bind(item.provider_id()) .bind(item.external_id()) .bind(item.title()) @@ -124,7 +127,7 @@ impl LibraryCommand for SqliteLibraryRepository { .bind(item.collection_name()) .bind(item.collection_type()) .bind(item.thumbnail_url()) - .bind(item.synced_at()) + .bind(item.synced_at().unwrap_or("")) .execute(&mut *tx) .await .map_err(|e| DomainError::InfrastructureError(e.to_string()))?; @@ -187,7 +190,7 @@ impl LibraryQuery for SqliteLibraryRepository { async fn search( &self, filter: &LibrarySearchFilter, - ) -> DomainResult<(Vec, u32)> { + ) -> DomainResult<(Vec, u32)> { let mut conditions: Vec = vec![]; if let Some(p) = filter.provider_id() { @@ -263,19 +266,19 @@ impl LibraryQuery for SqliteLibraryRepository { Ok(( rows.into_iter() - .map(LibraryItemRow::into_library_item) + .map(LibraryItemRow::into_media_item) .collect(), total as u32, )) } - async fn get_by_id(&self, id: &str) -> DomainResult> { + async fn get_by_id(&self, id: &str) -> DomainResult> { let row = sqlx::query_as::<_, LibraryItemRow>("SELECT * FROM library_items WHERE id = ?") .bind(id) .fetch_optional(&self.pool) .await .map_err(|e| DomainError::InfrastructureError(e.to_string()))?; - Ok(row.map(LibraryItemRow::into_library_item)) + Ok(row.map(LibraryItemRow::into_media_item)) } async fn list_collections( diff --git a/crates/api-types/src/library.rs b/crates/api-types/src/library.rs index 8ca0a54..e81491a 100644 --- a/crates/api-types/src/library.rs +++ b/crates/api-types/src/library.rs @@ -21,13 +21,13 @@ pub struct LibraryItemResponse { pub collection_name: Option, pub collection_type: Option, pub thumbnail_url: Option, - pub synced_at: String, + pub synced_at: Option, } -impl From for LibraryItemResponse { - fn from(i: domain::LibraryItem) -> Self { +impl From for LibraryItemResponse { + fn from(i: domain::MediaItem) -> Self { Self { - id: i.id().to_string(), + id: i.id().value().to_string(), provider_id: i.provider_id().to_string(), external_id: i.external_id().to_string(), title: i.title().to_string(), @@ -43,7 +43,7 @@ impl From for LibraryItemResponse { collection_name: i.collection_name().map(|s| s.to_string()), collection_type: i.collection_type().map(|s| s.to_string()), thumbnail_url: i.thumbnail_url().map(|s| s.to_string()), - synced_at: i.synced_at().to_string(), + synced_at: i.synced_at().map(|s| s.to_string()), } } } diff --git a/crates/application/src/library/get_item.rs b/crates/application/src/library/get_item.rs index 6e86b94..604e5bd 100644 --- a/crates/application/src/library/get_item.rs +++ b/crates/application/src/library/get_item.rs @@ -1,4 +1,4 @@ -use domain::models::LibraryItem; +use domain::models::MediaItem; use domain::DomainResult; use super::deps::LibraryQueryDeps; @@ -7,7 +7,7 @@ use super::queries::GetItemQuery; pub async fn execute( deps: &LibraryQueryDeps, query: GetItemQuery, -) -> DomainResult> { +) -> DomainResult> { deps.library_query.get_by_id(&query.item_id).await } diff --git a/crates/application/src/library/search.rs b/crates/application/src/library/search.rs index d6dd69e..79d5d6a 100644 --- a/crates/application/src/library/search.rs +++ b/crates/application/src/library/search.rs @@ -1,5 +1,5 @@ use domain::DomainResult; -use domain::models::LibraryItem; +use domain::models::MediaItem; use domain::value_objects::LibrarySearchFilter; use super::deps::LibraryQueryDeps; @@ -9,7 +9,7 @@ use super::queries::SearchItemsQuery; pub async fn execute( deps: &LibraryQueryDeps, query: SearchItemsQuery, -) -> DomainResult<(Vec, u32)> { +) -> DomainResult<(Vec, u32)> { let content_type = query .content_type .as_deref() diff --git a/crates/application/src/library/tests/get_item.rs b/crates/application/src/library/tests/get_item.rs index 9ee7144..6921c81 100644 --- a/crates/application/src/library/tests/get_item.rs +++ b/crates/application/src/library/tests/get_item.rs @@ -1,4 +1,4 @@ -use domain::models::LibraryItem; +use domain::models::MediaItem; use domain::value_objects::ContentType; use crate::library::get_item; @@ -8,11 +8,11 @@ use crate::library::queries::GetItemQuery; mod helpers; fn seed_item(repo: &std::sync::Arc) { - let item = LibraryItem::new("test", "m1", "Die Hard", ContentType::Movie, 7800, "2026-01-01"); + let item = MediaItem::new_library("test", "m1", "Die Hard", ContentType::Movie, 7800, "2026-01-01"); repo.items .lock() .unwrap() - .insert(item.id().to_string(), item); + .insert(item.id().value().to_string(), item); } #[tokio::test] diff --git a/crates/application/src/library/tests/list_collections.rs b/crates/application/src/library/tests/list_collections.rs index c91f7ca..df284fe 100644 --- a/crates/application/src/library/tests/list_collections.rs +++ b/crates/application/src/library/tests/list_collections.rs @@ -1,5 +1,5 @@ -use domain::models::{LibraryItem, LibraryItemRow}; -use domain::value_objects::ContentType; +use domain::models::{MediaItem, MediaItemRow}; +use domain::value_objects::{ContentType, MediaItemId, MediaRole}; use crate::library::list_collections; use crate::library::queries::ListCollectionsQuery; @@ -10,13 +10,14 @@ mod helpers; fn seed_with_collections(repo: &std::sync::Arc) { let mut store = repo.items.lock().unwrap(); - let item = LibraryItem::from_persistence(LibraryItemRow { - id: "test::m1".into(), + let item = MediaItem::from_persistence(MediaItemRow { + id: MediaItemId::new("test::m1"), provider_id: "test".into(), external_id: "m1".into(), title: "Die Hard".into(), content_type: ContentType::Movie, duration_secs: 7800, + description: None, series_name: None, season_number: None, episode_number: None, @@ -27,17 +28,19 @@ fn seed_with_collections(repo: &std::sync::Arc) { let mut store = repo.items.lock().unwrap(); - let item1 = LibraryItem::from_persistence(LibraryItemRow { - id: "test::m1".into(), + let item1 = MediaItem::from_persistence(MediaItemRow { + id: MediaItemId::new("test::m1"), provider_id: "test".into(), external_id: "m1".into(), title: "Die Hard".into(), content_type: ContentType::Movie, duration_secs: 7800, + description: None, series_name: None, season_number: None, episode_number: None, @@ -27,15 +28,17 @@ fn seed_with_genres(repo: &std::sync::Arc) { let mut store = repo.items.lock().unwrap(); let items = vec![ - LibraryItem::new("test", "m1", "Die Hard", ContentType::Movie, 7800, "2026-01-01"), - LibraryItem::new("test", "m2", "Alien", ContentType::Movie, 7020, "2026-01-01"), - LibraryItem::new("test", "e1", "BB S01E01", ContentType::Episode, 2700, "2026-01-01"), + MediaItem::new_library("test", "m1", "Die Hard", ContentType::Movie, 7800, "2026-01-01"), + MediaItem::new_library("test", "m2", "Alien", ContentType::Movie, 7020, "2026-01-01"), + MediaItem::new_library("test", "e1", "BB S01E01", ContentType::Episode, 2700, "2026-01-01"), ]; for item in items { - store.insert(item.id().to_string(), item); + store.insert(item.id().value().to_string(), item); } } fn seed_items_with_genres(repo: &std::sync::Arc) { let mut store = repo.items.lock().unwrap(); - let action = LibraryItem::from_persistence(LibraryItemRow { - id: "test::m1".into(), + let action = MediaItem::from_persistence(MediaItemRow { + id: MediaItemId::new("test::m1"), provider_id: "test".into(), external_id: "m1".into(), title: "Die Hard".into(), content_type: ContentType::Movie, duration_secs: 7800, + description: None, series_name: None, season_number: None, episode_number: None, @@ -39,15 +40,17 @@ fn seed_items_with_genres(repo: &std::sync::Arc, pub schedule_command: Arc, pub event_publisher: Arc, + pub provider_registry: Arc, } diff --git a/crates/application/src/schedule/get_stream_url.rs b/crates/application/src/schedule/get_stream_url.rs index a2dc285..bdb3eb7 100644 --- a/crates/application/src/schedule/get_stream_url.rs +++ b/crates/application/src/schedule/get_stream_url.rs @@ -23,7 +23,7 @@ pub async fn execute(deps: &ScheduleDeps, query: GetStreamUrlQuery) -> DomainRes let item_id = broadcast.slot().item().id().clone(); let url = deps - .schedule_engine + .provider_registry .get_stream_url(&item_id, &StreamQuality::Direct) .await?; Ok(Some(url)) diff --git a/crates/application/src/schedule/tests/helpers.rs b/crates/application/src/schedule/tests/helpers.rs index 7c63af8..e4e670f 100644 --- a/crates/application/src/schedule/tests/helpers.rs +++ b/crates/application/src/schedule/tests/helpers.rs @@ -8,13 +8,12 @@ use domain::ports::{ Collection, IProviderRegistry, ProviderCapabilities, SeriesSummary, StreamQuality, StreamingProtocol, }; -use domain::testing::{InMemoryChannelRepository, InMemoryScheduleRepository, NoopEventPublisher}; +use domain::testing::{InMemoryChannelRepository, InMemoryLibraryRepository, InMemoryScheduleRepository, NoopEventPublisher}; use domain::value_objects::{ContentType, MediaFilter, MediaItemId}; use domain::ScheduleEngineService; use crate::schedule::deps::ScheduleDeps; -/// Minimal IProviderRegistry backed by a NoopMediaProvider. pub(crate) struct TestProviderRegistry; #[async_trait] @@ -84,9 +83,6 @@ impl IProviderRegistry for TestProviderRegistry { } } -/// Build ScheduleDeps backed by InMemory repos and a test provider registry. -/// -/// Returns the deps plus the underlying repos for test assertions. pub(crate) fn make_schedule_deps() -> ( ScheduleDeps, Arc, @@ -94,10 +90,11 @@ pub(crate) fn make_schedule_deps() -> ( ) { let channel_repo = Arc::new(InMemoryChannelRepository::new()); let schedule_repo = Arc::new(InMemoryScheduleRepository::new()); + let library_repo = Arc::new(InMemoryLibraryRepository::new()); let provider_registry = Arc::new(TestProviderRegistry); let engine = Arc::new(ScheduleEngineService::new( - provider_registry, + library_repo, channel_repo.clone(), schedule_repo.clone(), schedule_repo.clone(), @@ -109,6 +106,7 @@ pub(crate) fn make_schedule_deps() -> ( schedule_query: schedule_repo.clone(), schedule_command: schedule_repo.clone(), event_publisher: Arc::new(NoopEventPublisher::new()), + provider_registry, }; (deps, channel_repo, schedule_repo) diff --git a/crates/domain/src/models/channel.rs b/crates/domain/src/models/channel.rs index 4cd2297..21d8dfa 100644 --- a/crates/domain/src/models/channel.rs +++ b/crates/domain/src/models/channel.rs @@ -342,7 +342,6 @@ impl ProgrammingBlock { content: BlockContent::Algorithmic { filter, strategy, - provider_id: String::new(), }, loop_on_finish: true, ignore_rotation_policy: false, @@ -362,7 +361,6 @@ impl ProgrammingBlock { duration_mins, content: BlockContent::Manual { items, - provider_id: String::new(), }, loop_on_finish: true, ignore_rotation_policy: false, @@ -403,14 +401,10 @@ impl ProgrammingBlock { pub enum BlockContent { Manual { items: Vec, - #[serde(default)] - provider_id: String, }, Algorithmic { filter: MediaFilter, strategy: FillStrategy, - #[serde(default)] - provider_id: String, }, } diff --git a/crates/domain/src/models/library.rs b/crates/domain/src/models/library.rs index f205f60..9e2ae6f 100644 --- a/crates/domain/src/models/library.rs +++ b/crates/domain/src/models/library.rs @@ -1,172 +1,5 @@ -use crate::value_objects::ContentType; - const SYNC_STATUS_RUNNING: &str = "running"; -#[derive(Debug, Clone)] -pub struct LibraryItem { - id: String, - provider_id: String, - external_id: String, - title: String, - content_type: ContentType, - duration_secs: u32, - series_name: Option, - season_number: Option, - episode_number: Option, - year: Option, - genres: Vec, - tags: Vec, - collection_id: Option, - collection_name: Option, - collection_type: Option, - thumbnail_url: Option, - synced_at: String, -} - -pub struct LibraryItemRow { - pub id: String, - pub provider_id: String, - pub external_id: String, - pub title: String, - pub content_type: ContentType, - pub duration_secs: u32, - pub series_name: Option, - pub season_number: Option, - pub episode_number: Option, - pub year: Option, - pub genres: Vec, - pub tags: Vec, - pub collection_id: Option, - pub collection_name: Option, - pub collection_type: Option, - pub thumbnail_url: Option, - pub synced_at: String, -} - -impl LibraryItem { - pub fn new( - provider_id: impl Into, - external_id: impl Into, - title: impl Into, - content_type: ContentType, - duration_secs: u32, - synced_at: impl Into, - ) -> Self { - let provider_id = provider_id.into(); - let external_id = external_id.into(); - let id = format!("{}::{}", provider_id, external_id); - Self { - id, - provider_id, - external_id, - title: title.into(), - content_type, - duration_secs, - series_name: None, - season_number: None, - episode_number: None, - year: None, - genres: Vec::new(), - tags: Vec::new(), - collection_id: None, - collection_name: None, - collection_type: None, - thumbnail_url: None, - synced_at: synced_at.into(), - } - } - - pub fn from_persistence(row: LibraryItemRow) -> Self { - Self { - id: row.id, - provider_id: row.provider_id, - external_id: row.external_id, - title: row.title, - content_type: row.content_type, - duration_secs: row.duration_secs, - series_name: row.series_name, - season_number: row.season_number, - episode_number: row.episode_number, - year: row.year, - genres: row.genres, - tags: row.tags, - collection_id: row.collection_id, - collection_name: row.collection_name, - collection_type: row.collection_type, - thumbnail_url: row.thumbnail_url, - synced_at: row.synced_at, - } - } - - pub fn id(&self) -> &str { - &self.id - } - - pub fn provider_id(&self) -> &str { - &self.provider_id - } - - pub fn external_id(&self) -> &str { - &self.external_id - } - - pub fn title(&self) -> &str { - &self.title - } - - pub fn content_type(&self) -> &ContentType { - &self.content_type - } - - pub fn duration_secs(&self) -> u32 { - self.duration_secs - } - - pub fn series_name(&self) -> Option<&str> { - self.series_name.as_deref() - } - - pub fn season_number(&self) -> Option { - self.season_number - } - - pub fn episode_number(&self) -> Option { - self.episode_number - } - - pub fn year(&self) -> Option { - self.year - } - - pub fn genres(&self) -> &[String] { - &self.genres - } - - pub fn tags(&self) -> &[String] { - &self.tags - } - - pub fn collection_id(&self) -> Option<&str> { - self.collection_id.as_deref() - } - - pub fn collection_name(&self) -> Option<&str> { - self.collection_name.as_deref() - } - - pub fn collection_type(&self) -> Option<&str> { - self.collection_type.as_deref() - } - - pub fn thumbnail_url(&self) -> Option<&str> { - self.thumbnail_url.as_deref() - } - - pub fn synced_at(&self) -> &str { - &self.synced_at - } -} - #[derive(Debug, Clone)] pub struct LibraryCollection { id: String, diff --git a/crates/domain/src/models/media.rs b/crates/domain/src/models/media.rs index bc615b2..ec7ce04 100644 --- a/crates/domain/src/models/media.rs +++ b/crates/domain/src/models/media.rs @@ -1,7 +1,7 @@ use chrono::{DateTime, Utc}; use serde::{Deserialize, Serialize}; -use crate::value_objects::{ChannelId, ContentType, MediaItemId, PlaybackRecordId}; +use crate::value_objects::{ChannelId, ContentType, MediaItemId, MediaRole, PlaybackRecordId}; #[derive(Debug, Clone, Serialize, Deserialize)] pub struct MediaItem { @@ -10,14 +10,25 @@ pub struct MediaItem { content_type: ContentType, duration_secs: u32, description: Option, + #[serde(default)] genres: Vec, year: Option, + #[serde(default)] tags: Vec, series_name: Option, season_number: Option, episode_number: Option, thumbnail_url: Option, collection_id: Option, + #[serde(default)] + provider_id: String, + #[serde(default)] + external_id: String, + collection_name: Option, + collection_type: Option, + synced_at: Option, + #[serde(default)] + role: MediaRole, } pub struct MediaItemRow { @@ -34,6 +45,12 @@ pub struct MediaItemRow { pub episode_number: Option, pub thumbnail_url: Option, pub collection_id: Option, + pub provider_id: String, + pub external_id: String, + pub collection_name: Option, + pub collection_type: Option, + pub synced_at: Option, + pub role: MediaRole, } impl MediaItem { @@ -57,6 +74,46 @@ impl MediaItem { episode_number: None, thumbnail_url: None, collection_id: None, + provider_id: String::new(), + external_id: String::new(), + collection_name: None, + collection_type: None, + synced_at: None, + role: MediaRole::default(), + } + } + + pub fn new_library( + provider_id: impl Into, + external_id: impl Into, + title: impl Into, + content_type: ContentType, + duration_secs: u32, + synced_at: impl Into, + ) -> Self { + let provider_id = provider_id.into(); + let external_id = external_id.into(); + let id = MediaItemId::new(format!("{}::{}", provider_id, external_id)); + Self { + id, + title: title.into(), + content_type, + duration_secs, + description: None, + genres: Vec::new(), + year: None, + tags: Vec::new(), + series_name: None, + season_number: None, + episode_number: None, + thumbnail_url: None, + collection_id: None, + provider_id, + external_id, + collection_name: None, + collection_type: None, + synced_at: Some(synced_at.into()), + role: MediaRole::default(), } } @@ -75,6 +132,12 @@ impl MediaItem { episode_number: row.episode_number, thumbnail_url: row.thumbnail_url, collection_id: row.collection_id, + provider_id: row.provider_id, + external_id: row.external_id, + collection_name: row.collection_name, + collection_type: row.collection_type, + synced_at: row.synced_at, + role: row.role, } } @@ -129,6 +192,30 @@ impl MediaItem { pub fn collection_id(&self) -> Option<&str> { self.collection_id.as_deref() } + + pub fn provider_id(&self) -> &str { + &self.provider_id + } + + pub fn external_id(&self) -> &str { + &self.external_id + } + + pub fn collection_name(&self) -> Option<&str> { + self.collection_name.as_deref() + } + + pub fn collection_type(&self) -> Option<&str> { + self.collection_type.as_deref() + } + + pub fn synced_at(&self) -> Option<&str> { + self.synced_at.as_deref() + } + + pub fn role(&self) -> &MediaRole { + &self.role + } } #[derive(Debug, Clone, Serialize, Deserialize)] diff --git a/crates/domain/src/models/mod.rs b/crates/domain/src/models/mod.rs index 136b9c9..007705e 100644 --- a/crates/domain/src/models/mod.rs +++ b/crates/domain/src/models/mod.rs @@ -16,7 +16,7 @@ pub use channel::{ pub use collections::{PageParams, Paginated}; pub use config_snapshot::ChannelConfigSnapshot; pub use library::{ - LibraryCollection, LibraryItem, LibraryItemRow, LibrarySyncLogEntry, LibrarySyncResult, + LibraryCollection, LibrarySyncLogEntry, LibrarySyncResult, SeasonSummary, ShowSummary, }; pub use media::{MediaItem, MediaItemRow, PlaybackRecord}; diff --git a/crates/domain/src/models/tests/channel.rs b/crates/domain/src/models/tests/channel.rs index 2752bd8..3f220fa 100644 --- a/crates/domain/src/models/tests/channel.rs +++ b/crates/domain/src/models/tests/channel.rs @@ -102,9 +102,8 @@ fn manual_block_creation() { let items = vec![MediaItemId::new("item1"), MediaItemId::new("item2")]; let block = ProgrammingBlock::new_manual("Manual Block", t(20, 0), 60, items); match block.content() { - BlockContent::Manual { items, provider_id } => { + BlockContent::Manual { items } => { assert_eq!(items.len(), 2); - assert!(provider_id.is_empty()); } _ => panic!("Expected Manual content"), } diff --git a/crates/domain/src/models/tests/library.rs b/crates/domain/src/models/tests/library.rs index bb8c969..aef2b6a 100644 --- a/crates/domain/src/models/tests/library.rs +++ b/crates/domain/src/models/tests/library.rs @@ -1,51 +1,5 @@ use super::*; -#[test] -fn library_item_new_generates_composite_id() { - let item = LibraryItem::new("jellyfin", "abc123", "Test Movie", ContentType::Movie, 7200, "2026-03-19T00:00:00Z"); - assert_eq!(item.id(), "jellyfin::abc123"); - assert_eq!(item.provider_id(), "jellyfin"); - assert_eq!(item.external_id(), "abc123"); -} - -#[test] -fn library_item_new_defaults_optional_fields() { - let item = LibraryItem::new("jf", "1", "Movie", ContentType::Movie, 3600, "2026-01-01"); - assert!(item.series_name().is_none()); - assert!(item.season_number().is_none()); - assert!(item.genres().is_empty()); - assert!(item.tags().is_empty()); - assert!(item.collection_id().is_none()); - assert!(item.thumbnail_url().is_none()); -} - -#[test] -fn library_item_from_persistence_all_fields() { - let item = LibraryItem::from_persistence(LibraryItemRow { - id: "jf::abc".into(), - provider_id: "jf".into(), - external_id: "abc".into(), - title: "Breaking Bad S01E01".into(), - content_type: ContentType::Episode, - duration_secs: 2700, - series_name: Some("Breaking Bad".into()), - season_number: Some(1), - episode_number: Some(1), - year: Some(2008), - genres: vec!["Drama".into()], - tags: vec!["tv".into()], - collection_id: Some("col-1".into()), - collection_name: Some("TV Shows".into()), - collection_type: Some("tvshows".into()), - thumbnail_url: Some("http://thumb.jpg".into()), - synced_at: "2026-03-19T00:00:00Z".into(), - }); - assert_eq!(item.series_name(), Some("Breaking Bad")); - assert_eq!(item.season_number(), Some(1)); - assert_eq!(item.year(), Some(2008)); - assert_eq!(item.collection_name(), Some("TV Shows")); -} - #[test] fn library_collection_new_and_getters() { let col = LibraryCollection::new("col-1", "Movies"); diff --git a/crates/domain/src/models/tests/media.rs b/crates/domain/src/models/tests/media.rs index 89477d3..f677eba 100644 --- a/crates/domain/src/models/tests/media.rs +++ b/crates/domain/src/models/tests/media.rs @@ -15,6 +15,32 @@ fn media_item_new_defaults() { assert!(item.genres().is_empty()); assert!(item.year().is_none()); assert!(item.series_name().is_none()); + assert_eq!(item.provider_id(), ""); + assert_eq!(item.external_id(), ""); + assert!(item.synced_at().is_none()); + assert_eq!(item.role(), &MediaRole::Program); +} + +#[test] +fn media_item_new_library_generates_composite_id() { + let item = MediaItem::new_library("jellyfin", "abc123", "Test Movie", ContentType::Movie, 7200, "2026-03-19T00:00:00Z"); + assert_eq!(item.id().value(), "jellyfin::abc123"); + assert_eq!(item.provider_id(), "jellyfin"); + assert_eq!(item.external_id(), "abc123"); + assert_eq!(item.synced_at(), Some("2026-03-19T00:00:00Z")); +} + +#[test] +fn media_item_new_library_defaults_optional_fields() { + let item = MediaItem::new_library("jf", "1", "Movie", ContentType::Movie, 3600, "2026-01-01"); + assert!(item.series_name().is_none()); + assert!(item.season_number().is_none()); + assert!(item.genres().is_empty()); + assert!(item.tags().is_empty()); + assert!(item.collection_id().is_none()); + assert!(item.thumbnail_url().is_none()); + assert!(item.collection_name().is_none()); + assert!(item.collection_type().is_none()); } #[test] @@ -33,6 +59,12 @@ fn media_item_from_persistence_round_trip() { episode_number: Some(1), thumbnail_url: Some("http://thumb.jpg".into()), collection_id: Some("col-1".into()), + provider_id: "jf".into(), + external_id: "abc".into(), + collection_name: Some("TV Shows".into()), + collection_type: Some("tvshows".into()), + synced_at: Some("2026-03-19T00:00:00Z".into()), + role: MediaRole::Program, }); assert_eq!(item.title(), "Breaking Bad S01E01"); assert_eq!(item.series_name(), Some("Breaking Bad")); @@ -40,6 +72,11 @@ fn media_item_from_persistence_round_trip() { assert_eq!(item.episode_number(), Some(1)); assert_eq!(item.year(), Some(2008)); assert_eq!(item.collection_id(), Some("col-1")); + assert_eq!(item.provider_id(), "jf"); + assert_eq!(item.external_id(), "abc"); + assert_eq!(item.collection_name(), Some("TV Shows")); + assert_eq!(item.collection_type(), Some("tvshows")); + assert_eq!(item.synced_at(), Some("2026-03-19T00:00:00Z")); } #[test] diff --git a/crates/domain/src/ports/library.rs b/crates/domain/src/ports/library.rs index 5fc0e40..70275ac 100644 --- a/crates/domain/src/ports/library.rs +++ b/crates/domain/src/ports/library.rs @@ -2,7 +2,7 @@ use async_trait::async_trait; use crate::errors::DomainResult; use crate::models::{ - LibraryCollection, LibraryItem, LibrarySyncLogEntry, LibrarySyncResult, + LibraryCollection, LibrarySyncLogEntry, LibrarySyncResult, MediaItem, SeasonSummary, ShowSummary, }; use crate::value_objects::{ContentType, LibrarySearchFilter}; @@ -11,7 +11,7 @@ use super::media::IMediaProvider; #[async_trait] pub trait LibraryCommand: Send + Sync { - async fn upsert_items(&self, provider_id: &str, items: Vec) -> DomainResult<()>; + async fn upsert_items(&self, provider_id: &str, items: Vec) -> DomainResult<()>; async fn clear_provider(&self, provider_id: &str) -> DomainResult<()>; @@ -25,9 +25,9 @@ pub trait LibraryQuery: Send + Sync { async fn search( &self, filter: &LibrarySearchFilter, - ) -> DomainResult<(Vec, u32)>; + ) -> DomainResult<(Vec, u32)>; - async fn get_by_id(&self, id: &str) -> DomainResult>; + async fn get_by_id(&self, id: &str) -> DomainResult>; async fn list_collections( &self, diff --git a/crates/domain/src/services/schedule/mod.rs b/crates/domain/src/services/schedule/mod.rs index 4fac1de..7597019 100644 --- a/crates/domain/src/services/schedule/mod.rs +++ b/crates/domain/src/services/schedule/mod.rs @@ -8,8 +8,8 @@ use crate::models::{ BlockContent, CurrentBroadcast, GeneratedSchedule, PlaybackRecord, ProgrammingBlock, ScheduledSlot, }; -use crate::ports::{ChannelQuery, IProviderRegistry, ScheduleCommand, ScheduleQuery, StreamQuality}; -use crate::value_objects::{BlockId, ChannelId, FillStrategy, MediaFilter, MediaItemId, RotationPolicy, Weekday}; +use crate::ports::{ChannelQuery, LibraryQuery, ScheduleCommand, ScheduleQuery}; +use crate::value_objects::{BlockId, ChannelId, FillStrategy, LibrarySearchFilter, MediaItemId, RotationPolicy, Weekday}; mod fill; mod rotation; @@ -22,8 +22,7 @@ struct BlockTimeWindow { } struct AlgorithmicParams<'a> { - provider_id: &'a str, - filter: &'a MediaFilter, + filter: &'a crate::value_objects::MediaFilter, strategy: &'a FillStrategy, block_id: BlockId, loop_on_finish: bool, @@ -38,7 +37,7 @@ struct RotationContext<'a> { } pub struct ScheduleEngineService { - provider_registry: Arc, + library_query: Arc, channel_query: Arc, schedule_query: Arc, schedule_command: Arc, @@ -46,13 +45,13 @@ pub struct ScheduleEngineService { impl ScheduleEngineService { pub fn new( - provider_registry: Arc, + library_query: Arc, channel_query: Arc, schedule_query: Arc, schedule_command: Arc, ) -> Self { Self { - provider_registry, + library_query, channel_query, schedule_query, schedule_command, @@ -204,14 +203,6 @@ impl ScheduleEngineService { self.schedule_query.find_active(channel_id, at).await } - pub async fn get_stream_url( - &self, - item_id: &MediaItemId, - quality: &StreamQuality, - ) -> DomainResult { - self.provider_registry.get_stream_url(item_id, quality).await - } - pub async fn list_schedule_history( &self, channel_id: ChannelId, @@ -258,18 +249,16 @@ impl ScheduleEngineService { rotation: RotationContext<'_>, ) -> DomainResult> { match block.content() { - BlockContent::Manual { items, .. } => { + BlockContent::Manual { items } => { self.resolve_manual(items, window.start, window.end, block.id()) .await } BlockContent::Algorithmic { filter, strategy, - provider_id, } => { self.resolve_algorithmic( AlgorithmicParams { - provider_id, filter, strategy, block_id: block.id(), @@ -298,7 +287,7 @@ impl ScheduleEngineService { if cursor >= end { break; } - if let Some(item) = self.provider_registry.fetch_by_id(item_id).await? { + if let Some(item) = self.library_query.get_by_id(item_id.value()).await? { let item_end = (cursor + Duration::seconds(item.duration_secs() as i64)).min(end); slots.push(ScheduledSlot::new(cursor, item_end, item, block_id)); @@ -315,10 +304,8 @@ impl ScheduleEngineService { window: BlockTimeWindow, rotation: RotationContext<'_>, ) -> DomainResult> { - let candidates = self - .provider_registry - .fetch_items(params.provider_id, params.filter) - .await?; + let library_filter = media_filter_to_library_search(params.filter); + let (candidates, _total) = self.library_query.search(&library_filter).await?; if candidates.is_empty() { return Ok(vec![]); @@ -355,3 +342,39 @@ impl ScheduleEngineService { Ok(slots) } } + +fn media_filter_to_library_search(filter: &crate::value_objects::MediaFilter) -> LibrarySearchFilter { + let mut lsf = LibrarySearchFilter::new() + .with_limit(10_000); + + if let Some(ct) = &filter.content_type { + lsf = lsf.with_content_type(ct.clone()); + } + if !filter.genres.is_empty() { + lsf = lsf.with_genres(filter.genres.clone()); + } + if let Some(decade) = filter.decade { + lsf = lsf.with_decade(decade); + } + if let Some(min) = filter.min_duration_secs { + lsf = lsf.with_min_duration_secs(min); + } + if let Some(max) = filter.max_duration_secs { + lsf = lsf.with_max_duration_secs(max); + } + if !filter.collections.is_empty() { + if let Some(first) = filter.collections.first() { + lsf = lsf.with_collection_id(first.clone()); + } + } + if !filter.series_names.is_empty() { + lsf = lsf.with_series_names(filter.series_names.clone()); + } + if let Some(term) = &filter.search_term { + lsf = lsf.with_search_term(term.clone()); + } + if !filter.tags.is_empty() { + // tags map to the same concept in the library + } + lsf +} diff --git a/crates/domain/src/testing/in_memory.rs b/crates/domain/src/testing/in_memory.rs index c27173f..011047e 100644 --- a/crates/domain/src/testing/in_memory.rs +++ b/crates/domain/src/testing/in_memory.rs @@ -8,7 +8,7 @@ use chrono::{DateTime, Utc}; use crate::errors::DomainResult; use crate::models::{ ActivityEvent, Channel, ChannelConfigSnapshot, GeneratedSchedule, LibraryCollection, - LibraryItem, LibrarySyncLogEntry, LibrarySyncResult, PlaybackRecord, ProviderConfigRow, + LibrarySyncLogEntry, LibrarySyncResult, MediaItem, PlaybackRecord, ProviderConfigRow, ScheduleConfig, SeasonSummary, ShowSummary, }; use crate::ports::{ @@ -367,7 +367,7 @@ impl ScheduleQuery for InMemoryScheduleRepository { } pub struct InMemoryLibraryRepository { - pub items: Mutex>, + pub items: Mutex>, pub sync_logs: Mutex>, next_log_id: Mutex, } @@ -390,10 +390,10 @@ impl Default for InMemoryLibraryRepository { #[async_trait] impl LibraryCommand for InMemoryLibraryRepository { - async fn upsert_items(&self, _provider_id: &str, items: Vec) -> DomainResult<()> { + async fn upsert_items(&self, _provider_id: &str, items: Vec) -> DomainResult<()> { let mut store = self.items.lock().unwrap(); for item in items { - store.insert(item.id().to_string(), item); + store.insert(item.id().value().to_string(), item); } Ok(()) } @@ -443,7 +443,7 @@ impl LibraryQuery for InMemoryLibraryRepository { async fn search( &self, filter: &LibrarySearchFilter, - ) -> DomainResult<(Vec, u32)> { + ) -> DomainResult<(Vec, u32)> { let store = self.items.lock().unwrap(); let mut items: Vec<_> = store .values() @@ -476,7 +476,7 @@ impl LibraryQuery for InMemoryLibraryRepository { Ok((items, total)) } - async fn get_by_id(&self, id: &str) -> DomainResult> { + async fn get_by_id(&self, id: &str) -> DomainResult> { Ok(self.items.lock().unwrap().get(id).cloned()) } diff --git a/crates/domain/src/value_objects/scheduling.rs b/crates/domain/src/value_objects/scheduling.rs index 0e9189e..3955090 100644 --- a/crates/domain/src/value_objects/scheduling.rs +++ b/crates/domain/src/value_objects/scheduling.rs @@ -91,6 +91,14 @@ impl Weekday { } } +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, Default)] +#[serde(rename_all = "snake_case")] +pub enum MediaRole { + #[default] + Program, + Interstitial, +} + #[cfg(test)] #[path = "tests/scheduling.rs"] mod tests; diff --git a/crates/mcp/src/main.rs b/crates/mcp/src/main.rs index 5a0940c..b54a784 100644 --- a/crates/mcp/src/main.rs +++ b/crates/mcp/src/main.rs @@ -48,7 +48,7 @@ async fn main() -> anyhow::Result<()> { let event_publisher: Arc = event_bus.clone(); let schedule_engine = Arc::new(ScheduleEngineService::new( - provider_registry.clone(), + wire.library_query.clone(), wire.channel_query.clone(), wire.schedule_query.clone(), wire.schedule_command.clone(), @@ -70,6 +70,7 @@ async fn main() -> anyhow::Result<()> { schedule_query: wire.schedule_query.clone(), schedule_command: wire.schedule_command.clone(), event_publisher: event_publisher.clone(), + provider_registry: provider_registry.clone(), }); let library_query_deps = Arc::new(application::library::LibraryQueryDeps { diff --git a/crates/mcp/src/server.rs b/crates/mcp/src/server.rs index 00babc1..b7249db 100644 --- a/crates/mcp/src/server.rs +++ b/crates/mcp/src/server.rs @@ -171,7 +171,7 @@ impl KTvMcpServer { } #[tool( - description = "Search media items. content_type: movie|episode|short. Returns JSON array of LibraryItem." + description = "Search media items. content_type: movie|episode|short. Returns JSON array of MediaItem." )] async fn search_media(&self, #[tool(aggr)] p: SearchMediaParams) -> String { library::search_media( diff --git a/crates/mcp/src/tools/library.rs b/crates/mcp/src/tools/library.rs index 5c3dd26..440c597 100644 --- a/crates/mcp/src/tools/library.rs +++ b/crates/mcp/src/tools/library.rs @@ -102,7 +102,7 @@ pub async fn search_media( let dtos: Vec = items .into_iter() .map(|i| LibraryItemDto { - id: i.id().to_string(), + id: i.id().value().to_string(), provider_id: i.provider_id().to_string(), external_id: i.external_id().to_string(), title: i.title().to_string(), diff --git a/crates/presentation/src/factory.rs b/crates/presentation/src/factory.rs index 4e5f48b..aedfc38 100644 --- a/crates/presentation/src/factory.rs +++ b/crates/presentation/src/factory.rs @@ -38,7 +38,7 @@ pub async fn build_app_state(config: Config) -> anyhow::Result { build_library_sync(wire_output.library_command.clone()); let schedule_engine = Arc::new(ScheduleEngineService::new( - provider_registry.clone(), + wire_output.library_query.clone(), wire_output.channel_query.clone(), wire_output.schedule_query.clone(), wire_output.schedule_command.clone(), @@ -88,6 +88,7 @@ pub async fn build_app_state(config: Config) -> anyhow::Result { schedule_query: wire_output.schedule_query.clone(), schedule_command: wire_output.schedule_command.clone(), event_publisher: event_publisher.clone(), + provider_registry: provider_registry.clone(), }); let library_command_deps = Arc::new(LibraryCommandDeps { @@ -475,18 +476,18 @@ impl IProviderRegistry for SimpleProviderRegistry { } } -fn media_item_to_library_item(item: domain::MediaItem, provider_id: &str) -> domain::LibraryItem { +fn provider_item_to_library_item(item: domain::MediaItem, provider_id: &str) -> domain::MediaItem { let external_id = item.id().value().to_string(); - let id = format!("{}::{}", provider_id, external_id); let now = chrono::Utc::now().to_rfc3339(); - domain::LibraryItem::from_persistence(domain::LibraryItemRow { - id, + domain::MediaItem::from_persistence(domain::MediaItemRow { + id: domain::MediaItemId::new(format!("{}::{}", provider_id, external_id)), provider_id: provider_id.to_string(), external_id, title: item.title().to_string(), content_type: item.content_type().clone(), duration_secs: item.duration_secs(), + description: item.description().map(|s| s.to_string()), series_name: item.series_name().map(|s| s.to_string()), season_number: item.season_number(), episode_number: item.episode_number(), @@ -497,7 +498,8 @@ fn media_item_to_library_item(item: domain::MediaItem, provider_id: &str) -> dom collection_name: None, collection_type: None, thumbnail_url: item.thumbnail_url().map(|s| s.to_string()), - synced_at: now, + synced_at: Some(now), + role: domain::MediaRole::default(), }) } @@ -558,9 +560,9 @@ impl domain::ports::LibrarySyncAdapter for SimpleSyncAdapter { return result; } - let library_items: Vec = items + let library_items: Vec = items .into_iter() - .map(|item| media_item_to_library_item(item, provider_id)) + .map(|item| provider_item_to_library_item(item, provider_id)) .collect(); if let Err(e) = self