using System.Threading.Channels; using ErsatzTV.Application.Playouts; using ErsatzTV.Core; using ErsatzTV.Core.Domain; using ErsatzTV.Core.Interfaces.Repositories; using ErsatzTV.Core.Scheduling; using ErsatzTV.Infrastructure.Data; using ErsatzTV.Infrastructure.Extensions; using Microsoft.EntityFrameworkCore; namespace ErsatzTV.Application.MediaCollections; public class UpdateCollectionCustomOrderHandler : IRequestHandler> { private readonly ChannelWriter _channel; private readonly IDbContextFactory _dbContextFactory; private readonly IMediaCollectionRepository _mediaCollectionRepository; public UpdateCollectionCustomOrderHandler( IDbContextFactory dbContextFactory, IMediaCollectionRepository mediaCollectionRepository, ChannelWriter channel) { _dbContextFactory = dbContextFactory; _mediaCollectionRepository = mediaCollectionRepository; _channel = channel; } public async Task> Handle( UpdateCollectionCustomOrder request, CancellationToken cancellationToken) { await using TvContext dbContext = await _dbContextFactory.CreateDbContextAsync(cancellationToken); Validation validation = await Validate(dbContext, request, cancellationToken); // Optimistic-concurrency check as a standalone Either AFTER validation, never via Apply (which // Join()s the error Seq and would flatten PreconditionFailedError to a 422) — issue #253 §7a. Either validated = LanguageExtensions.ToEither(validation) .Bind(c => c.CheckVersion(request.ExpectedVersions)); return await validated.Match( Right: c => ApplyUpdateRequest(dbContext, c, request, cancellationToken), Left: error => Task.FromResult>(error)); } private async Task> ApplyUpdateRequest( TvContext dbContext, Collection c, UpdateCollectionCustomOrder request, CancellationToken cancellationToken) { foreach (MediaItemCustomOrder updateItem in request.MediaItemCustomOrders) { Option maybeCollectionItem = c.CollectionItems .FirstOrDefault(ci => ci.MediaItemId == updateItem.MediaItemId); foreach (CollectionItem collectionItem in maybeCollectionItem) { collectionItem.CustomIndex = updateItem.CustomIndex; } } // Unconditional bump (issue #253 §7a / M1) then guarded save (→ 412 on a lost race). c.Version++; Either saveResult = await dbContext.SaveChangesWithConcurrencyGuard(cancellationToken); if (saveResult.IsLeft) { return saveResult; } // Refresh all playouts that use this collection. The old `SaveChangesAsync() > 0` gate is always // true once the version bumps unconditionally (M2), so run the refresh on any successful save. // Post-commit enqueue on CancellationToken.None so a late cancellation can't drop the rebuild // after the commit landed (#254 / §7b). foreach (int playoutId in await _mediaCollectionRepository .PlayoutIdsUsingCollection(request.CollectionId)) { await _channel.WriteAsync( new BuildPlayout(playoutId, PlayoutBuildMode.Refresh), CancellationToken.None); } return Unit.Default; } private static Task> Validate( TvContext dbContext, UpdateCollectionCustomOrder request, CancellationToken cancellationToken) => CollectionMustExist(dbContext, request, cancellationToken); private static Task> CollectionMustExist( TvContext dbContext, UpdateCollectionCustomOrder request, CancellationToken cancellationToken) => dbContext.Collections .Include(c => c.CollectionItems) .SelectOneAsync(c => c.Id, c => c.Id == request.CollectionId, cancellationToken) .Map(o => o.ToValidation("Collection does not exist.")); }