using System.Text.Json; using Microsoft.EntityFrameworkCore; using Microsoft.Extensions.Options; using TeleWave.Application.Broadcast.Scheduling; using TeleWave.Application.Common.Interfaces; using TeleWave.Application.Programming.Templates; using TeleWave.Application.Streaming; using TeleWave.Domain.Broadcast; using TeleWave.Domain.Broadcast.Scheduling; using TeleWave.Domain.Media; using TeleWave.Domain.Programming; using TeleWave.Domain.Programming.Planning; namespace TeleWave.Application.Programming.Planning; /// Что дала генерация: сколько записей добавлено и с какими предупреждениями. public sealed record GenerationReport( int Added, IReadOnlyList Warnings, bool ChannelSkipped = false ); /// /// Оркестратор генерации по сетке: собирает конфигурацию канала, вызывает чистый планировщик, /// материализует ленту и двигает курсоры слотов. /// /// Генерация одного канала сериализуется advisory-блокировкой: фоновый тик и ручное применение /// не должны читать одну точку продолжения и оба дописывать хвост. /// public sealed class GridScheduleGenerator( IAppDbContext dbContext, GroupExpander expander, BumperResolver bumperResolver, PostCheckRunner postChecks, IRandomSource random, IOptions options, IOptions streamingOptions ) { private readonly SchedulerOptions _options = options.Value; private readonly int _segmentSeconds = Math.Max(1, streamingOptions.Value.SegmentSeconds); private static readonly JsonSerializerOptions TraceJsonOptions = new() { PropertyNamingPolicy = JsonNamingPolicy.CamelCase, }; /// /// Достраивает горизонт канала, а при — пересобирает будущий /// хвост от текущего момента. Прошлое и идущая сейчас запись не трогаются никогда: зритель /// не должен обнаружить, что у него из-под носа вырезали программу. /// public async Task GenerateAsync( Guid channelId, DateTimeOffset now, bool rebuildFuture, CancellationToken cancellationToken ) { var channel = await dbContext .Channels.Include(c => c.BumperTemplates) .ThenInclude(t => t.Variants) .FirstOrDefaultAsync(c => c.Id == channelId, cancellationToken); if (channel is null || !channel.IsEnabled || channel.TemplateId is null) return new GenerationReport(0, [], ChannelSkipped: true); await using var transaction = await dbContext.BeginTransactionAsync(cancellationToken); await dbContext.AcquireChannelLockAsync(channelId, cancellationToken); var template = await dbContext .ScheduleTemplates.Include(t => t.Layers) .ThenInclude(l => l.Slots) .FirstOrDefaultAsync(t => t.Id == channel.TemplateId, cancellationToken); if (template is null) return new GenerationReport(0, [], ChannelSkipped: true); await CleanupAsync(channelId, now, cancellationToken); if (rebuildFuture) await dbContext .ScheduleEntries.Where(e => e.ChannelId == channelId && e.StartsAtUtc >= now) .ExecuteDeleteAsync(cancellationToken); // Точка продолжения — конец последней сохранённой записи. При пересборке будущее уже удалено, // поэтому она укажет на границу неизменяемого прошлого. var lastEnd = await dbContext .ScheduleEntries.Where(e => e.ChannelId == channelId) .MaxAsync(e => (DateTimeOffset?)e.EndsAtUtc, cancellationToken); var startUtc = lastEnd is { } end && end > now ? end : now; var horizonEnd = now.AddDays(Math.Max(1, _options.HorizonDays)); if (startUtc >= horizonEnd) { await dbContext.SaveChangesAsync(cancellationToken); await transaction.CommitAsync(cancellationToken); return new GenerationReport(0, []); } var input = await BuildInputAsync( channel, template, startUtc, horizonEnd, cancellationToken ); var result = Domain.Programming.Planning.SchedulePlanner.Plan(input, random); // Пост-проверки считаются по готовой ленте и только предупреждают — переигрывать что-либо // по их итогу мы намеренно не будем (см. 3.8). var postWarnings = await postChecks.RunAsync( result.Items, PlanningRules.FromJson(template.RulesJson), channel.UtcOffsetMinutes, channel.DayStartTime, cancellationToken ); // Заставки резолвятся после сборки ленты: пара соседей известна только теперь. var bumperAssets = await bumperResolver.ResolveAsync( channel, result.Items, cancellationToken ); var added = 0; foreach (var item in result.Items) { var assetId = item.MediaAssetId; if (item.Kind == PlannedItemKind.Bumper && !TryResolveBumper(item, bumperAssets, out assetId)) continue; // Без ассета запись стала бы дырой в ленте. dbContext.ScheduleEntries.Add( ScheduleEntry.FromSlot( channel.Id, assetId, ToEntryKind(item.Kind), item.StartsAtUtc, item.EndsAtUtc, item.ShowId, item.UnitIndex, item.SlotId, item.Trace is null ? null : JsonSerializer.Serialize(item.Trace, TraceJsonOptions) ) ); added++; } await SaveCursorsAsync(result.Cursors, cancellationToken); template.MarkApplied(); await dbContext.SaveChangesAsync(cancellationToken); await transaction.CommitAsync(cancellationToken); return new GenerationReport(added, [.. result.Warnings, .. postWarnings]); } /// /// Сухой прогон: считает, каким получился бы эфир по текущим правилам, но ничего не пишет — /// ни ленты, ни курсоров слотов, ни отметки о применении. Заставки остаются резервом известной /// длины: рендер долгий и может упасть, поэтому он делается только при реальном применении. /// public async Task PreviewAsync( Guid channelId, DateTimeOffset from, int days, CancellationToken cancellationToken ) { var channel = await dbContext .Channels.AsNoTracking() .Include(c => c.BumperTemplates) .ThenInclude(t => t.Variants) .FirstOrDefaultAsync(c => c.Id == channelId, cancellationToken); if (channel is null || channel.TemplateId is null) return null; var template = await dbContext .ScheduleTemplates.AsNoTracking() .Include(t => t.Layers) .ThenInclude(l => l.Slots) .FirstOrDefaultAsync(t => t.Id == channel.TemplateId, cancellationToken); if (template is null) return null; var horizonEnd = from.AddDays(Math.Clamp(days, 1, Math.Max(1, _options.HorizonDays))); var input = await BuildInputAsync(channel, template, from, horizonEnd, cancellationToken); var result = Domain.Programming.Planning.SchedulePlanner.Plan(input, random); var postWarnings = await postChecks.RunAsync( result.Items, PlanningRules.FromJson(template.RulesJson), channel.UtcOffsetMinutes, channel.DayStartTime, cancellationToken ); return result with { Warnings = [.. result.Warnings, .. postWarnings] }; } /// /// Чистит прошлое сверх окна хранения. Окно должно покрывать самое долгое остывание среди правил — /// история показов берётся из самой ленты, отдельного журнала нет. /// private Task CleanupAsync(Guid channelId, DateTimeOffset now, CancellationToken cancellationToken) { var cutoff = now.AddDays(-Math.Max(1, _options.RetentionDays)); return dbContext .ScheduleEntries.Where(e => e.ChannelId == channelId && e.EndsAtUtc < cutoff) .ExecuteDeleteAsync(cancellationToken); } private async Task SaveCursorsAsync( IReadOnlyList cursors, CancellationToken cancellationToken ) { if (cursors.Count == 0) return; var slotIds = cursors.Select(c => c.SlotId).Distinct().ToList(); var states = await dbContext .SlotStates.Where(s => slotIds.Contains(s.SlotId)) .ToDictionaryAsync(s => s.SlotId, cancellationToken); foreach (var cursor in cursors) { if (!states.TryGetValue(cursor.SlotId, out var state)) { state = SlotState.Create(cursor.SlotId); dbContext.SlotStates.Add(state); states[cursor.SlotId] = state; } if (cursor.ElementKind is { } kind && cursor.ElementId is { } elementId) state.MoveTo(kind, elementId, cursor.NextUnitIndex); else state.Reset(); } } private async Task BuildInputAsync( Channel channel, ScheduleTemplate template, DateTimeOffset startUtc, DateTimeOffset horizonEnd, CancellationToken cancellationToken ) { var scheduled = EffectiveGridBuilder.Build( template, channel.UtcOffsetMinutes, channel.DayStartTime, startUtc, horizonEnd ); // Стыки грузим целиком: их немного, а группы врезок надо развернуть тем же проходом, // что и группы контента. var junctions = await dbContext .JunctionTemplates.AsNoTracking() .Include(j => j.Elements) .Where(j => j.ChannelId == channel.Id) .ToDictionaryAsync(j => j.Id, cancellationToken); var groupIds = scheduled .Select(s => s.Slot.GroupId) .Where(id => id is not null) .Select(id => id!.Value) .Concat( junctions .Values.SelectMany(j => j.Elements) .Select(e => e.GroupId) .Where(id => id is not null) .Select(id => id!.Value) ) .Distinct() .ToList(); // Правила канала: детское время резолвится на каждый слот по его времени, а под потолок // повторов нужна история показов — её глубина и задаётся окном правила. var rules = PlanningRules.FromJson(template.RulesJson); var repeatLimit = rules?.MaxRepeatsInWindow is { WindowDays: > 0, Max: > 0 } limit ? new RepeatLimit(limit.WindowDays, limit.Max) : null; var elementsByGroup = await expander.ExpandAsync( groupIds, channel.Id, cancellationToken, repeatLimit is null ? null : startUtc.AddDays(-repeatLimit.WindowDays) ); var slotIds = scheduled.Select(s => s.Slot.Id).Distinct().ToList(); var states = await dbContext .SlotStates.AsNoTracking() .Where(s => slotIds.Contains(s.SlotId)) .ToDictionaryAsync(s => s.SlotId, cancellationToken); var slots = new List(); foreach (var item in scheduled) { var slot = item.Slot; var elements = slot.GroupId is { } groupId && elementsByGroup.TryGetValue(groupId, out var groupElements) ? groupElements : []; var repeatUnits = slot.SlotKind == SlotKind.Repeat ? await LoadRepeatUnitsAsync(channel, slot, item, cancellationToken) : null; var strategy = ToPlanningStrategy(SlotStrategy.FromJson(slot.StrategyJson)); var cursor = states.TryGetValue(slot.Id, out var state) ? new PlanningCursor(state.CurrentElementKind, state.CurrentElementId, state.NextUnitIndex) : null; slots.Add( new PlanningSlot( slot.Id, item.StartUtc, slot.TargetDurationMinutes, slot.SlotKind, slot.IsAnchor, slot.MaxDriftMinutes, slot.SnapToMinutes, slot.BlockMode, slot.BlockValue, slot.OverflowPolicy, strategy, elements, cursor, repeatUnits, BuildJunction(slot.JunctionBetweenId, junctions, elementsByGroup, channel), BuildJunction( slot.JunctionAfterId ?? template.DefaultJunctionId, junctions, elementsByGroup, channel ), rules?.AudienceAt( TimeOnly.FromDateTime( item.StartUtc.ToOffset(TimeSpan.FromMinutes(channel.UtcOffsetMinutes)).DateTime ) ), repeatLimit ) ); } var fallback = await LoadFallbackUnitsAsync(channel, template, cancellationToken); return new PlanningInput( channel.Id, startUtc, horizonEnd, slots, fallback, _segmentSeconds ); } /// Ассет заставки по зарезервированной записи; false — подобрать не удалось. private static bool TryResolveBumper( PlannedItem item, IReadOnlyDictionary bumperAssets, out Guid assetId ) { assetId = Guid.Empty; if (item.BumperTemplateId is not { } templateId) return false; // Подблок выбирает резолвер, поэтому ищем по блоку и паре шоу. foreach (var pair in bumperAssets) { if ( pair.Key.TemplateId == templateId && pair.Key.FromShowId == (item.FromShowId ?? Guid.Empty) && pair.Key.ToShowId == (item.ToShowId ?? Guid.Empty) ) { assetId = pair.Value; return true; } } return false; } /// Разворачивает шаблон стыка для планировщика, включая резерв под заставки. private PlanningJunction? BuildJunction( Guid? junctionId, IReadOnlyDictionary junctions, IReadOnlyDictionary> elementsByGroup, Channel channel ) { if (junctionId is not { } id || !junctions.TryGetValue(id, out var template)) return null; var elements = new List(); foreach (var element in template.Elements.OrderBy(e => e.Position)) { var conditions = JunctionConditions.FromJson(element.ConditionsJson) ?? new JunctionConditions(); if (element.Kind == JunctionElementKind.Bumper) { // Длительность задаётся блоком (по звуку) и выровнена на сегмент: планировщик // резервирует именно её, ассет подставит резолвер после сборки ленты. if ( element.BumperTemplateId is not { } bumperTemplateId || channel.FindBumperTemplate(bumperTemplateId) is not { } bumperTemplate ) continue; var seconds = BumperDuration.Aligned( BumperDuration.TemplateSeconds(bumperTemplate), _segmentSeconds ); elements.Add( new PlanningJunctionElement( element.Kind, [], element.AmountMode, element.AmountValue, element.IsRequired, conditions.OnlyOnElementChange, conditions.MinMinutesBetween, bumperTemplateId, TimeSpan.FromSeconds(seconds) ) ); continue; } if ( element.GroupId is not { } groupId || !elementsByGroup.TryGetValue(groupId, out var groupElements) ) continue; elements.Add( new PlanningJunctionElement( element.Kind, groupElements.SelectMany(e => e.Units).ToList(), element.AmountMode, element.AmountValue, element.IsRequired, conditions.OnlyOnElementChange, conditions.MinMinutesBetween ) ); } return elements.Count == 0 ? null : new PlanningJunction(id, elements); } /// /// Что играло в точке, на которую ссылается слот-повтор. Читается уже записанная лента того же /// канала — стратегий и состояния повтору не нужно. /// private async Task> LoadRepeatUnitsAsync( Channel channel, Slot slot, ScheduledSlot scheduled, CancellationToken cancellationToken ) { var source = RepeatSource.FromJson(slot.RepeatSourceJson); if (source is null) return []; var offset = TimeSpan.FromMinutes(channel.UtcOffsetMinutes); var sourceDate = scheduled.BroadcastDate.AddDays(-Math.Max(1, source.DaysAgo)); var from = EffectiveGridBuilder.ToUtc(sourceDate, source.Time, offset, channel.DayStartTime); var to = from.AddMinutes(Math.Max(1, source.DurationMinutes)); var entries = await dbContext .ScheduleEntries.AsNoTracking() .Where(e => e.ChannelId == channel.Id && e.Kind == ScheduleEntryKind.Program && e.StartsAtUtc >= from && e.StartsAtUtc < to ) .OrderBy(e => e.StartsAtUtc) .Select(e => new { e.MediaAssetId, e.ShowId, e.EpisodeIndex, Duration = e.EndsAtUtc - e.StartsAtUtc, }) .ToListAsync(cancellationToken); return entries .Select(e => new PlanningUnit( e.MediaAssetId, e.Duration, e.ShowId ?? Guid.Empty, e.EpisodeIndex ?? 0 )) .ToList(); } /// /// Чем закрывать паузы. Сначала аварийная группа шаблона, затем филлер канала — последний /// зациклен, поэтому короткая единица закрывает любой остаток лучше длинной. /// private async Task> LoadFallbackUnitsAsync( Channel channel, ScheduleTemplate template, CancellationToken cancellationToken ) { var units = new List(); if (template.FallbackGroupId is { } groupId) { var expanded = await expander.ExpandAsync([groupId], channel.Id, cancellationToken); if (expanded.TryGetValue(groupId, out var elements)) units.AddRange(elements.SelectMany(e => e.Units)); } if (channel.FillerAssetId is { } fillerId) { var filler = await dbContext .MediaAssets.AsNoTracking() .Where(a => a.Id == fillerId && a.Status == MediaAssetStatus.Ready && a.Duration != null ) .Select(a => new { a.Id, a.Duration }) .FirstOrDefaultAsync(cancellationToken); if (filler is not null) units.Add(new PlanningUnit(filler.Id, filler.Duration!.Value, Guid.Empty, 0)); } // Короткие вперёд: чем короче единица, тем меньше остаётся незакрытым остаток паузы. return units.OrderBy(u => u.Duration).ToList(); } private static PlanningStrategy ToPlanningStrategy(SlotStrategy? strategy) { if (strategy is null) return new PlanningStrategy(SlotStrategyKind.Sequential); var kind = strategy.Type switch { SlotStrategyType.RandomWithCooldown => SlotStrategyKind.RandomWithCooldown, SlotStrategyType.Fixed => SlotStrategyKind.Fixed, _ => SlotStrategyKind.Sequential, }; return new PlanningStrategy( kind, strategy.RestartOnEnd, strategy.CooldownDays, strategy.Fallback == CooldownFallback.IgnoreCooldown, strategy.FixedElementId ); } private static ScheduleEntryKind ToEntryKind(PlannedItemKind kind) => kind switch { PlannedItemKind.SignOff => ScheduleEntryKind.SignOff, PlannedItemKind.Fallback => ScheduleEntryKind.Fallback, PlannedItemKind.Ad or PlannedItemKind.Promo => ScheduleEntryKind.Ad, PlannedItemKind.Bumper => ScheduleEntryKind.Bumper, _ => ScheduleEntryKind.Program, }; }