using System.Collections.Concurrent;
using Microsoft.EntityFrameworkCore;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Hosting;
using Microsoft.Extensions.Logging;
using Microsoft.Extensions.Options;
using TeleWave.Application.Common.Interfaces;
using TeleWave.Domain.Media;
namespace TeleWave.Infrastructure.Media;
///
/// Обработчик медиа: прогоняет ассеты через ffmpeg, до
/// файлов одновременно. Источник истины — статус в БД: единственный диспетчер последовательно и
/// атомарно захватывает следующий Pending (помечает Processing), поэтому два транскода никогда не
/// возьмут один ассет; сам транскод запускается в фоне с ограничением по числу слотов. Рестарт/краш
/// ничего не теряет — незавершённые подхватываются из БД (прерванные Processing на старте сбрасываются
/// в Pending). БД-контекст держится короткими отрезками (пометить статус), сам транскод идёт вне
/// scope, чтобы не держать соединение минутами.
///
public sealed class MediaProcessingBackgroundService(
IMediaProcessingQueue queue,
IServiceScopeFactory scopeFactory,
MediaPathResolver paths,
IMediaProcessor processor,
IOptions storageOptions,
IOptions mediaOptions,
ILogger logger
) : BackgroundService
{
// Периодически перепроверяем БД, даже если сигнал не пришёл — страховка на любой случай.
private static readonly TimeSpan IdlePoll = TimeSpan.FromSeconds(30);
private readonly StorageOptions _storage = storageOptions.Value;
private readonly int _maxParallel = Math.Max(1, mediaOptions.Value.MaxParallelTranscodes);
protected override async Task ExecuteAsync(CancellationToken stoppingToken)
{
paths.EnsureDirectories();
await ResetInterruptedAsync(stoppingToken);
// Слоты параллелизма: не запускаем больше _maxParallel транскодов одновременно.
using var slots = new SemaphoreSlim(_maxParallel, _maxParallel);
var inFlight = new ConcurrentDictionary();
while (!stoppingToken.IsCancellationRequested)
{
try
{
// Захватываем и раздаём по слотам всю накопившуюся работу из БД.
while (!stoppingToken.IsCancellationRequested)
{
await slots.WaitAsync(stoppingToken);
// Слот уже захвачен — любой сбой захвата ассета (транзиентная ошибка БД и т.п.)
// обязан вернуть слот, иначе после нескольких ошибок семафор исчерпается и
// диспетчер зависнет навсегда (сервис формально жив, но ничего не обрабатывает).
(Guid Id, string Extension)? claim;
try
{
claim = await ClaimNextAsync(stoppingToken);
}
catch
{
slots.Release();
throw;
}
if (claim is not { } job)
{
slots.Release();
break;
}
var task = Task.Run(
async () =>
{
try
{
await ProcessClaimedAsync(job.Id, job.Extension, stoppingToken);
}
finally
{
slots.Release();
}
},
CancellationToken.None
);
inFlight[task] = 0;
_ = task.ContinueWith(t => inFlight.TryRemove(t, out _), TaskScheduler.Default);
}
// Работы нет — ждём сигнала о новой либо периодического опроса.
using var wake = CancellationTokenSource.CreateLinkedTokenSource(stoppingToken);
wake.CancelAfter(IdlePoll);
try
{
await queue.WaitAsync(wake.Token);
}
catch (OperationCanceledException) when (!stoppingToken.IsCancellationRequested)
{
// Тайм-аут опроса — просто перепроверяем БД.
}
}
catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested)
{
break;
}
catch (Exception ex)
{
logger.LogError(ex, "Ошибка цикла обработки медиа");
await Task.Delay(TimeSpan.FromSeconds(5), stoppingToken);
}
}
// Даём уже запущенным транскодам корректно завершиться (или отмениться) на остановке.
try
{
await Task.WhenAll(inFlight.Keys.ToArray());
}
catch
{
// Ошибки/отмена отдельных задач уже залогированы внутри ProcessClaimedAsync.
}
}
/// Сброс прерванных рестартом задач (Processing → Pending) на старте.
private async Task ResetInterruptedAsync(CancellationToken cancellationToken)
{
await using var scope = scopeFactory.CreateAsyncScope();
var db = scope.ServiceProvider.GetRequiredService();
var interrupted = await db
.MediaAssets.Where(x => x.Status == MediaAssetStatus.Processing)
.ToListAsync(cancellationToken);
if (interrupted.Count == 0)
return;
foreach (var asset in interrupted)
asset.ResetToPending();
await db.SaveChangesAsync(cancellationToken);
}
///
/// Атомарно захватывает самый ранний Pending: помечает его Processing и возвращает (id, расширение),
/// либо null если работы нет. Вызывается только диспетчером последовательно, поэтому два транскода
/// не возьмут один ассет.
///
private async Task<(Guid Id, string Extension)?> ClaimNextAsync(
CancellationToken cancellationToken
)
{
await using var scope = scopeFactory.CreateAsyncScope();
var db = scope.ServiceProvider.GetRequiredService();
var asset = await db
.MediaAssets.Where(x => x.Status == MediaAssetStatus.Pending)
.OrderBy(x => x.CreatedAt)
.FirstOrDefaultAsync(cancellationToken);
if (asset is null)
return null;
asset.MarkProcessing();
await db.SaveChangesAsync(cancellationToken);
return (asset.Id, asset.OriginalExtension);
}
/// Обрабатывает уже захваченный (Processing) ассет: транскод → Ready/Failed.
private async Task ProcessClaimedAsync(
Guid assetId,
string extension,
CancellationToken cancellationToken
)
{
try
{
var result = await processor.ProcessAsync(assetId, extension, cancellationToken);
await CompleteAsync(assetId, result, cancellationToken);
if (!_storage.KeepOriginals)
DeleteOriginal(assetId, extension);
logger.LogInformation(
"Ассет {AssetId} обработан: {Segments} сегментов, {Seconds:0.#}с",
assetId,
result.SegmentCount,
result.Duration.TotalSeconds
);
}
catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested)
{
throw;
}
catch (Exception ex)
{
logger.LogError(ex, "Обработка ассета {AssetId} провалилась", assetId);
await FailAsync(assetId, ex.Message, CancellationToken.None);
}
}
private async Task CompleteAsync(
Guid assetId,
MediaProcessingResult result,
CancellationToken cancellationToken
)
{
await using var scope = scopeFactory.CreateAsyncScope();
var db = scope.ServiceProvider.GetRequiredService();
var asset = await db.MediaAssets.FirstOrDefaultAsync(
x => x.Id == assetId,
cancellationToken
);
if (asset is null)
return;
asset.MarkReady(
result.Duration,
result.SegmentSeconds,
result.SegmentCount,
result.Width,
result.Height,
result.VideoCodec,
result.AudioCodec,
result.RelativePath
);
await db.SaveChangesAsync(cancellationToken);
}
private async Task FailAsync(Guid assetId, string error, CancellationToken cancellationToken)
{
await using var scope = scopeFactory.CreateAsyncScope();
var db = scope.ServiceProvider.GetRequiredService();
var asset = await db.MediaAssets.FirstOrDefaultAsync(
x => x.Id == assetId,
cancellationToken
);
if (asset is null)
return;
asset.MarkFailed(error);
await db.SaveChangesAsync(cancellationToken);
}
private void DeleteOriginal(Guid assetId, string extension)
{
var original = paths.OriginalPath(assetId, extension);
if (File.Exists(original))
File.Delete(original);
}
}