fix mysql playout builds (#2352)

* more cancellation tokens and fixes

* so much cancellation token

* fix mysql playout builds
This commit is contained in:
Jason Dove
2025-08-27 18:09:56 +00:00
committed by GitHub
parent 1c07df5bc3
commit 9462156148
103 changed files with 2628 additions and 1666 deletions
@@ -74,7 +74,7 @@ public class SynchronizeEmbyShowByIdHandler : IRequestHandler<SynchronizeEmbySho
SynchronizeEmbyShowById request,
CancellationToken cancellationToken) =>
(await ValidateConnection(request), await EmbyLibraryMustExist(request, cancellationToken),
await EmbyShowMustExist(request))
await EmbyShowMustExist(request, cancellationToken))
.Apply((connectionParameters, embyLibrary, showTitleItemId) =>
new RequestParameters(
connectionParameters,
@@ -121,8 +121,9 @@ public class SynchronizeEmbyShowByIdHandler : IRequestHandler<SynchronizeEmbySho
.Map(v => v.ToValidation<BaseError>($"Emby library {request.EmbyLibraryId} does not exist."));
private Task<Validation<BaseError, EmbyShowTitleItemIdResult>> EmbyShowMustExist(
SynchronizeEmbyShowById request) =>
_embyTelevisionRepository.GetShowTitleItemId(request.EmbyLibraryId, request.ShowId)
SynchronizeEmbyShowById request,
CancellationToken cancellationToken) =>
_embyTelevisionRepository.GetShowTitleItemId(request.EmbyLibraryId, request.ShowId, cancellationToken)
.Map(v => v.ToValidation<BaseError>(
$"Jellyfin show {request.ShowId} does not exist in library {request.EmbyLibraryId}."));
@@ -34,7 +34,7 @@ public class
SynchronizeJellyfinShowById request,
CancellationToken cancellationToken)
{
Validation<BaseError, RequestParameters> validation = await Validate(request);
Validation<BaseError, RequestParameters> validation = await Validate(request, cancellationToken);
return await validation.Match(
parameters => Synchronize(parameters, cancellationToken),
error => Task.FromResult<Either<BaseError, string>>(error.Join()));
@@ -71,9 +71,11 @@ public class
return result.Map(_ => $"Show '{parameters.ShowTitle}' in {parameters.Library.Name}");
}
private async Task<Validation<BaseError, RequestParameters>> Validate(SynchronizeJellyfinShowById request) =>
private async Task<Validation<BaseError, RequestParameters>> Validate(
SynchronizeJellyfinShowById request,
CancellationToken cancellationToken) =>
(await ValidateConnection(request), await JellyfinLibraryMustExist(request),
await JellyfinShowMustExist(request))
await JellyfinShowMustExist(request, cancellationToken))
.Apply((connectionParameters, jellyfinLibrary, showTitleItemId) =>
new RequestParameters(
connectionParameters,
@@ -119,8 +121,9 @@ public class
.Map(v => v.ToValidation<BaseError>($"Jellyfin library {request.JellyfinLibraryId} does not exist."));
private Task<Validation<BaseError, JellyfinShowTitleItemIdResult>> JellyfinShowMustExist(
SynchronizeJellyfinShowById request) =>
_jellyfinTelevisionRepository.GetShowTitleItemId(request.JellyfinLibraryId, request.ShowId)
SynchronizeJellyfinShowById request,
CancellationToken cancellationToken) =>
_jellyfinTelevisionRepository.GetShowTitleItemId(request.JellyfinLibraryId, request.ShowId, cancellationToken)
.Map(v => v.ToValidation<BaseError>(
$"Jellyfin show {request.ShowId} does not exist in library {request.JellyfinLibraryId}."));
@@ -107,7 +107,7 @@ public class SynchronizePlexNetworksHandler : IRequestHandler<SynchronizePlexNet
if (result.IsRight)
{
parameters.Library.LastNetworksScan = DateTime.UtcNow;
await _plexTelevisionRepository.UpdateLastNetworksScan(parameters.Library);
await _plexTelevisionRepository.UpdateLastNetworksScan(parameters.Library, cancellationToken);
}
return result;
@@ -33,7 +33,7 @@ public class SynchronizePlexShowByIdHandler : IRequestHandler<SynchronizePlexSho
SynchronizePlexShowById request,
CancellationToken cancellationToken)
{
Validation<BaseError, RequestParameters> validation = await Validate(request);
Validation<BaseError, RequestParameters> validation = await Validate(request, cancellationToken);
return await validation.Match(
parameters => Synchronize(parameters, cancellationToken),
error => Task.FromResult<Either<BaseError, string>>(error.Join()));
@@ -70,8 +70,11 @@ public class SynchronizePlexShowByIdHandler : IRequestHandler<SynchronizePlexSho
return result.Map(_ => $"Show '{parameters.ShowTitle}' in {parameters.Library.Name}");
}
private async Task<Validation<BaseError, RequestParameters>> Validate(SynchronizePlexShowById request) =>
(await ValidateConnection(request), await PlexLibraryMustExist(request), await PlexShowMustExist(request))
private async Task<Validation<BaseError, RequestParameters>> Validate(
SynchronizePlexShowById request,
CancellationToken cancellationToken) =>
(await ValidateConnection(request), await PlexLibraryMustExist(request),
await PlexShowMustExist(request, cancellationToken))
.Apply((connectionParameters, plexLibrary, titleKey) =>
new RequestParameters(
connectionParameters,
@@ -117,8 +120,9 @@ public class SynchronizePlexShowByIdHandler : IRequestHandler<SynchronizePlexSho
.Map(v => v.ToValidation<BaseError>($"Plex library {request.PlexLibraryId} does not exist."));
private Task<Validation<BaseError, PlexShowTitleKeyResult>> PlexShowMustExist(
SynchronizePlexShowById request) =>
_plexTelevisionRepository.GetShowTitleKey(request.PlexLibraryId, request.ShowId)
SynchronizePlexShowById request,
CancellationToken cancellationToken) =>
_plexTelevisionRepository.GetShowTitleKey(request.PlexLibraryId, request.ShowId, cancellationToken)
.Map(v => v.ToValidation<BaseError>(
$"Plex show {request.ShowId} does not exist in library {request.PlexLibraryId}."));
@@ -14,7 +14,7 @@ public interface ILocalMetadataProvider
Task<bool> RefreshSidecarMetadata(OtherVideo otherVideo, string nfoFileName);
Task<bool> RefreshTagMetadata(Song song);
Task<bool> RefreshTagMetadata(Image image, double? durationSeconds);
Task<bool> RefreshTagMetadata(RemoteStream remoteStream);
Task<bool> RefreshTagMetadata(RemoteStream remoteStream, CancellationToken cancellationToken);
Task<bool> RefreshFallbackMetadata(Movie movie);
Task<bool> RefreshFallbackMetadata(Episode episode);
Task<bool> RefreshFallbackMetadata(Artist artist, string artistFolder);
@@ -22,6 +22,6 @@ public interface ILocalMetadataProvider
Task<bool> RefreshFallbackMetadata(OtherVideo otherVideo);
Task<bool> RefreshFallbackMetadata(Song song);
Task<bool> RefreshFallbackMetadata(Image image);
Task<bool> RefreshFallbackMetadata(RemoteStream remoteStream);
Task<bool> RefreshFallbackMetadata(RemoteStream remoteStream, CancellationToken cancellationToken);
Task<bool> RefreshFallbackMetadata(Show televisionShow, string showFolder);
}
@@ -224,13 +224,13 @@ public class LocalMetadataProvider : ILocalMetadataProvider
return await RefreshFallbackMetadata(image);
}
public async Task<bool> RefreshTagMetadata(RemoteStream remoteStream) =>
public async Task<bool> RefreshTagMetadata(RemoteStream remoteStream, CancellationToken cancellationToken) =>
// Option<RemoteStreamMetadata> maybeMetadata = LoadRemoteStreamMetadata(remoteStream);
// foreach (RemoteStreamMetadata metadata in maybeMetadata)
// {
// return await ApplyMetadataUpdate(remoteStream, metadata);
// }
await RefreshFallbackMetadata(remoteStream);
await RefreshFallbackMetadata(remoteStream, cancellationToken);
public Task<bool> RefreshFallbackMetadata(Movie movie) =>
ApplyMetadataUpdate(movie, _fallbackMetadataProvider.GetFallbackMetadata(movie));
@@ -274,12 +274,12 @@ public class LocalMetadataProvider : ILocalMetadataProvider
return false;
}
public async Task<bool> RefreshFallbackMetadata(RemoteStream remoteStream)
public async Task<bool> RefreshFallbackMetadata(RemoteStream remoteStream, CancellationToken cancellationToken)
{
Option<RemoteStreamMetadata> maybeMetadata = _fallbackMetadataProvider.GetFallbackMetadata(remoteStream);
foreach (RemoteStreamMetadata metadata in maybeMetadata)
{
return await ApplyMetadataUpdate(remoteStream, metadata);
return await ApplyMetadataUpdate(remoteStream, metadata, cancellationToken);
}
return false;
@@ -1239,7 +1239,10 @@ public class LocalMetadataProvider : ILocalMetadataProvider
return await _metadataRepository.Add(metadata);
}
private async Task<bool> ApplyMetadataUpdate(RemoteStream remoteStream, RemoteStreamMetadata metadata)
private async Task<bool> ApplyMetadataUpdate(
RemoteStream remoteStream,
RemoteStreamMetadata metadata,
CancellationToken cancellationToken)
{
Option<RemoteStreamMetadata> maybeMetadata = Optional(remoteStream.RemoteStreamMetadata).Flatten().HeadOrNone();
foreach (RemoteStreamMetadata existing in maybeMetadata)
@@ -1265,7 +1268,7 @@ public class LocalMetadataProvider : ILocalMetadataProvider
existing,
metadata,
(_, _) => Task.FromResult(false),
_remoteStreamRepository.AddTag,
(metadata1, tag) => _remoteStreamRepository.AddTag(metadata1, tag, cancellationToken),
(_, _) => Task.FromResult(false),
(_, _) => Task.FromResult(false));
@@ -91,7 +91,7 @@ public abstract class MediaServerTelevisionLibraryScanner<TConnectionParameters,
CancellationToken cancellationToken)
{
var incomingItemIds = new List<string>();
List<TEtag> existingShows = await televisionRepository.GetExistingShows(library);
List<TEtag> existingShows = await televisionRepository.GetExistingShows(library, cancellationToken);
await foreach ((TShow incoming, int totalShowCount) in showEntries.WithCancellation(cancellationToken))
{
@@ -147,9 +147,9 @@ public abstract class MediaServerTelevisionLibraryScanner<TConnectionParameters,
return error;
}
await televisionRepository.SetEtag(result.Item, MediaServerEtag(incoming));
await televisionRepository.SetEtag(result.Item, MediaServerEtag(incoming), cancellationToken);
Option<int> flagResult = await televisionRepository.FlagNormal(library, result.Item);
Option<int> flagResult = await televisionRepository.FlagNormal(library, result.Item, cancellationToken);
if (flagResult.IsSome)
{
result.IsUpdated = true;
@@ -173,7 +173,7 @@ public abstract class MediaServerTelevisionLibraryScanner<TConnectionParameters,
{
// trash shows that are no longer present on the media server
var fileNotFoundItemIds = existingShows.Map(s => s.MediaServerItemId).Except(incomingItemIds).ToList();
List<int> ids = await televisionRepository.FlagFileNotFoundShows(library, fileNotFoundItemIds);
List<int> ids = await televisionRepository.FlagFileNotFoundShows(library, fileNotFoundItemIds, cancellationToken);
await _mediator.Publish(
new ScannerProgressUpdate(library.Id, null, null, ids.ToArray(), Array.Empty<int>()),
cancellationToken);
@@ -296,7 +296,7 @@ public abstract class MediaServerTelevisionLibraryScanner<TConnectionParameters,
CancellationToken cancellationToken)
{
var incomingItemIds = new List<string>();
List<TEtag> existingSeasons = await televisionRepository.GetExistingSeasons(library, show);
List<TEtag> existingSeasons = await televisionRepository.GetExistingSeasons(library, show, cancellationToken);
await foreach ((TSeason incoming, int _) in seasonEntries.WithCancellation(cancellationToken))
{
@@ -346,9 +346,9 @@ public abstract class MediaServerTelevisionLibraryScanner<TConnectionParameters,
return error;
}
await televisionRepository.SetEtag(result.Item, MediaServerEtag(incoming));
await televisionRepository.SetEtag(result.Item, MediaServerEtag(incoming), cancellationToken);
Option<int> flagResult = await televisionRepository.FlagNormal(library, result.Item);
Option<int> flagResult = await televisionRepository.FlagNormal(library, result.Item, cancellationToken);
if (flagResult.IsSome)
{
result.IsUpdated = true;
@@ -372,7 +372,7 @@ public abstract class MediaServerTelevisionLibraryScanner<TConnectionParameters,
// trash seasons that are no longer present on the media server
var fileNotFoundItemIds = existingSeasons.Map(s => s.MediaServerItemId).Except(incomingItemIds).ToList();
List<int> ids = await televisionRepository.FlagFileNotFoundSeasons(library, fileNotFoundItemIds);
List<int> ids = await televisionRepository.FlagFileNotFoundSeasons(library, fileNotFoundItemIds, cancellationToken);
await _mediator.Publish(
new ScannerProgressUpdate(library.Id, null, null, ids.ToArray(), Array.Empty<int>()),
cancellationToken);
@@ -393,7 +393,7 @@ public abstract class MediaServerTelevisionLibraryScanner<TConnectionParameters,
CancellationToken cancellationToken)
{
var incomingItemIds = new List<string>();
List<TEtag> existingEpisodes = await televisionRepository.GetExistingEpisodes(library, season);
List<TEtag> existingEpisodes = await televisionRepository.GetExistingEpisodes(library, season, cancellationToken);
await foreach ((TEpisode incoming, int _) in episodeEntries.WithCancellation(cancellationToken))
{
@@ -413,7 +413,8 @@ public abstract class MediaServerTelevisionLibraryScanner<TConnectionParameters,
existingEpisodes,
incoming,
localPath,
deepScan))
deepScan,
cancellationToken))
{
continue;
}
@@ -485,11 +486,11 @@ public abstract class MediaServerTelevisionLibraryScanner<TConnectionParameters,
foreach (MediaItemScanResult<TEpisode> result in maybeEpisode.RightToSeq())
{
await televisionRepository.SetEtag(result.Item, MediaServerEtag(incoming));
await televisionRepository.SetEtag(result.Item, MediaServerEtag(incoming), cancellationToken);
if (_localFileSystem.FileExists(result.LocalPath))
{
Option<int> flagResult = await televisionRepository.FlagNormal(library, result.Item);
Option<int> flagResult = await televisionRepository.FlagNormal(library, result.Item, cancellationToken);
if (flagResult.IsSome)
{
result.IsUpdated = true;
@@ -497,7 +498,7 @@ public abstract class MediaServerTelevisionLibraryScanner<TConnectionParameters,
}
else if (ServerSupportsRemoteStreaming)
{
Option<int> flagResult = await televisionRepository.FlagRemoteOnly(library, result.Item);
Option<int> flagResult = await televisionRepository.FlagRemoteOnly(library, result.Item, cancellationToken);
if (flagResult.IsSome)
{
result.IsUpdated = true;
@@ -505,7 +506,7 @@ public abstract class MediaServerTelevisionLibraryScanner<TConnectionParameters,
}
else
{
Option<int> flagResult = await televisionRepository.FlagUnavailable(library, result.Item);
Option<int> flagResult = await televisionRepository.FlagUnavailable(library, result.Item, cancellationToken);
if (flagResult.IsSome)
{
result.IsUpdated = true;
@@ -528,7 +529,7 @@ public abstract class MediaServerTelevisionLibraryScanner<TConnectionParameters,
// trash episodes that are no longer present on the media server
var fileNotFoundItemIds = existingEpisodes.Map(m => m.MediaServerItemId).Except(incomingItemIds).ToList();
List<int> ids = await televisionRepository.FlagFileNotFoundEpisodes(library, fileNotFoundItemIds);
List<int> ids = await televisionRepository.FlagFileNotFoundEpisodes(library, fileNotFoundItemIds, cancellationToken);
await _mediator.Publish(
new ScannerProgressUpdate(library.Id, null, null, ids.ToArray(), Array.Empty<int>()),
cancellationToken);
@@ -544,7 +545,8 @@ public abstract class MediaServerTelevisionLibraryScanner<TConnectionParameters,
List<TEtag> existingEpisodes,
TEpisode incoming,
string localPath,
bool deepScan)
bool deepScan,
CancellationToken cancellationToken)
{
// deep scan will always pull every episode
if (deepScan)
@@ -575,10 +577,10 @@ public abstract class MediaServerTelevisionLibraryScanner<TConnectionParameters,
{
if (existingState is not MediaItemState.RemoteOnly)
{
foreach (int id in await televisionRepository.FlagRemoteOnly(library, incoming))
foreach (int id in await televisionRepository.FlagRemoteOnly(library, incoming, cancellationToken))
{
await _mediator.Publish(
new ScannerProgressUpdate(library.Id, null, null, new[] { id }, Array.Empty<int>()),
new ScannerProgressUpdate(library.Id, null, null, [id], []),
CancellationToken.None);
}
}
@@ -587,10 +589,10 @@ public abstract class MediaServerTelevisionLibraryScanner<TConnectionParameters,
{
if (existingState is not MediaItemState.Unavailable)
{
foreach (int id in await televisionRepository.FlagUnavailable(library, incoming))
foreach (int id in await televisionRepository.FlagUnavailable(library, incoming, cancellationToken))
{
await _mediator.Publish(
new ScannerProgressUpdate(library.Id, null, null, new[] { id }, Array.Empty<int>()),
new ScannerProgressUpdate(library.Id, null, null, [id], []),
CancellationToken.None);
}
}
@@ -182,7 +182,7 @@ public class RemoteStreamFolderScanner : LocalFolderScanner, IRemoteStreamFolder
.BindT(video => ParseRemoteStreamDefinition(video, deserializer, cancellationToken))
.BindT(video => UpdateStatistics(video, ffmpegPath, ffprobePath))
.BindT(video => UpdateLibraryFolderId(video, knownFolder))
.BindT(UpdateMetadata)
.BindT(video => UpdateMetadata(video, cancellationToken))
//.BindT(video => UpdateThumbnail(video, cancellationToken))
//.BindT(UpdateSubtitles)
.BindT(FlagNormal);
@@ -216,7 +216,7 @@ public class RemoteStreamFolderScanner : LocalFolderScanner, IRemoteStreamFolder
}
}
foreach (string path in await _remoteStreamRepository.FindRemoteStreamPaths(libraryPath))
foreach (string path in await _remoteStreamRepository.FindRemoteStreamPaths(libraryPath, cancellationToken))
{
if (!_localFileSystem.FileExists(path))
{
@@ -234,7 +234,7 @@ public class RemoteStreamFolderScanner : LocalFolderScanner, IRemoteStreamFolder
else if (Path.GetFileName(path).StartsWith("._", StringComparison.OrdinalIgnoreCase))
{
_logger.LogInformation("Removing dot underscore file at {Path}", path);
List<int> remoteStreamIds = await _remoteStreamRepository.DeleteByPath(libraryPath, path);
List<int> remoteStreamIds = await _remoteStreamRepository.DeleteByPath(libraryPath, path, cancellationToken);
await _mediator.Publish(
new ScannerProgressUpdate(
libraryPath.LibraryId,
@@ -336,7 +336,7 @@ public class RemoteStreamFolderScanner : LocalFolderScanner, IRemoteStreamFolder
if (updated)
{
await _remoteStreamRepository.UpdateDefinition(remoteStream);
await _remoteStreamRepository.UpdateDefinition(remoteStream, cancellationToken);
result.IsUpdated = true;
}
@@ -350,7 +350,8 @@ public class RemoteStreamFolderScanner : LocalFolderScanner, IRemoteStreamFolder
}
private async Task<Either<BaseError, MediaItemScanResult<RemoteStream>>> UpdateMetadata(
MediaItemScanResult<RemoteStream> result)
MediaItemScanResult<RemoteStream> result,
CancellationToken cancellationToken)
{
try
{
@@ -371,7 +372,7 @@ public class RemoteStreamFolderScanner : LocalFolderScanner, IRemoteStreamFolder
remoteStream.RemoteStreamMetadata ??= [];
_logger.LogDebug("Refreshing {Attribute} for {Path}", "Metadata", path);
if (await _localMetadataProvider.RefreshTagMetadata(remoteStream))
if (await _localMetadataProvider.RefreshTagMetadata(remoteStream, cancellationToken))
{
result.IsUpdated = true;
}
@@ -58,7 +58,7 @@ public class PlexNetworkScanner(
var keepIds = new System.Collections.Generic.HashSet<int>();
await foreach ((PlexShow item, int _) in items)
{
PlexShowAddTagResult result = await plexTelevisionRepository.AddTag(library, item, tag);
PlexShowAddTagResult result = await plexTelevisionRepository.AddTag(library, item, tag, cancellationToken);
foreach (int existing in result.Existing)
{
@@ -74,7 +74,7 @@ public class PlexNetworkScanner(
cancellationToken.ThrowIfCancellationRequested();
}
List<int> removedIds = await plexTelevisionRepository.RemoveAllTags(library, tag, keepIds);
List<int> removedIds = await plexTelevisionRepository.RemoveAllTags(library, tag, keepIds, cancellationToken);
var changedIds = removedIds.Concat(addedIds).Distinct().ToList();
if (changedIds.Count > 0)