Update scheduling parameters and refactor channel endpoints: extend HorizonDays to 7 and RetentionDays to 90 in appsettings.json. Consolidate channel-related endpoint logic by removing obsolete files and enhancing the ShowEndpoints with audience and genre management capabilities. Improve error handling and streamline command handlers for channel operations.
This commit is contained in:
@@ -0,0 +1,508 @@
|
||||
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;
|
||||
|
||||
/// <summary>Что дала генерация: сколько записей добавлено и с какими предупреждениями.</summary>
|
||||
public sealed record GenerationReport(
|
||||
int Added,
|
||||
IReadOnlyList<PlanningWarning> Warnings,
|
||||
bool ChannelSkipped = false
|
||||
);
|
||||
|
||||
/// <summary>
|
||||
/// Оркестратор генерации по сетке: собирает конфигурацию канала, вызывает чистый планировщик,
|
||||
/// материализует ленту и двигает курсоры слотов.
|
||||
///
|
||||
/// Генерация одного канала сериализуется advisory-блокировкой: фоновый тик и ручное применение
|
||||
/// не должны читать одну точку продолжения и оба дописывать хвост.
|
||||
/// </summary>
|
||||
public sealed class GridScheduleGenerator(
|
||||
IAppDbContext dbContext,
|
||||
GroupExpander expander,
|
||||
BumperResolver bumperResolver,
|
||||
IRandomSource random,
|
||||
IOptions<SchedulerOptions> options,
|
||||
IOptions<StreamingOptions> 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,
|
||||
};
|
||||
|
||||
/// <summary>
|
||||
/// Достраивает горизонт канала, а при <paramref name="rebuildFuture"/> — пересобирает будущий
|
||||
/// хвост от текущего момента. Прошлое и идущая сейчас запись не трогаются никогда: зритель
|
||||
/// не должен обнаружить, что у него из-под носа вырезали программу.
|
||||
/// </summary>
|
||||
public async Task<GenerationReport> 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);
|
||||
|
||||
// Заставки резолвятся после сборки ленты: пара соседей известна только теперь.
|
||||
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);
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Чистит прошлое сверх окна хранения. Окно должно покрывать самое долгое остывание среди правил —
|
||||
/// история показов берётся из самой ленты, отдельного журнала нет.
|
||||
/// </summary>
|
||||
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<PlanningCursorUpdate> 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<PlanningInput> 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 elementsByGroup = await expander.ExpandAsync(groupIds, channel.Id, cancellationToken);
|
||||
|
||||
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<PlanningSlot>();
|
||||
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
|
||||
)
|
||||
)
|
||||
);
|
||||
}
|
||||
|
||||
var fallback = await LoadFallbackUnitsAsync(channel, template, cancellationToken);
|
||||
|
||||
return new PlanningInput(
|
||||
channel.Id,
|
||||
startUtc,
|
||||
horizonEnd,
|
||||
slots,
|
||||
fallback,
|
||||
_segmentSeconds
|
||||
);
|
||||
}
|
||||
|
||||
/// <summary>Ассет заставки по зарезервированной записи; false — подобрать не удалось.</summary>
|
||||
private static bool TryResolveBumper(
|
||||
PlannedItem item,
|
||||
IReadOnlyDictionary<BumperKey, Guid> 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;
|
||||
}
|
||||
|
||||
/// <summary>Разворачивает шаблон стыка для планировщика, включая резерв под заставки.</summary>
|
||||
private PlanningJunction? BuildJunction(
|
||||
Guid? junctionId,
|
||||
IReadOnlyDictionary<Guid, JunctionTemplate> junctions,
|
||||
IReadOnlyDictionary<Guid, IReadOnlyList<PlanningElement>> elementsByGroup,
|
||||
Channel channel
|
||||
)
|
||||
{
|
||||
if (junctionId is not { } id || !junctions.TryGetValue(id, out var template))
|
||||
return null;
|
||||
|
||||
var elements = new List<PlanningJunctionElement>();
|
||||
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);
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Что играло в точке, на которую ссылается слот-повтор. Читается уже записанная лента того же
|
||||
/// канала — стратегий и состояния повтору не нужно.
|
||||
/// </summary>
|
||||
private async Task<IReadOnlyList<PlanningUnit>> 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();
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Чем закрывать паузы. Сначала аварийная группа шаблона, затем филлер канала — последний
|
||||
/// зациклен, поэтому короткая единица закрывает любой остаток лучше длинной.
|
||||
/// </summary>
|
||||
private async Task<IReadOnlyList<PlanningUnit>> LoadFallbackUnitsAsync(
|
||||
Channel channel,
|
||||
ScheduleTemplate template,
|
||||
CancellationToken cancellationToken
|
||||
)
|
||||
{
|
||||
var units = new List<PlanningUnit>();
|
||||
|
||||
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,
|
||||
};
|
||||
}
|
||||
Reference in New Issue
Block a user