use std::sync::Arc; use chrono::{DateTime, Datelike, Duration, TimeZone, Utc}; use chrono_tz::Tz; use crate::errors::{DomainError, DomainResult}; use crate::models::{ BlockContent, CurrentBroadcast, GeneratedSchedule, PlaybackRecord, ProgrammingBlock, ScheduledSlot, }; use crate::ports::{ChannelQuery, LibraryQuery, ScheduleCommand, ScheduleQuery}; use crate::models::MediaItem; use crate::value_objects::{BlockId, ChannelId, FillStrategy, InterstitialRule, LibrarySearchFilter, MediaFilter, MediaItemId, MediaRole, MidRollRule, RotationPolicy, Weekday}; mod fill; mod rotation; const SCHEDULE_DURATION_DAYS: i64 = 7; struct BlockTimeWindow { start: DateTime, end: DateTime, } struct AlgorithmicParams<'a> { filter: &'a crate::value_objects::MediaFilter, strategy: &'a FillStrategy, block_id: BlockId, loop_on_finish: bool, ignore_rotation_policy: bool, } struct RotationContext<'a> { history: &'a [PlaybackRecord], policy: &'a RotationPolicy, generation: u32, last_item_id: Option<&'a MediaItemId>, } pub struct ScheduleEngineService { library_query: Arc, channel_query: Arc, schedule_query: Arc, schedule_command: Arc, } impl ScheduleEngineService { pub fn new( library_query: Arc, channel_query: Arc, schedule_query: Arc, schedule_command: Arc, ) -> Self { Self { library_query, channel_query, schedule_query, schedule_command, } } pub async fn generate_schedule( &self, channel_id: ChannelId, from: DateTime, ) -> DomainResult { let channel = self.load_channel(channel_id).await?; let valid_until = from + Duration::days(SCHEDULE_DURATION_DAYS); let mut schedule = self .build_schedule(&channel, channel.schedule_config(), from, valid_until) .await?; if let Some(gap_filter) = channel.gap_filler() { let filler_items = self.query_gap_fillers(gap_filter).await?; if !filler_items.is_empty() { let generation = schedule.generation(); let mut slots = schedule.into_slots(); Self::fill_gaps(&mut slots, &filler_items, from, valid_until); schedule = GeneratedSchedule::new( channel_id, from, valid_until, generation, slots, ); } } self.schedule_command.save(&schedule).await?; for slot in schedule.slots() { let record = PlaybackRecord::new(channel_id, slot.item().id().clone(), schedule.generation()); self.schedule_command.save_playback_record(&record).await?; } Ok(schedule) } pub async fn preview_schedule( &self, channel_id: ChannelId, from: DateTime, duration_hours: u32, ) -> DomainResult { let channel = self.load_channel(channel_id).await?; let valid_until = from + Duration::hours(duration_hours as i64); self.build_schedule(&channel, channel.schedule_config(), from, valid_until) .await } pub async fn preview_config( &self, channel_id: ChannelId, config: &crate::models::ScheduleConfig, from: DateTime, duration_hours: u32, ) -> DomainResult { let channel = self.load_channel(channel_id).await?; let valid_until = from + Duration::hours(duration_hours as i64); self.build_schedule(&channel, config, from, valid_until) .await } async fn load_channel( &self, channel_id: ChannelId, ) -> DomainResult { self.channel_query .find_by_id(channel_id) .await? .ok_or(DomainError::ChannelNotFound(channel_id)) } async fn build_schedule( &self, channel: &crate::models::Channel, config: &crate::models::ScheduleConfig, valid_from: DateTime, valid_until: DateTime, ) -> DomainResult { let channel_id = 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 start_date = valid_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, ) -> Option { schedule .slots() .iter() .find(|s| s.start_at() <= now && now < s.end_at()) .map(|slot| { CurrentBroadcast::new( slot.clone(), (now - slot.start_at()).num_seconds() as u32, ) }) } pub async fn get_latest_schedule( &self, channel_id: ChannelId, ) -> DomainResult> { self.schedule_query.find_latest(channel_id).await } pub async fn get_active_schedule( &self, channel_id: ChannelId, at: DateTime, ) -> DomainResult> { self.schedule_query.find_active(channel_id, at).await } pub async fn list_schedule_history( &self, channel_id: ChannelId, ) -> DomainResult> { self.schedule_query.list_schedule_history(channel_id).await } pub async fn get_schedule_by_id( &self, channel_id: ChannelId, schedule_id: crate::value_objects::ScheduleId, ) -> DomainResult> { self.schedule_query .get_schedule_by_id(channel_id, schedule_id) .await } pub async fn delete_schedules_after( &self, channel_id: ChannelId, target_generation: u32, ) -> DomainResult<()> { self.schedule_command .delete_schedules_after(channel_id, target_generation) .await } pub fn get_epg( schedule: &GeneratedSchedule, from: DateTime, until: DateTime, ) -> Vec<&ScheduledSlot> { schedule .slots() .iter() .filter(|s| s.start_at() < until && s.end_at() > from) .collect() } async fn resolve_block( &self, block: &ProgrammingBlock, window: BlockTimeWindow, rotation: RotationContext<'_>, ) -> DomainResult> { let program_slots = match block.content() { BlockContent::Manual { items } => { self.resolve_manual(items, window.start, window.end, block.id()) .await? } BlockContent::Algorithmic { filter, strategy, } => { self.resolve_algorithmic( AlgorithmicParams { filter, strategy, block_id: block.id(), loop_on_finish: block.loop_on_finish(), ignore_rotation_policy: block.ignore_rotation_policy(), }, BlockTimeWindow { start: window.start, end: window.end, }, rotation, ) .await? } }; let with_interstitials = if let Some(rule) = block.interstitial_rule() { self.insert_interstitials(program_slots, rule, block.id(), window.end) .await? } else { program_slots }; if let Some(rule) = block.mid_roll_rule() { self.apply_mid_rolls(with_interstitials, rule, block.id(), window.end) .await } else { Ok(with_interstitials) } } async fn insert_interstitials( &self, program_slots: Vec, rule: &InterstitialRule, block_id: BlockId, block_end: DateTime, ) -> DomainResult> { if program_slots.len() < 2 { return Ok(program_slots); } let mut filter = media_filter_to_library_search(rule.pool_filter()); filter = filter.with_role(MediaRole::Interstitial); let (interstitials, _) = self.library_query.search(&filter).await?; if interstitials.is_empty() { return Ok(program_slots); } let mut result: Vec = Vec::new(); let mut cursor = program_slots[0].start_at(); let mut interstitial_idx = 0; for (i, slot) in program_slots.iter().enumerate() { if cursor >= block_end { break; } let program_duration = slot.item().duration_secs(); let program_end = (cursor + Duration::seconds(program_duration as i64)).min(block_end); result.push(ScheduledSlot::new( cursor, program_end, slot.item().clone(), block_id, )); cursor = program_end; let should_insert = i + 1 < program_slots.len() && program_duration >= rule.min_gap_secs() && cursor < block_end; if should_insert { let interstitial = &interstitials[interstitial_idx % interstitials.len()]; let interstitial_end = (cursor + Duration::seconds(interstitial.duration_secs() as i64)) .min(block_end); if interstitial_end > cursor { result.push(ScheduledSlot::new( cursor, interstitial_end, interstitial.clone(), block_id, )); cursor = interstitial_end; interstitial_idx += 1; } } } Ok(result) } async fn apply_mid_rolls( &self, slots: Vec, rule: &MidRollRule, block_id: BlockId, block_end: DateTime, ) -> DomainResult> { let interval_secs = rule.fallback_interval_mins() as u64 * 60; if interval_secs == 0 { return Ok(slots); } let mut filter = media_filter_to_library_search(rule.pool_filter()); filter = filter.with_role(MediaRole::Interstitial); let (break_items, _) = self.library_query.search(&filter).await?; if break_items.is_empty() { return Ok(slots); } let mut result: Vec = Vec::new(); let mut break_idx = 0usize; for slot in &slots { let item_duration = slot.item().duration_secs() as u64; if item_duration < interval_secs { result.push(slot.clone()); continue; } let break_points = Self::compute_break_points( slot.item(), rule.prefer_chapters(), interval_secs, ); if break_points.is_empty() { result.push(slot.clone()); continue; } let mut cursor = slot.start_at(); let mut prev_offset = 0u64; for bp in &break_points { if cursor >= block_end { break; } let segment_duration = bp - prev_offset; let segment_end = (cursor + Duration::seconds(segment_duration as i64)).min(block_end); result.push(ScheduledSlot::new( cursor, segment_end, slot.item().clone(), block_id, )); cursor = segment_end; prev_offset = *bp; if cursor < block_end { let break_item = &break_items[break_idx % break_items.len()]; let break_end_secs = rule.break_duration_secs().min(break_item.duration_secs()); let break_end = (cursor + Duration::seconds(break_end_secs as i64)).min(block_end); if break_end > cursor { result.push(ScheduledSlot::new( cursor, break_end, break_item.clone(), block_id, )); cursor = break_end; break_idx += 1; } } } let remaining = item_duration - prev_offset; if remaining > 0 && cursor < block_end { let tail_end = (cursor + Duration::seconds(remaining as i64)).min(block_end); result.push(ScheduledSlot::new( cursor, tail_end, slot.item().clone(), block_id, )); } } Ok(result) } fn compute_break_points( item: &MediaItem, prefer_chapters: bool, interval_secs: u64, ) -> Vec { if prefer_chapters && !item.chapters().is_empty() { item.chapters() .iter() .filter_map(|ch| { let end = ch.end_secs() as u64; if end > 0 && end < item.duration_secs() as u64 { Some(end) } else { None } }) .collect() } else { let duration = item.duration_secs() as u64; let mut points = Vec::new(); let mut offset = interval_secs; while offset < duration { points.push(offset); offset += interval_secs; } points } } async fn resolve_manual( &self, item_ids: &[MediaItemId], start: DateTime, end: DateTime, block_id: BlockId, ) -> DomainResult> { let mut slots = Vec::new(); let mut cursor = start; for item_id in item_ids { if cursor >= end { break; } 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)); cursor = item_end; } } Ok(slots) } async fn resolve_algorithmic( &self, params: AlgorithmicParams<'_>, window: BlockTimeWindow, rotation: RotationContext<'_>, ) -> DomainResult> { let library_filter = media_filter_to_library_search(params.filter).with_role(MediaRole::Program); let (candidates, _total) = self.library_query.search(&library_filter).await?; if candidates.is_empty() { return Ok(vec![]); } let pool = if params.ignore_rotation_policy { candidates.clone() } else { rotation::apply_rotation_policy(&candidates, rotation.history, rotation.policy, rotation.generation) }; let target_secs = (window.end - window.start).num_seconds() as u32; let selected = fill::fill_block( &candidates, &pool, target_secs, params.strategy, rotation.last_item_id, params.loop_on_finish, rotation.history, ); let mut slots = Vec::new(); let mut cursor = window.start; for item in selected { if cursor >= window.end { break; } let item_end = (cursor + Duration::seconds(item.duration_secs() as i64)).min(window.end); slots.push(ScheduledSlot::new(cursor, item_end, item.clone(), params.block_id)); cursor = item_end; } Ok(slots) } async fn query_gap_fillers(&self, gap_filter: &MediaFilter) -> DomainResult> { let filter = media_filter_to_library_search(gap_filter).with_role(MediaRole::Interstitial); let (items, _) = self.library_query.search(&filter).await?; Ok(items) } fn fill_gaps( slots: &mut Vec, fillers: &[MediaItem], valid_from: DateTime, valid_until: DateTime, ) { if fillers.is_empty() { return; } let gap_block_id = BlockId::generate(); let mut gap_slots: Vec = Vec::new(); let mut filler_idx = 0usize; let mut boundaries: Vec<(DateTime, DateTime)> = Vec::new(); if slots.is_empty() { boundaries.push((valid_from, valid_until)); } else { if slots[0].start_at() > valid_from { boundaries.push((valid_from, slots[0].start_at())); } for pair in slots.windows(2) { if pair[1].start_at() > pair[0].end_at() { boundaries.push((pair[0].end_at(), pair[1].start_at())); } } if let Some(last) = slots.last() && last.end_at() < valid_until { boundaries.push((last.end_at(), valid_until)); } } for (gap_start, gap_end) in boundaries { let mut cursor = gap_start; while cursor < gap_end { let filler = &fillers[filler_idx % fillers.len()]; let filler_end = (cursor + Duration::seconds(filler.duration_secs() as i64)).min(gap_end); if filler_end <= cursor { break; } gap_slots.push(ScheduledSlot::new( cursor, filler_end, filler.clone(), gap_block_id, )); cursor = filler_end; filler_idx += 1; } } slots.append(&mut gap_slots); slots.sort_by_key(|s| s.start_at()); } } 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 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 } #[cfg(all(test, feature = "test-helpers"))] #[path = "tests/integration.rs"] mod integration_tests;