Refactor media asset handling: update ListMediaAssetsQuery to accept an array of statuses, modify MediaEndpoints to handle multiple statuses, and enhance MediaProcessingQueue for improved signal handling. Update frontend components to support new status filtering and improve user experience with asset selection.

This commit is contained in:
Leonid Pershin
2026-07-24 19:22:22 +03:00
parent 8281832e1d
commit 6e7db4a6a9
12 changed files with 303 additions and 66 deletions
@@ -95,7 +95,7 @@ public static class MediaEndpoints
private static async Task<IResult> List(
int page,
int pageSize,
MediaAssetStatus? status,
MediaAssetStatus[]? status,
string? search,
ISender sender,
CancellationToken cancellationToken
@@ -105,7 +105,7 @@ public static class MediaEndpoints
new ListMediaAssetsQuery(
page <= 0 ? 1 : page,
pageSize <= 0 ? 20 : pageSize,
status,
status ?? [],
search
),
cancellationToken
@@ -1,12 +1,15 @@
namespace TeleWave.Application.Common.Interfaces;
/// <summary>
/// Очередь фоновой обработки ассетов. Продюсеры (endpoint загрузки, inbox-сканер) ставят id ассета,
/// единственный фоновый потребитель обрабатывает по одному за раз.
/// Сигнал «появилась работа» для фонового обработчика. Источник истины — статус ассета в БД
/// (обработчик всегда берёт следующий Pending из базы), а очередь лишь будит его без задержки;
/// поэтому потеря сигнала при рестарте не теряет задачи — они подхватываются из БД.
/// </summary>
public interface IMediaProcessingQueue
{
/// <summary>Разбудить обработчик: появился ассет в статусе Pending.</summary>
void Enqueue(Guid assetId);
IAsyncEnumerable<Guid> DequeueAllAsync(CancellationToken cancellationToken);
/// <summary>Ждать сигнала о новой работе (с дренажом накопленных).</summary>
ValueTask WaitAsync(CancellationToken cancellationToken);
}
@@ -7,6 +7,6 @@ namespace TeleWave.Application.Media.ListMedia;
public sealed record ListMediaAssetsQuery(
int Page,
int PageSize,
MediaAssetStatus? Status,
IReadOnlyList<MediaAssetStatus> Statuses,
string? Search
) : IQuery<PagedList<MediaAssetDto>>;
@@ -15,8 +15,8 @@ public sealed class ListMediaAssetsQueryHandler(IAppDbContext dbContext)
{
var q = dbContext.MediaAssets.AsNoTracking();
if (query.Status is { } status)
q = q.Where(x => x.Status == status);
if (query.Statuses.Count > 0)
q = q.Where(x => query.Statuses.Contains(x.Status));
if (!string.IsNullOrWhiteSpace(query.Search))
{
@@ -9,10 +9,11 @@ using TeleWave.Domain.Media;
namespace TeleWave.Infrastructure.Media;
/// <summary>
/// Единственный потребитель очереди обработки: по одному ассету за раз прогоняет через ffmpeg.
/// На старте восстанавливает прерванные задачи (Pending/Processing) — переживает рестарт/краш.
/// БД-контекст держится короткими отрезками (пометить статус), сам транскод идёт вне scope, чтобы
/// не держать соединение открытым минутами.
/// Единственный обработчик медиа: по одному ассету за раз прогоняет через ffmpeg. Источник истины —
/// статус в БД: сервис в цикле берёт следующий Pending из базы и обрабатывает, а очередь лишь будит
/// его без задержки. Поэтому рестарт/краш ничего не теряет — незавершённые задачи подхватываются из
/// БД (прерванные Processing на старте сбрасываются в Pending). БД-контекст держится короткими
/// отрезками (пометить статус), сам транскод идёт вне scope, чтобы не держать соединение минутами.
/// </summary>
public sealed class MediaProcessingBackgroundService(
IMediaProcessingQueue queue,
@@ -23,18 +24,34 @@ public sealed class MediaProcessingBackgroundService(
ILogger<MediaProcessingBackgroundService> logger
) : BackgroundService
{
// Периодически перепроверяем БД, даже если сигнал не пришёл — страховка на любой случай.
private static readonly TimeSpan IdlePoll = TimeSpan.FromSeconds(30);
private readonly StorageOptions _storage = storageOptions.Value;
protected override async Task ExecuteAsync(CancellationToken stoppingToken)
{
paths.EnsureDirectories();
await RecoverPendingAsync(stoppingToken);
await ResetInterruptedAsync(stoppingToken);
await foreach (var assetId in queue.DequeueAllAsync(stoppingToken))
while (!stoppingToken.IsCancellationRequested)
{
try
{
await ProcessAsync(assetId, stoppingToken);
// Разобрать всю накопившуюся работу из БД.
while (await NextPendingIdAsync(stoppingToken) is { } assetId)
await ProcessAsync(assetId, stoppingToken);
// Работы нет — ждём сигнала о новой либо периодического опроса.
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)
{
@@ -42,30 +59,41 @@ public sealed class MediaProcessingBackgroundService(
}
catch (Exception ex)
{
logger.LogError(ex, "Необработанная ошибка обработки ассета {AssetId}", assetId);
logger.LogError(ex, "Ошибка цикла обработки медиа");
await Task.Delay(TimeSpan.FromSeconds(5), stoppingToken);
}
}
}
private async Task RecoverPendingAsync(CancellationToken cancellationToken)
/// <summary>Сброс прерванных рестартом задач (Processing → Pending) на старте.</summary>
private async Task ResetInterruptedAsync(CancellationToken cancellationToken)
{
await using var scope = scopeFactory.CreateAsyncScope();
var db = scope.ServiceProvider.GetRequiredService<IAppDbContext>();
var pending = await db.MediaAssets
.Where(x =>
x.Status == MediaAssetStatus.Pending || x.Status == MediaAssetStatus.Processing
)
var interrupted = await db.MediaAssets
.Where(x => x.Status == MediaAssetStatus.Processing)
.ToListAsync(cancellationToken);
if (interrupted.Count == 0)
return;
foreach (var asset in pending.Where(x => x.Status == MediaAssetStatus.Processing))
foreach (var asset in interrupted)
asset.ResetToPending();
await db.SaveChangesAsync(cancellationToken);
}
if (pending.Count > 0)
await db.SaveChangesAsync(cancellationToken);
/// <summary>Id самого раннего ассета в статусе Pending, либо null если работы нет.</summary>
private async Task<Guid?> NextPendingIdAsync(CancellationToken cancellationToken)
{
await using var scope = scopeFactory.CreateAsyncScope();
var db = scope.ServiceProvider.GetRequiredService<IAppDbContext>();
foreach (var asset in pending)
queue.Enqueue(asset.Id);
var id = await db.MediaAssets
.Where(x => x.Status == MediaAssetStatus.Pending)
.OrderBy(x => x.CreatedAt)
.Select(x => (Guid?)x.Id)
.FirstOrDefaultAsync(cancellationToken);
return id;
}
private async Task ProcessAsync(Guid assetId, CancellationToken cancellationToken)
@@ -107,7 +135,7 @@ public sealed class MediaProcessingBackgroundService(
var db = scope.ServiceProvider.GetRequiredService<IAppDbContext>();
var asset = await db.MediaAssets.FirstOrDefaultAsync(x => x.Id == assetId, cancellationToken);
if (asset is null || asset.Status is MediaAssetStatus.Ready or MediaAssetStatus.Failed)
if (asset is null || asset.Status != MediaAssetStatus.Pending)
return null;
asset.MarkProcessing();
@@ -3,7 +3,7 @@ using TeleWave.Application.Common.Interfaces;
namespace TeleWave.Infrastructure.Media;
/// <summary>Неограниченная in-memory очередь id ассетов на обработку (один потребитель).</summary>
/// <summary>Сигнальная очередь-будильник поверх Channel (id ассета используется лишь как сигнал).</summary>
public sealed class MediaProcessingQueue : IMediaProcessingQueue
{
private readonly Channel<Guid> _channel = Channel.CreateUnbounded<Guid>(
@@ -12,6 +12,10 @@ public sealed class MediaProcessingQueue : IMediaProcessingQueue
public void Enqueue(Guid assetId) => _channel.Writer.TryWrite(assetId);
public IAsyncEnumerable<Guid> DequeueAllAsync(CancellationToken cancellationToken) =>
_channel.Reader.ReadAllAsync(cancellationToken);
public async ValueTask WaitAsync(CancellationToken cancellationToken)
{
await _channel.Reader.ReadAsync(cancellationToken);
// Сдренировать накопившиеся сигналы — работу всё равно берём из БД пачкой.
while (_channel.Reader.TryRead(out _)) { }
}
}