diff --git a/crates/domain/src/models/channel.rs b/crates/domain/src/models/channel.rs index 7fd5ccc..40923ab 100644 --- a/crates/domain/src/models/channel.rs +++ b/crates/domain/src/models/channel.rs @@ -282,6 +282,13 @@ impl ScheduleConfig { &self.day_blocks } + pub fn find_block_mut(&mut self, block_id: BlockId) -> Option<&mut ProgrammingBlock> { + self.day_blocks + .values_mut() + .flatten() + .find(|b| b.id() == block_id) + } + pub fn insert_day(&mut self, day: Weekday, blocks: Vec) { self.day_blocks.insert(day, blocks); } @@ -436,6 +443,14 @@ impl ProgrammingBlock { self.mid_roll_rule = Some(rule); self } + + pub fn set_interstitial_rule(&mut self, rule: Option) { + self.interstitial_rule = rule; + } + + pub fn set_mid_roll_rule(&mut self, rule: Option) { + self.mid_roll_rule = rule; + } } #[derive(Debug, Clone, Serialize, Deserialize)] diff --git a/crates/domain/src/services/schedule/mod.rs b/crates/domain/src/services/schedule/mod.rs index defe15f..b27227a 100644 --- a/crates/domain/src/services/schedule/mod.rs +++ b/crates/domain/src/services/schedule/mod.rs @@ -180,6 +180,211 @@ impl ScheduleEngineService { Ok(schedule) } + pub async fn preview_schedule( + &self, + channel_id: ChannelId, + from: DateTime, + duration_hours: u32, + ) -> DomainResult { + let channel = self + .channel_query + .find_by_id(channel_id) + .await? + .ok_or(DomainError::ChannelNotFound(channel_id))?; + + let tz: Tz = channel + .timezone() + .parse() + .map_err(|_| DomainError::TimezoneError(channel.timezone().to_owned()))?; + + let history = self + .schedule_query + .find_playback_history(channel_id) + .await?; + let latest_schedule = self.schedule_query.find_latest(channel_id).await?; + let generation = latest_schedule + .as_ref() + .map(|s| s.generation() + 1) + .unwrap_or(1); + + let mut block_continuity = self + .schedule_query + .find_last_slot_per_block(channel_id) + .await?; + + let valid_from = from; + let valid_until = from + Duration::hours(duration_hours as i64); + let start_date = from.with_timezone(&tz).date_naive(); + let end_date = valid_until.with_timezone(&tz).date_naive(); + + let mut slots: Vec = Vec::new(); + let mut current_date = start_date; + + while current_date <= end_date { + let weekday = Weekday::from(current_date.weekday()); + + for block in channel.schedule_config().blocks_for(weekday) { + let naive_start = current_date.and_time(block.start_time()); + + let block_start_utc = match tz.from_local_datetime(&naive_start).earliest() { + Some(dt) => dt.with_timezone(&Utc), + None => continue, + }; + + let block_end_utc = + block_start_utc + Duration::minutes(block.duration_mins() as i64); + + let slot_start = block_start_utc.max(valid_from); + let slot_end = block_end_utc.min(valid_until); + + if slot_end <= slot_start { + continue; + } + + let last_item_id = block_continuity.get(&block.id()); + + let mut block_slots = self + .resolve_block( + block, + BlockTimeWindow { + start: slot_start, + end: slot_end, + }, + RotationContext { + history: &history, + policy: channel.rotation_policy(), + generation, + last_item_id, + }, + ) + .await?; + + if let Some(last_slot) = block_slots.last() { + block_continuity.insert(block.id(), last_slot.item().id().clone()); + } + + slots.append(&mut block_slots); + } + + current_date = current_date.succ_opt().ok_or_else(|| { + DomainError::validation("Date overflow during schedule generation") + })?; + } + + slots.sort_by_key(|s| s.start_at()); + + Ok(GeneratedSchedule::new( + channel_id, + valid_from, + valid_until, + generation, + slots, + )) + } + + pub async fn preview_config( + &self, + channel_id: ChannelId, + config: &crate::models::ScheduleConfig, + from: DateTime, + duration_hours: u32, + ) -> DomainResult { + let channel = self + .channel_query + .find_by_id(channel_id) + .await? + .ok_or(DomainError::ChannelNotFound(channel_id))?; + + let tz: Tz = channel + .timezone() + .parse() + .map_err(|_| DomainError::TimezoneError(channel.timezone().to_owned()))?; + + let history = self + .schedule_query + .find_playback_history(channel_id) + .await?; + let latest_schedule = self.schedule_query.find_latest(channel_id).await?; + let generation = latest_schedule + .as_ref() + .map(|s| s.generation() + 1) + .unwrap_or(1); + + let mut block_continuity = self + .schedule_query + .find_last_slot_per_block(channel_id) + .await?; + + let valid_from = from; + let valid_until = from + Duration::hours(duration_hours as i64); + let start_date = from.with_timezone(&tz).date_naive(); + let end_date = valid_until.with_timezone(&tz).date_naive(); + + let mut slots: Vec = Vec::new(); + let mut current_date = start_date; + + while current_date <= end_date { + let weekday = Weekday::from(current_date.weekday()); + + for block in config.blocks_for(weekday) { + let naive_start = current_date.and_time(block.start_time()); + + let block_start_utc = match tz.from_local_datetime(&naive_start).earliest() { + Some(dt) => dt.with_timezone(&Utc), + None => continue, + }; + + let block_end_utc = + block_start_utc + Duration::minutes(block.duration_mins() as i64); + + let slot_start = block_start_utc.max(valid_from); + let slot_end = block_end_utc.min(valid_until); + + if slot_end <= slot_start { + continue; + } + + let last_item_id = block_continuity.get(&block.id()); + + let mut block_slots = self + .resolve_block( + block, + BlockTimeWindow { + start: slot_start, + end: slot_end, + }, + RotationContext { + history: &history, + policy: channel.rotation_policy(), + generation, + last_item_id, + }, + ) + .await?; + + if let Some(last_slot) = block_slots.last() { + block_continuity.insert(block.id(), last_slot.item().id().clone()); + } + + slots.append(&mut block_slots); + } + + current_date = current_date.succ_opt().ok_or_else(|| { + DomainError::validation("Date overflow during schedule generation") + })?; + } + + slots.sort_by_key(|s| s.start_at()); + + Ok(GeneratedSchedule::new( + channel_id, + valid_from, + valid_until, + generation, + slots, + )) + } + pub fn get_current_broadcast( schedule: &GeneratedSchedule, now: DateTime, diff --git a/crates/mcp/src/server.rs b/crates/mcp/src/server.rs index fe1af8b..6f5cf03 100644 --- a/crates/mcp/src/server.rs +++ b/crates/mcp/src/server.rs @@ -74,6 +74,83 @@ pub struct ListGenresParams { pub content_type: Option, } +#[derive(Debug, Deserialize, JsonSchema)] +pub struct BrowseLibraryParams { + /// Filter by content type: movie, episode, short + pub content_type: Option, + /// Filter by genre names + pub genres: Option>, + /// Full-text search term + pub search_term: Option, + /// Filter by series names + pub series_names: Option>, + /// Filter by collection IDs + pub collections: Option>, + /// Filter by decade (e.g. 1990) + pub decade: Option, + /// Filter by role: program or interstitial + pub role: Option, + /// Minimum duration in seconds + pub min_duration_secs: Option, + /// Maximum duration in seconds + pub max_duration_secs: Option, + /// Max items to return (default 50) + pub limit: Option, +} + +#[derive(Debug, Deserialize, JsonSchema)] +pub struct PreviewScheduleParams { + pub channel_id: String, + /// Hours to preview (default 24) + pub duration_hours: Option, +} + +#[derive(Debug, Deserialize, JsonSchema)] +pub struct PreviewConfigParams { + pub channel_id: String, + /// Full ScheduleConfig as JSON + pub schedule_config_json: String, + /// Hours to preview (default 24) + pub duration_hours: Option, +} + +#[derive(Debug, Deserialize, JsonSchema)] +pub struct SuggestScheduleParams { + /// Desired genres for the block + pub genres: Option>, + /// Content type: movie, episode, short + pub content_type: Option, + /// Name for the programming block + pub block_name: Option, + /// Start time in HH:MM format (default "20:00") + pub start_time: Option, + /// Block duration in minutes (default 180) + pub duration_mins: Option, +} + +#[derive(Debug, Deserialize, JsonSchema)] +pub struct SetInterstitialRuleParams { + pub channel_id: String, + pub block_id: String, + /// InterstitialRule as JSON. Omit or null to clear the rule. + pub rule_json: Option, +} + +#[derive(Debug, Deserialize, JsonSchema)] +pub struct SetMidRollRuleParams { + pub channel_id: String, + pub block_id: String, + /// MidRollRule as JSON. Omit or null to clear the rule. + pub rule_json: Option, +} + +#[derive(Debug, Deserialize, JsonSchema)] +pub struct SetGapFillerParams { + pub channel_id: String, + /// MediaFilter as JSON. Omit or null to clear gap filler. + pub filter_json: Option, +} + fn parse_uuid(s: &str) -> Result { s.parse::() .map_err(|_| serde_json::json!({"error": format!("invalid UUID: {s}")}).to_string()) @@ -201,6 +278,151 @@ impl KTvMcpServer { ) .await } + + #[tool( + description = "Browse the library with rich filters. Supports genre, decade, series, content type, role (program/interstitial), duration range. Returns concise summaries." + )] + async fn browse_library(&self, #[tool(aggr)] p: BrowseLibraryParams) -> String { + library::browse_library( + &self.library_command_deps, + library::BrowseParams { + content_type: p.content_type, + genres: p.genres.unwrap_or_default(), + search_term: p.search_term, + series_names: p.series_names.unwrap_or_default(), + collections: p.collections.unwrap_or_default(), + decade: p.decade, + role: p.role, + min_duration_secs: p.min_duration_secs, + max_duration_secs: p.max_duration_secs, + limit: p.limit, + }, + ) + .await + } + + #[tool( + description = "Get aggregate library statistics: total items, items by genre/content type/role, series with episode counts, recently synced items, total duration." + )] + async fn library_stats(&self) -> String { + library::library_stats(&self.library_query, &self.library_command_deps).await + } + + #[tool( + description = "Analyze a channel's schedule: most-played items, genre distribution, rotation coverage, upcoming gaps." + )] + async fn analyze_schedule(&self, #[tool(aggr)] p: ChannelIdParam) -> String { + match parse_uuid(&p.channel_id) { + Ok(id) => { + schedule::analyze_schedule(&self.channel_query, &self.schedule_query, id).await + } + Err(e) => e, + } + } + + #[tool( + description = "Dry-run schedule generation: preview the next N hours of slots without persisting. Uses the channel's current config." + )] + async fn preview_schedule(&self, #[tool(aggr)] p: PreviewScheduleParams) -> String { + match parse_uuid(&p.channel_id) { + Ok(id) => schedule::preview_schedule(&self.schedule_deps, id, p.duration_hours).await, + Err(e) => e, + } + } + + #[tool( + description = "Preview a ScheduleConfig without applying it: provide schedule_config_json and see what slots would be generated. Does not persist." + )] + async fn preview_config(&self, #[tool(aggr)] p: PreviewConfigParams) -> String { + match parse_uuid(&p.channel_id) { + Ok(id) => { + schedule::preview_config( + &self.schedule_deps, + id, + &p.schedule_config_json, + p.duration_hours, + ) + .await + } + Err(e) => e, + } + } + + #[tool( + description = "Suggest a ScheduleConfig block based on constraints (genres, content type, time). Returns a concrete block config." + )] + async fn suggest_schedule(&self, #[tool(aggr)] p: SuggestScheduleParams) -> String { + schedule::suggest_schedule( + &self.library_query, + p.genres.unwrap_or_default(), + p.content_type, + p.block_name, + p.start_time, + p.duration_mins, + ) + .await + } + + #[tool( + description = "Set or clear an interstitial rule on a programming block. Pass rule_json with {pool_filter, strategy, min_gap_secs} or omit to clear." + )] + async fn set_interstitial_rule(&self, #[tool(aggr)] p: SetInterstitialRuleParams) -> String { + let channel_id = match parse_uuid(&p.channel_id) { + Ok(id) => id, + Err(e) => return e, + }; + let block_id = match parse_uuid(&p.block_id) { + Ok(id) => id, + Err(e) => return e, + }; + schedule::set_interstitial_rule( + &self.channel_cmd_deps, + channel_id, + self.owner_id, + block_id, + p.rule_json, + ) + .await + } + + #[tool( + description = "Set or clear a mid-roll break rule on a programming block. Pass rule_json with {prefer_chapters, fallback_interval_mins, break_duration_secs, pool_filter} or omit to clear." + )] + async fn set_mid_roll_rule(&self, #[tool(aggr)] p: SetMidRollRuleParams) -> String { + let channel_id = match parse_uuid(&p.channel_id) { + Ok(id) => id, + Err(e) => return e, + }; + let block_id = match parse_uuid(&p.block_id) { + Ok(id) => id, + Err(e) => return e, + }; + schedule::set_mid_roll_rule( + &self.channel_cmd_deps, + channel_id, + self.owner_id, + block_id, + p.rule_json, + ) + .await + } + + #[tool( + description = "Set or clear the gap filler on a channel. Pass filter_json with a MediaFilter or omit to clear." + )] + async fn set_gap_filler(&self, #[tool(aggr)] p: SetGapFillerParams) -> String { + let channel_id = match parse_uuid(&p.channel_id) { + Ok(id) => id, + Err(e) => return e, + }; + schedule::set_gap_filler( + &self.channel_cmd_deps, + channel_id, + self.owner_id, + p.filter_json, + ) + .await + } } #[tool(tool_box)] @@ -214,7 +436,8 @@ impl ServerHandler for KTvMcpServer { version: env!("CARGO_PKG_VERSION").into(), }, instructions: Some( - "K-TV MCP server. Create channels, define programming blocks, generate schedules. \ + "K-TV MCP server — creative programming director for linear TV channels. \ + Browse the library, analyze schedules, preview configs, configure interstitials and gap fillers. \ All operations run as the user configured via MCP_USER_ID." .into(), ), diff --git a/crates/mcp/src/tools/library.rs b/crates/mcp/src/tools/library.rs index 281019b..5dbd979 100644 --- a/crates/mcp/src/tools/library.rs +++ b/crates/mcp/src/tools/library.rs @@ -1,3 +1,4 @@ +use std::collections::HashMap; use std::sync::Arc; use application::library::SearchItemsQuery; @@ -6,6 +7,7 @@ use serde::Serialize; use crate::error::{domain_err, ok_json}; const DEFAULT_SEARCH_LIMIT: u32 = 50; +const STATS_SEARCH_LIMIT: u32 = 10_000; #[derive(Serialize)] struct CollectionDto { @@ -30,6 +32,7 @@ struct LibraryItemDto { tags: Vec, collection_id: Option, thumbnail_url: Option, + role: String, } #[derive(Serialize)] @@ -38,6 +41,59 @@ struct SearchResult { total: u32, } +#[derive(Serialize)] +struct BrowseResult { + items: Vec, + total: u32, + summary: String, +} + +#[derive(Serialize)] +struct BrowseItemDto { + id: String, + title: String, + content_type: String, + duration_mins: u32, + series_name: Option, + season_episode: Option, + year: Option, + genres: Vec, + role: String, +} + +#[derive(Serialize)] +struct LibraryStats { + total_items: u32, + by_content_type: HashMap, + by_role: HashMap, + genres: Vec, + series: Vec, + total_duration_hours: f64, + recently_synced: Vec, +} + +#[derive(Serialize)] +struct GenreStat { + genre: String, + count: u32, +} + +#[derive(Serialize)] +struct SeriesStat { + name: String, + episode_count: u32, + season_count: u32, + genres: Vec, +} + +#[derive(Serialize)] +struct RecentSyncDto { + provider_id: String, + started_at: String, + status: String, + items_found: u32, +} + fn content_type_to_str(ct: &domain::ContentType) -> &'static str { match ct { domain::ContentType::Movie => "movie", @@ -46,6 +102,33 @@ fn content_type_to_str(ct: &domain::ContentType) -> &'static str { } } +fn role_to_str(r: &domain::MediaRole) -> &'static str { + match r { + domain::MediaRole::Program => "program", + domain::MediaRole::Interstitial => "interstitial", + } +} + +fn item_to_dto(i: &domain::MediaItem) -> LibraryItemDto { + LibraryItemDto { + 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(), + content_type: content_type_to_str(i.content_type()).to_string(), + duration_secs: i.duration_secs(), + series_name: i.series_name().map(|s| s.to_string()), + season_number: i.season_number(), + episode_number: i.episode_number(), + year: i.year(), + genres: i.genres().to_vec(), + tags: i.tags().to_vec(), + collection_id: i.collection_id().map(|s| s.to_string()), + thumbnail_url: i.thumbnail_url().map(|s| s.to_string()), + role: role_to_str(i.role()).to_string(), + } +} + pub async fn list_collections(library_query: &Arc) -> String { match library_query.list_collections(None).await { Ok(cols) => { @@ -103,27 +186,209 @@ pub async fn search_media( }; match application::library::search::execute(library_command_deps, query).await { Ok((items, total)) => { - let dtos: Vec = items - .into_iter() - .map(|i| LibraryItemDto { - 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(), - content_type: content_type_to_str(i.content_type()).to_string(), - duration_secs: i.duration_secs(), - series_name: i.series_name().map(|s| s.to_string()), - season_number: i.season_number(), - episode_number: i.episode_number(), - year: i.year(), - genres: i.genres().to_vec(), - tags: i.tags().to_vec(), - collection_id: i.collection_id().map(|s| s.to_string()), - thumbnail_url: i.thumbnail_url().map(|s| s.to_string()), - }) - .collect(); + let dtos: Vec = items.iter().map(item_to_dto).collect(); ok_json(&SearchResult { items: dtos, total }) } Err(e) => domain_err(e), } } + +pub struct BrowseParams { + pub content_type: Option, + pub genres: Vec, + pub search_term: Option, + pub series_names: Vec, + pub collections: Vec, + pub decade: Option, + pub role: Option, + pub min_duration_secs: Option, + pub max_duration_secs: Option, + pub limit: Option, +} + +pub async fn browse_library( + library_command_deps: &Arc, + params: BrowseParams, +) -> String { + let BrowseParams { + content_type, + genres, + search_term, + series_names, + collections, + decade, + role, + min_duration_secs, + max_duration_secs, + limit, + } = params; + let ct_str = content_type.clone(); + let query = SearchItemsQuery { + provider_id: None, + content_type, + genres: genres.clone(), + search_term: search_term.clone(), + series_names: series_names.clone(), + collection_id: collections.first().cloned(), + season_number: None, + decade, + offset: 0, + limit: limit.unwrap_or(DEFAULT_SEARCH_LIMIT), + }; + match application::library::search::execute(library_command_deps, query).await { + Ok((items, total)) => { + let filtered: Vec<_> = items + .iter() + .filter(|i| { + if let Some(ref r) = role { + let item_role = role_to_str(i.role()); + if item_role != r.as_str() { + return false; + } + } + if let Some(min) = min_duration_secs { + if i.duration_secs() < min { + return false; + } + } + if let Some(max) = max_duration_secs { + if i.duration_secs() > max { + return false; + } + } + true + }) + .collect(); + + let dtos: Vec = filtered + .iter() + .map(|i| { + let se = match (i.season_number(), i.episode_number()) { + (Some(s), Some(e)) => Some(format!("S{s:02}E{e:02}")), + _ => None, + }; + BrowseItemDto { + id: i.id().value().to_string(), + title: i.title().to_string(), + content_type: content_type_to_str(i.content_type()).to_string(), + duration_mins: i.duration_secs() / 60, + series_name: i.series_name().map(|s| s.to_string()), + season_episode: se, + year: i.year(), + genres: i.genres().to_vec(), + role: role_to_str(i.role()).to_string(), + } + }) + .collect(); + + let mut parts = Vec::new(); + parts.push(format!("{} items matched", filtered.len())); + if total > filtered.len() as u32 { + parts.push(format!("({total} total before role/duration filter)")); + } + if let Some(ref ct) = ct_str { + parts.push(format!("type={ct}")); + } + if !genres.is_empty() { + parts.push(format!("genres={}", genres.join(","))); + } + if let Some(ref t) = search_term { + parts.push(format!("search=\"{t}\"")); + } + + ok_json(&BrowseResult { + items: dtos, + total: filtered.len() as u32, + summary: parts.join(", "), + }) + } + Err(e) => domain_err(e), + } +} + +pub async fn library_stats( + library_query: &Arc, + library_command_deps: &Arc, +) -> String { + let all_query = SearchItemsQuery { + provider_id: None, + content_type: None, + genres: vec![], + search_term: None, + series_names: vec![], + collection_id: None, + season_number: None, + decade: None, + offset: 0, + limit: STATS_SEARCH_LIMIT, + }; + + let (items, total) = match application::library::search::execute(library_command_deps, all_query) + .await + { + Ok(r) => r, + Err(e) => return domain_err(e), + }; + + let mut by_content_type: HashMap = HashMap::new(); + let mut by_role: HashMap = HashMap::new(); + let mut genre_counts: HashMap = HashMap::new(); + let mut total_duration_secs: u64 = 0; + + for item in &items { + *by_content_type + .entry(content_type_to_str(item.content_type()).to_string()) + .or_default() += 1; + *by_role + .entry(role_to_str(item.role()).to_string()) + .or_default() += 1; + for genre in item.genres() { + *genre_counts.entry(genre.clone()).or_default() += 1; + } + total_duration_secs += item.duration_secs() as u64; + } + + let mut genres: Vec = genre_counts + .into_iter() + .map(|(genre, count)| GenreStat { genre, count }) + .collect(); + genres.sort_by_key(|g| std::cmp::Reverse(g.count)); + + let shows = match library_query.list_shows(None, None, &[]).await { + Ok(s) => s, + Err(e) => return domain_err(e), + }; + + let series: Vec = shows + .iter() + .map(|s| SeriesStat { + name: s.series_name().to_string(), + episode_count: s.episode_count(), + season_count: s.season_count(), + genres: s.genres().to_vec(), + }) + .collect(); + + let recently_synced = match library_query.latest_sync_status().await { + Ok(logs) => logs + .iter() + .map(|l| RecentSyncDto { + provider_id: l.provider_id().to_string(), + started_at: l.started_at().to_string(), + status: l.status().to_string(), + items_found: l.items_found(), + }) + .collect(), + Err(_) => vec![], + }; + + ok_json(&LibraryStats { + total_items: total, + by_content_type, + by_role, + genres, + series, + total_duration_hours: total_duration_secs as f64 / 3600.0, + recently_synced, + }) +} diff --git a/crates/mcp/src/tools/schedule.rs b/crates/mcp/src/tools/schedule.rs index abc9c32..a855d7d 100644 --- a/crates/mcp/src/tools/schedule.rs +++ b/crates/mcp/src/tools/schedule.rs @@ -1,11 +1,11 @@ +use std::collections::HashMap; use std::sync::Arc; -use application::schedule::{ - GenerateScheduleCommand, GetCurrentBroadcastQuery, ScheduleDeps, -}; +use application::channels::ChannelCommandDeps; +use application::schedule::{GenerateScheduleCommand, GetCurrentBroadcastQuery, ScheduleDeps}; use chrono::Utc; -use domain::ScheduledSlot; -use domain::value_objects::ChannelId; +use domain::value_objects::{ChannelId, MediaFilter}; +use domain::{InterstitialRule, MidRollRule, ScheduleConfig, ScheduledSlot}; use serde::Serialize; use uuid::Uuid; @@ -51,3 +51,537 @@ pub async fn get_current_broadcast(deps: &Arc, channel_id: Uuid) - Err(e) => domain_err(e), } } + +#[derive(Serialize)] +struct ScheduleAnalysis { + channel_id: String, + channel_name: String, + has_schedule: bool, + schedule_valid_from: Option, + schedule_valid_until: Option, + total_slots: usize, + total_hours: f64, + most_played: Vec, + genre_distribution: HashMap, + block_coverage: Vec, + upcoming_gaps: Vec, +} + +#[derive(Serialize)] +struct ItemPlayCount { + title: String, + content_type: String, + count: u32, +} + +#[derive(Serialize)] +struct BlockCoverage { + block_name: String, + slot_count: usize, + total_minutes: f64, +} + +#[derive(Serialize)] +struct GapInfo { + from: String, + to: String, + duration_mins: f64, +} + +pub async fn analyze_schedule( + channel_query: &Arc, + schedule_query: &Arc, + channel_id: Uuid, +) -> String { + let cid = ChannelId::from(channel_id); + + let channel = match channel_query.find_by_id(cid).await { + Ok(Some(c)) => c, + Ok(None) => { + return serde_json::json!({"error": "Channel not found"}).to_string(); + } + Err(e) => return domain_err(e), + }; + + let schedule = match schedule_query.find_latest(cid).await { + Ok(s) => s, + Err(e) => return domain_err(e), + }; + + let (has_schedule, valid_from, valid_until, slots) = match &schedule { + Some(s) => ( + true, + Some(s.valid_from().to_rfc3339()), + Some(s.valid_until().to_rfc3339()), + s.slots(), + ), + None => (false, None, None, [].as_slice()), + }; + + let mut title_counts: HashMap = HashMap::new(); + let mut genre_distribution: HashMap = HashMap::new(); + let mut block_slots: HashMap = HashMap::new(); + let mut total_secs: f64 = 0.0; + + for slot in slots { + let duration = + (slot.end_at() - slot.start_at()).num_seconds() as f64; + total_secs += duration; + + let title_key = slot.item().title().to_string(); + let entry = title_counts + .entry(title_key) + .or_insert_with(|| { + let ct = match slot.item().content_type() { + domain::ContentType::Movie => "movie", + domain::ContentType::Episode => "episode", + domain::ContentType::Short => "short", + }; + (ct.to_string(), 0) + }); + entry.1 += 1; + + for genre in slot.item().genres() { + *genre_distribution.entry(genre.clone()).or_default() += 1; + } + + let block_name = slot.source_block_id().to_string(); + let block_entry = block_slots.entry(block_name).or_insert((0, 0.0)); + block_entry.0 += 1; + block_entry.1 += duration / 60.0; + } + + let config = channel.schedule_config(); + let block_name_map: HashMap = config + .all_blocks() + .map(|b| (b.id().to_string(), b.name().to_string())) + .collect(); + + let mut most_played: Vec = title_counts + .into_iter() + .map(|(title, (ct, count))| ItemPlayCount { + title, + content_type: ct, + count, + }) + .collect(); + most_played.sort_by_key(|i| std::cmp::Reverse(i.count)); + most_played.truncate(20); + + let block_coverage: Vec = block_slots + .into_iter() + .map(|(block_id, (slot_count, total_minutes))| { + let name = block_name_map + .get(&block_id) + .cloned() + .unwrap_or(block_id); + BlockCoverage { + block_name: name, + slot_count, + total_minutes, + } + }) + .collect(); + + let mut upcoming_gaps = Vec::new(); + let now = Utc::now(); + let future_slots: Vec<&ScheduledSlot> = slots + .iter() + .filter(|s| s.end_at() > now) + .collect(); + + for window in future_slots.windows(2) { + let gap_secs = (window[1].start_at() - window[0].end_at()).num_seconds(); + if gap_secs > 60 { + upcoming_gaps.push(GapInfo { + from: window[0].end_at().to_rfc3339(), + to: window[1].start_at().to_rfc3339(), + duration_mins: gap_secs as f64 / 60.0, + }); + } + } + upcoming_gaps.truncate(10); + + ok_json(&ScheduleAnalysis { + channel_id: channel_id.to_string(), + channel_name: channel.name().to_string(), + has_schedule, + schedule_valid_from: valid_from, + schedule_valid_until: valid_until, + total_slots: slots.len(), + total_hours: total_secs / 3600.0, + most_played, + genre_distribution, + block_coverage, + upcoming_gaps, + }) +} + +#[derive(Serialize)] +struct PreviewResult { + slot_count: usize, + total_hours: f64, + slots: Vec, +} + +#[derive(Serialize)] +struct PreviewSlotDto { + start_at: String, + end_at: String, + title: String, + content_type: String, + duration_mins: f64, + block_id: String, +} + +pub async fn preview_schedule( + schedule_deps: &Arc, + channel_id: Uuid, + duration_hours: Option, +) -> String { + let cid = ChannelId::from(channel_id); + let hours = duration_hours.unwrap_or(24); + + match schedule_deps + .schedule_engine + .preview_schedule(cid, Utc::now(), hours) + .await + { + Ok(schedule) => { + let total_secs: f64 = schedule + .slots() + .iter() + .map(|s| (s.end_at() - s.start_at()).num_seconds() as f64) + .sum(); + + let slots: Vec = schedule + .slots() + .iter() + .map(|s| PreviewSlotDto { + start_at: s.start_at().to_rfc3339(), + end_at: s.end_at().to_rfc3339(), + title: s.item().title().to_string(), + content_type: match s.item().content_type() { + domain::ContentType::Movie => "movie", + domain::ContentType::Episode => "episode", + domain::ContentType::Short => "short", + } + .to_string(), + duration_mins: (s.end_at() - s.start_at()).num_seconds() as f64 / 60.0, + block_id: s.source_block_id().to_string(), + }) + .collect(); + + ok_json(&PreviewResult { + slot_count: slots.len(), + total_hours: total_secs / 3600.0, + slots, + }) + } + Err(e) => domain_err(e), + } +} + +pub async fn preview_config( + schedule_deps: &Arc, + channel_id: Uuid, + config_json: &str, + duration_hours: Option, +) -> String { + let config: ScheduleConfig = match serde_json::from_str(config_json) { + Ok(c) => c, + Err(e) => { + return serde_json::json!({"error": format!("invalid schedule config: {e}")}) + .to_string(); + } + }; + + let cid = ChannelId::from(channel_id); + let hours = duration_hours.unwrap_or(24); + + match schedule_deps + .schedule_engine + .preview_config(cid, &config, Utc::now(), hours) + .await + { + Ok(schedule) => { + let total_secs: f64 = schedule + .slots() + .iter() + .map(|s| (s.end_at() - s.start_at()).num_seconds() as f64) + .sum(); + + let slots: Vec = schedule + .slots() + .iter() + .map(|s| PreviewSlotDto { + start_at: s.start_at().to_rfc3339(), + end_at: s.end_at().to_rfc3339(), + title: s.item().title().to_string(), + content_type: match s.item().content_type() { + domain::ContentType::Movie => "movie", + domain::ContentType::Episode => "episode", + domain::ContentType::Short => "short", + } + .to_string(), + duration_mins: (s.end_at() - s.start_at()).num_seconds() as f64 / 60.0, + block_id: s.source_block_id().to_string(), + }) + .collect(); + + ok_json(&PreviewResult { + slot_count: slots.len(), + total_hours: total_secs / 3600.0, + slots, + }) + } + Err(e) => domain_err(e), + } +} + +pub async fn suggest_schedule( + library_query: &Arc, + genres: Vec, + content_type: Option, + time_block_name: Option, + start_time: Option, + duration_mins: Option, +) -> String { + let ct = match content_type + .as_deref() + .map(application::library::parse_content_type) + .transpose() + { + Ok(ct) => ct, + Err(e) => return domain_err(e), + }; + + let available_genres = if genres.is_empty() { + match library_query.list_genres(ct.as_ref(), None).await { + Ok(g) => g, + Err(e) => return domain_err(e), + } + } else { + genres.clone() + }; + + let filter = MediaFilter { + content_type: ct.clone(), + genres: genres.clone(), + ..Default::default() + }; + + let block_name = time_block_name.unwrap_or_else(|| { + if !genres.is_empty() { + format!("{} Block", genres.join("/")) + } else if let Some(ref c) = ct { + let label = match c { + domain::ContentType::Movie => "Movie", + domain::ContentType::Episode => "Episode", + domain::ContentType::Short => "Short", + }; + format!("{label} Block") + } else { + "Programming Block".to_string() + } + }); + + let start = start_time.unwrap_or_else(|| "20:00".to_string()); + let dur = duration_mins.unwrap_or(180); + + let strategy = if ct.as_ref() == Some(&domain::ContentType::Movie) { + "best_fit" + } else if ct.as_ref() == Some(&domain::ContentType::Episode) { + "sequential" + } else { + "random" + }; + + let suggestion = serde_json::json!({ + "suggested_block": { + "name": block_name, + "start_time": start, + "duration_mins": dur, + "content": { + "type": "algorithmic", + "filter": filter, + "strategy": strategy, + }, + "loop_on_finish": true, + }, + "available_genres": available_genres, + "notes": format!( + "Suggested a {} block with {} strategy. Adjust start_time, duration_mins, and filter as needed. Apply via update_channel with schedule_config_json.", + block_name, strategy + ), + }); + + ok_json(&suggestion) +} + +pub async fn set_interstitial_rule( + channel_cmd_deps: &Arc, + channel_id: Uuid, + owner_id: Uuid, + block_id: Uuid, + rule_json: Option, +) -> String { + let cid = ChannelId::from(channel_id); + + let channel = match channel_cmd_deps + .channel_query + .find_by_id(cid) + .await + { + Ok(Some(c)) => c, + Ok(None) => { + return serde_json::json!({"error": "Channel not found"}).to_string(); + } + Err(e) => return domain_err(e), + }; + + if channel.owner_id() != domain::value_objects::UserId::from(owner_id) { + return serde_json::json!({"error": "Forbidden"}).to_string(); + } + + let rule: Option = match rule_json { + Some(json) => match serde_json::from_str(&json) { + Ok(r) => Some(r), + Err(e) => { + return serde_json::json!({"error": format!("invalid interstitial rule: {e}")}) + .to_string(); + } + }, + None => None, + }; + + let mut config = channel.schedule_config().clone(); + let bid = domain::value_objects::BlockId::from(block_id); + + match config.find_block_mut(bid) { + Some(block) => { + block.set_interstitial_rule(rule); + } + None => { + return serde_json::json!({"error": "Block not found in schedule config"}).to_string(); + } + } + + let cmd = application::channels::UpdateChannelCommand { + channel_id: cid, + owner_id: owner_id.into(), + name: None, + description: None, + timezone: None, + schedule_config: Some(config), + rotation_policy: None, + auto_schedule: None, + gap_filler: None, + }; + + match application::channels::update::execute(channel_cmd_deps, cmd).await { + Ok(ch) => ok_json(&ch), + Err(e) => domain_err(e), + } +} + +pub async fn set_mid_roll_rule( + channel_cmd_deps: &Arc, + channel_id: Uuid, + owner_id: Uuid, + block_id: Uuid, + rule_json: Option, +) -> String { + let cid = ChannelId::from(channel_id); + + let channel = match channel_cmd_deps + .channel_query + .find_by_id(cid) + .await + { + Ok(Some(c)) => c, + Ok(None) => { + return serde_json::json!({"error": "Channel not found"}).to_string(); + } + Err(e) => return domain_err(e), + }; + + if channel.owner_id() != domain::value_objects::UserId::from(owner_id) { + return serde_json::json!({"error": "Forbidden"}).to_string(); + } + + let rule: Option = match rule_json { + Some(json) => match serde_json::from_str(&json) { + Ok(r) => Some(r), + Err(e) => { + return serde_json::json!({"error": format!("invalid mid-roll rule: {e}")}) + .to_string(); + } + }, + None => None, + }; + + let mut config = channel.schedule_config().clone(); + let bid = domain::value_objects::BlockId::from(block_id); + + match config.find_block_mut(bid) { + Some(block) => { + block.set_mid_roll_rule(rule); + } + None => { + return serde_json::json!({"error": "Block not found in schedule config"}).to_string(); + } + } + + let cmd = application::channels::UpdateChannelCommand { + channel_id: cid, + owner_id: owner_id.into(), + name: None, + description: None, + timezone: None, + schedule_config: Some(config), + rotation_policy: None, + auto_schedule: None, + gap_filler: None, + }; + + match application::channels::update::execute(channel_cmd_deps, cmd).await { + Ok(ch) => ok_json(&ch), + Err(e) => domain_err(e), + } +} + +pub async fn set_gap_filler( + channel_cmd_deps: &Arc, + channel_id: Uuid, + owner_id: Uuid, + filter_json: Option, +) -> String { + let filter: Option = match filter_json { + Some(json) => match serde_json::from_str(&json) { + Ok(f) => Some(f), + Err(e) => { + return serde_json::json!({"error": format!("invalid media filter: {e}")}) + .to_string(); + } + }, + None => None, + }; + + let cmd = application::channels::UpdateChannelCommand { + channel_id: channel_id.into(), + owner_id: owner_id.into(), + name: None, + description: None, + timezone: None, + schedule_config: None, + rotation_policy: None, + auto_schedule: None, + gap_filler: Some(filter), + }; + + match application::channels::update::execute(channel_cmd_deps, cmd).await { + Ok(ch) => ok_json(&ch), + Err(e) => domain_err(e), + } +}