use cancellation tokens in many places (#2350)

* use cancellation tokens everywhere

* more cancellation tokens
This commit is contained in:
Jason Dove
2025-08-27 03:20:35 +00:00
committed by GitHub
parent a6198892f0
commit 1c07df5bc3
356 changed files with 3430 additions and 2454 deletions
@@ -29,19 +29,21 @@ public class SynchronizeEmbyCollectionsHandler : IRequestHandler<SynchronizeEmby
SynchronizeEmbyCollections request,
CancellationToken cancellationToken)
{
Validation<BaseError, RequestParameters> validation = await Validate(request);
Validation<BaseError, RequestParameters> validation = await Validate(request, cancellationToken);
return await validation.Match(
SynchronizeCollections,
error => Task.FromResult<Either<BaseError, Unit>>(error.Join()));
}
private async Task<Validation<BaseError, RequestParameters>> Validate(SynchronizeEmbyCollections request)
private async Task<Validation<BaseError, RequestParameters>> Validate(
SynchronizeEmbyCollections request,
CancellationToken cancellationToken)
{
Task<Validation<BaseError, ConnectionParameters>> mediaSource = MediaSourceMustExist(request)
Task<Validation<BaseError, ConnectionParameters>> mediaSource = MediaSourceMustExist(request, cancellationToken)
.BindT(MediaSourceMustHaveActiveConnection)
.BindT(MediaSourceMustHaveApiKey);
return (await mediaSource, await ValidateLibraryRefreshInterval())
return (await mediaSource, await ValidateLibraryRefreshInterval(cancellationToken))
.Apply((connectionParameters, libraryRefreshInterval) => new RequestParameters(
connectionParameters,
connectionParameters.MediaSource,
@@ -49,14 +51,15 @@ public class SynchronizeEmbyCollectionsHandler : IRequestHandler<SynchronizeEmby
libraryRefreshInterval));
}
private Task<Validation<BaseError, int>> ValidateLibraryRefreshInterval() =>
_configElementRepository.GetValue<int>(ConfigElementKey.LibraryRefreshInterval)
private Task<Validation<BaseError, int>> ValidateLibraryRefreshInterval(CancellationToken cancellationToken) =>
_configElementRepository.GetValue<int>(ConfigElementKey.LibraryRefreshInterval, cancellationToken)
.FilterT(lri => lri is >= 0 and < 1_000_000)
.Map(lri => lri.ToValidation<BaseError>("Library refresh interval is invalid"));
private Task<Validation<BaseError, EmbyMediaSource>> MediaSourceMustExist(
SynchronizeEmbyCollections request) =>
_mediaSourceRepository.GetEmby(request.EmbyMediaSourceId)
SynchronizeEmbyCollections request,
CancellationToken cancellationToken) =>
_mediaSourceRepository.GetEmby(request.EmbyMediaSourceId, cancellationToken)
.Map(o => o.ToValidation<BaseError>("Emby media source does not exist."));
private static Validation<BaseError, ConnectionParameters> MediaSourceMustHaveActiveConnection(
@@ -44,7 +44,7 @@ public class SynchronizeEmbyLibraryByIdHandler : IRequestHandler<SynchronizeEmby
public async Task<Either<BaseError, string>>
Handle(SynchronizeEmbyLibraryById 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()));
@@ -107,8 +107,10 @@ public class SynchronizeEmbyLibraryByIdHandler : IRequestHandler<SynchronizeEmby
}
private async Task<Validation<BaseError, RequestParameters>> Validate(
SynchronizeEmbyLibraryById request) =>
(await ValidateConnection(request), await EmbyLibraryMustExist(request), await ValidateLibraryRefreshInterval())
SynchronizeEmbyLibraryById request,
CancellationToken cancellationToken) =>
(await ValidateConnection(request), await EmbyLibraryMustExist(request, cancellationToken),
await ValidateLibraryRefreshInterval(cancellationToken))
.Apply((connectionParameters, embyLibrary, libraryRefreshInterval) =>
new RequestParameters(
connectionParameters,
@@ -149,12 +151,13 @@ public class SynchronizeEmbyLibraryByIdHandler : IRequestHandler<SynchronizeEmby
}
private Task<Validation<BaseError, EmbyLibrary>> EmbyLibraryMustExist(
SynchronizeEmbyLibraryById request) =>
_mediaSourceRepository.GetEmbyLibrary(request.EmbyLibraryId)
SynchronizeEmbyLibraryById request,
CancellationToken cancellationToken) =>
_mediaSourceRepository.GetEmbyLibrary(request.EmbyLibraryId, cancellationToken)
.Map(v => v.ToValidation<BaseError>($"Emby library {request.EmbyLibraryId} does not exist."));
private Task<Validation<BaseError, int>> ValidateLibraryRefreshInterval() =>
_configElementRepository.GetValue<int>(ConfigElementKey.LibraryRefreshInterval)
private Task<Validation<BaseError, int>> ValidateLibraryRefreshInterval(CancellationToken cancellationToken) =>
_configElementRepository.GetValue<int>(ConfigElementKey.LibraryRefreshInterval, cancellationToken)
.FilterT(lri => lri is >= 0 and < 1_000_000)
.Map(lri => lri.ToValidation<BaseError>("Library refresh interval is invalid"));
@@ -33,7 +33,7 @@ public class SynchronizeEmbyShowByIdHandler : IRequestHandler<SynchronizeEmbySho
SynchronizeEmbyShowById 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 SynchronizeEmbyShowByIdHandler : IRequestHandler<SynchronizeEmbySho
return result.Map(_ => $"Show '{parameters.ShowTitle}' in {parameters.Library.Name}");
}
private async Task<Validation<BaseError, RequestParameters>> Validate(SynchronizeEmbyShowById request) =>
(await ValidateConnection(request), await EmbyLibraryMustExist(request), await EmbyShowMustExist(request))
private async Task<Validation<BaseError, RequestParameters>> Validate(
SynchronizeEmbyShowById request,
CancellationToken cancellationToken) =>
(await ValidateConnection(request), await EmbyLibraryMustExist(request, cancellationToken),
await EmbyShowMustExist(request))
.Apply((connectionParameters, embyLibrary, showTitleItemId) =>
new RequestParameters(
connectionParameters,
@@ -112,8 +115,9 @@ public class SynchronizeEmbyShowByIdHandler : IRequestHandler<SynchronizeEmbySho
}
private Task<Validation<BaseError, EmbyLibrary>> EmbyLibraryMustExist(
SynchronizeEmbyShowById request) =>
_mediaSourceRepository.GetEmbyLibrary(request.EmbyLibraryId)
SynchronizeEmbyShowById request,
CancellationToken cancellationToken) =>
_mediaSourceRepository.GetEmbyLibrary(request.EmbyLibraryId, cancellationToken)
.Map(v => v.ToValidation<BaseError>($"Emby library {request.EmbyLibraryId} does not exist."));
private Task<Validation<BaseError, EmbyShowTitleItemIdResult>> EmbyShowMustExist(
@@ -31,19 +31,21 @@ public class
SynchronizeJellyfinCollections request,
CancellationToken cancellationToken)
{
Validation<BaseError, RequestParameters> validation = await Validate(request);
Validation<BaseError, RequestParameters> validation = await Validate(request, cancellationToken);
return await validation.Match(
SynchronizeCollections,
error => Task.FromResult<Either<BaseError, Unit>>(error.Join()));
}
private async Task<Validation<BaseError, RequestParameters>> Validate(SynchronizeJellyfinCollections request)
private async Task<Validation<BaseError, RequestParameters>> Validate(
SynchronizeJellyfinCollections request,
CancellationToken cancellationToken)
{
Task<Validation<BaseError, ConnectionParameters>> mediaSource = MediaSourceMustExist(request)
.BindT(MediaSourceMustHaveActiveConnection)
.BindT(MediaSourceMustHaveApiKey);
return (await mediaSource, await ValidateLibraryRefreshInterval())
return (await mediaSource, await ValidateLibraryRefreshInterval(cancellationToken))
.Apply((connectionParameters, libraryRefreshInterval) => new RequestParameters(
connectionParameters,
connectionParameters.MediaSource,
@@ -51,8 +53,8 @@ public class
libraryRefreshInterval));
}
private Task<Validation<BaseError, int>> ValidateLibraryRefreshInterval() =>
_configElementRepository.GetValue<int>(ConfigElementKey.LibraryRefreshInterval)
private Task<Validation<BaseError, int>> ValidateLibraryRefreshInterval(CancellationToken cancellationToken) =>
_configElementRepository.GetValue<int>(ConfigElementKey.LibraryRefreshInterval, cancellationToken)
.FilterT(lri => lri is >= 0 and < 1_000_000)
.Map(lri => lri.ToValidation<BaseError>("Library refresh interval is invalid"));
@@ -45,7 +45,7 @@ public class
public async Task<Either<BaseError, string>>
Handle(SynchronizeJellyfinLibraryById 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()));
@@ -108,9 +108,10 @@ public class
}
private async Task<Validation<BaseError, RequestParameters>> Validate(
SynchronizeJellyfinLibraryById request) =>
SynchronizeJellyfinLibraryById request,
CancellationToken cancellationToken) =>
(await ValidateConnection(request), await JellyfinLibraryMustExist(request),
await ValidateLibraryRefreshInterval())
await ValidateLibraryRefreshInterval(cancellationToken))
.Apply((connectionParameters, jellyfinLibrary, libraryRefreshInterval) =>
new RequestParameters(
connectionParameters,
@@ -155,8 +156,8 @@ public class
_mediaSourceRepository.GetJellyfinLibrary(request.JellyfinLibraryId)
.Map(v => v.ToValidation<BaseError>($"Jellyfin library {request.JellyfinLibraryId} does not exist."));
private Task<Validation<BaseError, int>> ValidateLibraryRefreshInterval() =>
_configElementRepository.GetValue<int>(ConfigElementKey.LibraryRefreshInterval)
private Task<Validation<BaseError, int>> ValidateLibraryRefreshInterval(CancellationToken cancellationToken) =>
_configElementRepository.GetValue<int>(ConfigElementKey.LibraryRefreshInterval, cancellationToken)
.FilterT(lri => lri is >= 0 and < 1_000_000)
.Map(lri => lri.ToValidation<BaseError>("Library refresh interval is invalid"));
@@ -50,7 +50,7 @@ public class ScanLocalLibraryHandler : IRequestHandler<ScanLocalLibrary, Either<
}
public Task<Either<BaseError, string>> Handle(ScanLocalLibrary request, CancellationToken cancellationToken) =>
Validate(request)
Validate(request, cancellationToken)
.MapT(parameters => PerformScan(parameters, cancellationToken).Map(_ => parameters.LocalLibrary.Name))
.Bind(v => v.ToEitherAsync());
@@ -172,18 +172,20 @@ public class ScanLocalLibraryHandler : IRequestHandler<ScanLocalLibrary, Either<
}
await _mediator.Publish(
new ScannerProgressUpdate(localLibrary.Id, localLibrary.Name, 0, Array.Empty<int>(), Array.Empty<int>()),
new ScannerProgressUpdate(localLibrary.Id, localLibrary.Name, 0, [], []),
cancellationToken);
return Unit.Default;
}
private async Task<Validation<BaseError, RequestParameters>> Validate(ScanLocalLibrary request)
private async Task<Validation<BaseError, RequestParameters>> Validate(
ScanLocalLibrary request,
CancellationToken cancellationToken)
{
Validation<BaseError, LocalLibrary> libraryResult = await LocalLibraryMustExist(request);
Validation<BaseError, string> ffprobePathResult = await ValidateFFprobePath();
Validation<BaseError, string> ffmpegPathResult = await ValidateFFmpegPath();
Validation<BaseError, int> refreshIntervalResult = await ValidateLibraryRefreshInterval();
Validation<BaseError, string> ffprobePathResult = await ValidateFFprobePath(cancellationToken);
Validation<BaseError, string> ffmpegPathResult = await ValidateFFmpegPath(cancellationToken);
Validation<BaseError, int> refreshIntervalResult = await ValidateLibraryRefreshInterval(cancellationToken);
return (libraryResult, ffprobePathResult, ffmpegPathResult, refreshIntervalResult)
.Apply((library, ffprobePath, ffmpegPath, libraryRefreshInterval) => new RequestParameters(
@@ -199,20 +201,20 @@ public class ScanLocalLibraryHandler : IRequestHandler<ScanLocalLibrary, Either<
.Map(maybeLibrary => maybeLibrary.OfType<LocalLibrary>().HeadOrNone())
.Map(v => v.ToValidation<BaseError>($"Local library {request.LibraryId} does not exist."));
private Task<Validation<BaseError, string>> ValidateFFprobePath() =>
_configElementRepository.GetValue<string>(ConfigElementKey.FFprobePath)
private Task<Validation<BaseError, string>> ValidateFFprobePath(CancellationToken cancellationToken) =>
_configElementRepository.GetValue<string>(ConfigElementKey.FFprobePath, cancellationToken)
.FilterT(File.Exists)
.Map(ffprobePath =>
ffprobePath.ToValidation<BaseError>("FFprobe path does not exist on the file system"));
private Task<Validation<BaseError, string>> ValidateFFmpegPath() =>
_configElementRepository.GetValue<string>(ConfigElementKey.FFmpegPath)
private Task<Validation<BaseError, string>> ValidateFFmpegPath(CancellationToken cancellationToken) =>
_configElementRepository.GetValue<string>(ConfigElementKey.FFmpegPath, cancellationToken)
.FilterT(File.Exists)
.Map(ffmpegPath =>
ffmpegPath.ToValidation<BaseError>("FFmpeg path does not exist on the file system"));
private Task<Validation<BaseError, int>> ValidateLibraryRefreshInterval() =>
_configElementRepository.GetValue<int>(ConfigElementKey.LibraryRefreshInterval)
private Task<Validation<BaseError, int>> ValidateLibraryRefreshInterval(CancellationToken cancellationToken) =>
_configElementRepository.GetValue<int>(ConfigElementKey.LibraryRefreshInterval, cancellationToken)
.FilterT(lri => lri is >= 0 and < 1_000_000)
.Map(lri => lri.ToValidation<BaseError>("Library refresh interval is invalid"));
@@ -29,19 +29,21 @@ public class SynchronizePlexCollectionsHandler : IRequestHandler<SynchronizePlex
SynchronizePlexCollections request,
CancellationToken cancellationToken)
{
Validation<BaseError, RequestParameters> validation = await Validate(request);
Validation<BaseError, RequestParameters> validation = await Validate(request, cancellationToken);
return await validation.Match(
p => SynchronizeCollections(p, cancellationToken),
error => Task.FromResult<Either<BaseError, Unit>>(error.Join()));
}
private async Task<Validation<BaseError, RequestParameters>> Validate(SynchronizePlexCollections request)
private async Task<Validation<BaseError, RequestParameters>> Validate(
SynchronizePlexCollections request,
CancellationToken cancellationToken)
{
Task<Validation<BaseError, ConnectionParameters>> mediaSource = MediaSourceMustExist(request)
Task<Validation<BaseError, ConnectionParameters>> mediaSource = MediaSourceMustExist(request, cancellationToken)
.BindT(MediaSourceMustHaveActiveConnection)
.BindT(MediaSourceMustHaveToken);
return (await mediaSource, await ValidateLibraryRefreshInterval())
return (await mediaSource, await ValidateLibraryRefreshInterval(cancellationToken))
.Apply((connectionParameters, libraryRefreshInterval) => new RequestParameters(
connectionParameters,
connectionParameters.PlexMediaSource,
@@ -49,14 +51,15 @@ public class SynchronizePlexCollectionsHandler : IRequestHandler<SynchronizePlex
libraryRefreshInterval));
}
private Task<Validation<BaseError, int>> ValidateLibraryRefreshInterval() =>
_configElementRepository.GetValue<int>(ConfigElementKey.LibraryRefreshInterval)
private Task<Validation<BaseError, int>> ValidateLibraryRefreshInterval(CancellationToken cancellationToken) =>
_configElementRepository.GetValue<int>(ConfigElementKey.LibraryRefreshInterval, cancellationToken)
.FilterT(lri => lri is >= 0 and < 1_000_000)
.Map(lri => lri.ToValidation<BaseError>("Library refresh interval is invalid"));
private Task<Validation<BaseError, PlexMediaSource>> MediaSourceMustExist(
SynchronizePlexCollections request) =>
_mediaSourceRepository.GetPlex(request.PlexMediaSourceId)
SynchronizePlexCollections request,
CancellationToken cancellationToken) =>
_mediaSourceRepository.GetPlex(request.PlexMediaSourceId, cancellationToken)
.Map(o => o.ToValidation<BaseError>("Plex media source does not exist."));
private static Validation<BaseError, ConnectionParameters> MediaSourceMustHaveActiveConnection(
@@ -46,7 +46,7 @@ public class SynchronizePlexLibraryByIdHandler : IRequestHandler<SynchronizePlex
SynchronizePlexLibraryById 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()));
@@ -117,8 +117,11 @@ public class SynchronizePlexLibraryByIdHandler : IRequestHandler<SynchronizePlex
return parameters.Library.Name;
}
private async Task<Validation<BaseError, RequestParameters>> Validate(SynchronizePlexLibraryById request) =>
(await ValidateConnection(request), await PlexLibraryMustExist(request), await ValidateLibraryRefreshInterval())
private async Task<Validation<BaseError, RequestParameters>> Validate(
SynchronizePlexLibraryById request,
CancellationToken cancellationToken) =>
(await ValidateConnection(request), await PlexLibraryMustExist(request),
await ValidateLibraryRefreshInterval(cancellationToken))
.Apply((connectionParameters, plexLibrary, libraryRefreshInterval) =>
new RequestParameters(
connectionParameters,
@@ -163,8 +166,8 @@ public class SynchronizePlexLibraryByIdHandler : IRequestHandler<SynchronizePlex
_mediaSourceRepository.GetPlexLibrary(request.PlexLibraryId)
.Map(v => v.ToValidation<BaseError>($"Plex library {request.PlexLibraryId} does not exist."));
private Task<Validation<BaseError, int>> ValidateLibraryRefreshInterval() =>
_configElementRepository.GetValue<int>(ConfigElementKey.LibraryRefreshInterval)
private Task<Validation<BaseError, int>> ValidateLibraryRefreshInterval(CancellationToken cancellationToken) =>
_configElementRepository.GetValue<int>(ConfigElementKey.LibraryRefreshInterval, cancellationToken)
.FilterT(lri => lri is >= 0 and < 1_000_000)
.Map(lri => lri.ToValidation<BaseError>("Library refresh interval is invalid"));
@@ -32,19 +32,22 @@ public class SynchronizePlexNetworksHandler : IRequestHandler<SynchronizePlexNet
SynchronizePlexNetworks request,
CancellationToken cancellationToken)
{
Validation<BaseError, RequestParameters> validation = await Validate(request);
Validation<BaseError, RequestParameters> validation = await Validate(request, cancellationToken);
return await validation.Match(
p => SynchronizeNetworks(p, cancellationToken),
error => Task.FromResult<Either<BaseError, Unit>>(error.Join()));
}
private async Task<Validation<BaseError, RequestParameters>> Validate(SynchronizePlexNetworks request)
private async Task<Validation<BaseError, RequestParameters>> Validate(
SynchronizePlexNetworks request,
CancellationToken cancellationToken)
{
Task<Validation<BaseError, ConnectionParameters>> mediaSource = MediaSourceMustExist(request)
.BindT(MediaSourceMustHaveActiveConnection)
.BindT(MediaSourceMustHaveToken);
return (await mediaSource, await PlexLibraryMustExist(request), await ValidateLibraryRefreshInterval())
return (await mediaSource, await PlexLibraryMustExist(request),
await ValidateLibraryRefreshInterval(cancellationToken))
.Apply((connectionParameters, plexLibrary, libraryRefreshInterval) => new RequestParameters(
connectionParameters,
plexLibrary,
@@ -57,8 +60,8 @@ public class SynchronizePlexNetworksHandler : IRequestHandler<SynchronizePlexNet
_mediaSourceRepository.GetPlexLibrary(request.PlexLibraryId)
.Map(v => v.ToValidation<BaseError>($"Plex library {request.PlexLibraryId} does not exist."));
private Task<Validation<BaseError, int>> ValidateLibraryRefreshInterval() =>
_configElementRepository.GetValue<int>(ConfigElementKey.LibraryRefreshInterval)
private Task<Validation<BaseError, int>> ValidateLibraryRefreshInterval(CancellationToken cancellationToken) =>
_configElementRepository.GetValue<int>(ConfigElementKey.LibraryRefreshInterval, cancellationToken)
.FilterT(lri => lri is >= 0 and < 1_000_000)
.Map(lri => lri.ToValidation<BaseError>("Library refresh interval is invalid"));
@@ -131,6 +131,7 @@ public class EmbyMovieLibraryScanner :
protected override Task<Either<BaseError, MediaItemScanResult<EmbyMovie>>> UpdateMetadata(
MediaItemScanResult<EmbyMovie> result,
MovieMetadata fullMetadata) =>
MovieMetadata fullMetadata,
CancellationToken cancellationToken) =>
Task.FromResult<Either<BaseError, MediaItemScanResult<EmbyMovie>>>(result);
}
@@ -237,7 +237,8 @@ public class EmbyTelevisionLibraryScanner : MediaServerTelevisionLibraryScanner<
protected override Task<Either<BaseError, MediaItemScanResult<EmbyEpisode>>> UpdateMetadata(
MediaItemScanResult<EmbyEpisode> result,
EpisodeMetadata fullMetadata) =>
EpisodeMetadata fullMetadata,
CancellationToken cancellationToken) =>
Task.FromResult<Either<BaseError, MediaItemScanResult<EmbyEpisode>>>(result);
private async Task<Either<BaseError, Unit>> ScanSingleShowInternal(
@@ -4,5 +4,5 @@ namespace ErsatzTV.Scanner.Core.Interfaces.Metadata;
public interface ILocalChaptersProvider : IDisposable
{
Task<bool> UpdateChapters(MediaItem mediaItem, Option<string> localPath);
Task<bool> UpdateChapters(MediaItem mediaItem, Option<string> localPath, CancellationToken cancellationToken);
}
@@ -4,5 +4,9 @@ namespace ErsatzTV.Scanner.Core.Interfaces.Metadata;
public interface ILocalSubtitlesProvider : IDisposable
{
Task<bool> UpdateSubtitles(MediaItem mediaItem, Option<string> localPath, bool saveFullPath);
Task<bool> UpdateSubtitles(
MediaItem mediaItem,
Option<string> localPath,
bool saveFullPath,
CancellationToken cancellationToken);
}
@@ -132,6 +132,7 @@ public class JellyfinMovieLibraryScanner :
protected override Task<Either<BaseError, MediaItemScanResult<JellyfinMovie>>> UpdateMetadata(
MediaItemScanResult<JellyfinMovie> result,
MovieMetadata fullMetadata) =>
MovieMetadata fullMetadata,
CancellationToken cancellationToken) =>
Task.FromResult<Either<BaseError, MediaItemScanResult<JellyfinMovie>>>(result);
}
@@ -236,7 +236,8 @@ public class JellyfinTelevisionLibraryScanner : MediaServerTelevisionLibraryScan
protected override Task<Either<BaseError, MediaItemScanResult<JellyfinEpisode>>> UpdateMetadata(
MediaItemScanResult<JellyfinEpisode> result,
EpisodeMetadata fullMetadata) =>
EpisodeMetadata fullMetadata,
CancellationToken cancellationToken) =>
Task.FromResult<Either<BaseError, MediaItemScanResult<JellyfinEpisode>>>(result);
private async Task<Either<BaseError, Unit>> ScanSingleShowInternal(
@@ -119,7 +119,10 @@ public class ImageFolderScanner : LocalFolderScanner, IImageFolderScanner
cancellationToken);
string imageFolder = folderQueue.Dequeue();
Option<int> maybeParentFolder = await _libraryRepository.GetParentFolderId(libraryPath, imageFolder);
Option<int> maybeParentFolder = await _libraryRepository.GetParentFolderId(
libraryPath,
imageFolder,
cancellationToken);
foldersCompleted++;
@@ -190,7 +193,7 @@ public class ImageFolderScanner : LocalFolderScanner, IImageFolderScanner
foreach (string file in allFiles.OrderBy(identity))
{
Either<BaseError, MediaItemScanResult<Image>> maybeVideo = await _imageRepository
.GetOrAdd(libraryPath, knownFolder, file)
.GetOrAdd(libraryPath, knownFolder, file, cancellationToken)
.BindT(video => UpdateStatistics(video, ffmpegPath, ffprobePath))
.BindT(video => UpdateLibraryFolderId(video, knownFolder))
.BindT(video => UpdateMetadata(video, durationSeconds))
@@ -27,7 +27,10 @@ public partial class LocalChaptersProvider : ILocalChaptersProvider
_logger = logger;
}
public async Task<bool> UpdateChapters(MediaItem mediaItem, Option<string> localPath)
public async Task<bool> UpdateChapters(
MediaItem mediaItem,
Option<string> localPath,
CancellationToken cancellationToken)
{
try
{
@@ -39,7 +42,7 @@ public partial class LocalChaptersProvider : ILocalChaptersProvider
if (chapters.Count > 0)
{
_logger.LogDebug("Located {Count} external chapters for {Path}", chapters.Count, mediaItemPath);
return await _metadataRepository.UpdateChapters(version, chapters);
return await _metadataRepository.UpdateChapters(version, chapters, cancellationToken);
}
return false;
@@ -31,11 +31,15 @@ public class LocalSubtitlesProvider : ILocalSubtitlesProvider
_logger = logger;
}
public async Task<bool> UpdateSubtitles(MediaItem mediaItem, Option<string> localPath, bool saveFullPath)
public async Task<bool> UpdateSubtitles(
MediaItem mediaItem,
Option<string> localPath,
bool saveFullPath,
CancellationToken cancellationToken)
{
if (_languageCodes.Count == 0)
{
await _slim.WaitAsync();
await _slim.WaitAsync(cancellationToken);
try
{
_languageCodes.AddRange(await _mediaItemRepository.GetAllKnownCultures());
@@ -79,7 +83,7 @@ public class LocalSubtitlesProvider : ILocalSubtitlesProvider
var subtitles = subtitleStreams.Map(Subtitle.FromMediaStream).ToList();
string mediaItemPath = await localPath.IfNoneAsync(() => mediaItem.GetHeadVersion().MediaFiles.Head().Path);
subtitles.AddRange(LocateExternalSubtitles(_languageCodes, mediaItemPath, saveFullPath));
bool updateResult = await _metadataRepository.UpdateSubtitles(metadata, subtitles);
bool updateResult = await _metadataRepository.UpdateSubtitles(metadata, subtitles, cancellationToken);
if (!updateResult)
{
_logger.LogError("Failed to save {Count} subtitles to database", subtitles.Count);
@@ -111,7 +111,7 @@ public abstract class MediaServerMovieLibraryScanner<TConnectionParameters, TLib
if (ServerReturnsStatisticsWithMetadata)
{
maybeMovie = await movieRepository
.GetOrAdd(library, incoming, deepScan)
.GetOrAdd(library, incoming, deepScan, cancellationToken)
.MapT(result =>
{
result.LocalPath = localPath;
@@ -122,13 +122,14 @@ public abstract class MediaServerMovieLibraryScanner<TConnectionParameters, TLib
library,
existing,
incoming,
deepScan))
.BindT(UpdateChapters);
deepScan,
cancellationToken))
.BindT(existing => UpdateChapters(existing, cancellationToken));
}
else
{
maybeMovie = await movieRepository
.GetOrAdd(library, incoming, deepScan)
.GetOrAdd(library, incoming, deepScan, cancellationToken)
.MapT(result =>
{
result.LocalPath = localPath;
@@ -140,7 +141,8 @@ public abstract class MediaServerMovieLibraryScanner<TConnectionParameters, TLib
existing,
incoming,
deepScan,
None))
None,
cancellationToken))
.BindT(existing => UpdateStatistics(
connectionParameters,
library,
@@ -148,8 +150,8 @@ public abstract class MediaServerMovieLibraryScanner<TConnectionParameters, TLib
incoming,
deepScan,
None))
.BindT(UpdateSubtitles)
.BindT(UpdateChapters);
.BindT(existing => UpdateSubtitles(existing, cancellationToken))
.BindT(existing => UpdateChapters(existing, cancellationToken));
}
if (maybeMovie.IsLeft)
@@ -255,7 +257,8 @@ public abstract class MediaServerMovieLibraryScanner<TConnectionParameters, TLib
protected abstract Task<Either<BaseError, MediaItemScanResult<TMovie>>> UpdateMetadata(
MediaItemScanResult<TMovie> result,
MovieMetadata fullMetadata);
MovieMetadata fullMetadata,
CancellationToken cancellationToken);
private async Task<bool> ShouldScanItem(
IMediaServerMovieRepository<TLibrary, TMovie, TEtag> movieRepository,
@@ -340,7 +343,8 @@ public abstract class MediaServerMovieLibraryScanner<TConnectionParameters, TLib
TLibrary library,
MediaItemScanResult<TMovie> result,
TMovie incoming,
bool deepScan)
bool deepScan,
CancellationToken cancellationToken)
{
Option<Tuple<MovieMetadata, MediaVersion>> maybeMetadataAndStatistics = await GetFullMetadataAndStatistics(
connectionParameters,
@@ -356,7 +360,8 @@ public abstract class MediaServerMovieLibraryScanner<TConnectionParameters, TLib
result,
incoming,
deepScan,
fullMetadata);
fullMetadata,
cancellationToken);
foreach (BaseError error in metadataResult.LeftToSeq())
{
@@ -396,7 +401,8 @@ public abstract class MediaServerMovieLibraryScanner<TConnectionParameters, TLib
MediaItemScanResult<TMovie> result,
TMovie incoming,
bool deepScan,
Option<MovieMetadata> maybeFullMetadata)
Option<MovieMetadata> maybeFullMetadata,
CancellationToken cancellationToken)
{
if (maybeFullMetadata.IsNone)
{
@@ -407,7 +413,7 @@ public abstract class MediaServerMovieLibraryScanner<TConnectionParameters, TLib
{
// TODO: move some of this code into this scanner
// will have to merge JF, Emby, Plex logic
return await UpdateMetadata(result, fullMetadata);
return await UpdateMetadata(result, fullMetadata, cancellationToken);
}
return result;
@@ -449,7 +455,8 @@ public abstract class MediaServerMovieLibraryScanner<TConnectionParameters, TLib
private async Task<Either<BaseError, MediaItemScanResult<TMovie>>> UpdateSubtitles(
MediaItemScanResult<TMovie> existing)
MediaItemScanResult<TMovie> existing,
CancellationToken cancellationToken)
{
try
{
@@ -462,7 +469,7 @@ public abstract class MediaServerMovieLibraryScanner<TConnectionParameters, TLib
.Map(Subtitle.FromMediaStream)
.ToList();
if (await _metadataRepository.UpdateSubtitles(metadata, subtitles))
if (await _metadataRepository.UpdateSubtitles(metadata, subtitles, cancellationToken))
{
return existing;
}
@@ -477,7 +484,8 @@ public abstract class MediaServerMovieLibraryScanner<TConnectionParameters, TLib
}
private async Task<Either<BaseError, MediaItemScanResult<TMovie>>> UpdateChapters(
MediaItemScanResult<TMovie> existing)
MediaItemScanResult<TMovie> existing,
CancellationToken cancellationToken)
{
try
{
@@ -487,7 +495,7 @@ public abstract class MediaServerMovieLibraryScanner<TConnectionParameters, TLib
return existing;
}
if (await _localChaptersProvider.UpdateChapters(existing.Item, Some(existing.LocalPath)))
if (await _localChaptersProvider.UpdateChapters(existing.Item, Some(existing.LocalPath), cancellationToken))
{
existing.IsUpdated = true;
}
@@ -96,8 +96,8 @@ public abstract class MediaServerOtherVideoLibraryScanner<TConnectionParameters,
library.Id,
library.Name,
percentCompletion,
Array.Empty<int>(),
Array.Empty<int>()),
[],
[]),
cancellationToken);
string localPath = getLocalPath(incoming);
@@ -118,7 +118,7 @@ public abstract class MediaServerOtherVideoLibraryScanner<TConnectionParameters,
if (ServerReturnsStatisticsWithMetadata)
{
maybeOtherVideo = await otherVideoRepository
.GetOrAdd(library, incoming, deepScan)
.GetOrAdd(library, incoming, deepScan, cancellationToken)
.MapT(result =>
{
result.LocalPath = localPath;
@@ -129,13 +129,14 @@ public abstract class MediaServerOtherVideoLibraryScanner<TConnectionParameters,
library,
existing,
incoming,
deepScan))
.BindT(UpdateChapters);
deepScan,
cancellationToken))
.BindT(existing => UpdateChapters(existing, cancellationToken));
}
else
{
maybeOtherVideo = await otherVideoRepository
.GetOrAdd(library, incoming, deepScan)
.GetOrAdd(library, incoming, deepScan, cancellationToken)
.MapT(result =>
{
result.LocalPath = localPath;
@@ -147,7 +148,8 @@ public abstract class MediaServerOtherVideoLibraryScanner<TConnectionParameters,
existing,
incoming,
deepScan,
None))
None,
cancellationToken))
.BindT(existing => UpdateStatistics(
connectionParameters,
library,
@@ -155,8 +157,8 @@ public abstract class MediaServerOtherVideoLibraryScanner<TConnectionParameters,
incoming,
deepScan,
None))
.BindT(UpdateSubtitles)
.BindT(UpdateChapters);
.BindT(existing => UpdateSubtitles(existing, cancellationToken))
.BindT(existing => UpdateChapters(existing, cancellationToken));
}
if (maybeOtherVideo.IsLeft)
@@ -262,7 +264,8 @@ public abstract class MediaServerOtherVideoLibraryScanner<TConnectionParameters,
protected abstract Task<Either<BaseError, MediaItemScanResult<TOtherVideo>>> UpdateMetadata(
MediaItemScanResult<TOtherVideo> result,
OtherVideoMetadata fullMetadata);
OtherVideoMetadata fullMetadata,
CancellationToken cancellationToken);
private async Task<bool> ShouldScanItem(
IMediaServerOtherVideoRepository<TLibrary, TOtherVideo, TEtag> otherVideoRepository,
@@ -349,7 +352,8 @@ public abstract class MediaServerOtherVideoLibraryScanner<TConnectionParameters,
TLibrary library,
MediaItemScanResult<TOtherVideo> result,
TOtherVideo incoming,
bool deepScan)
bool deepScan,
CancellationToken cancellationToken)
{
Option<Tuple<OtherVideoMetadata, MediaVersion>> maybeMetadataAndStatistics = await GetFullMetadataAndStatistics(
connectionParameters,
@@ -365,7 +369,8 @@ public abstract class MediaServerOtherVideoLibraryScanner<TConnectionParameters,
result,
incoming,
deepScan,
fullMetadata);
fullMetadata,
cancellationToken);
foreach (BaseError error in metadataResult.LeftToSeq())
{
@@ -405,7 +410,8 @@ public abstract class MediaServerOtherVideoLibraryScanner<TConnectionParameters,
MediaItemScanResult<TOtherVideo> result,
TOtherVideo incoming,
bool deepScan,
Option<OtherVideoMetadata> maybeFullMetadata)
Option<OtherVideoMetadata> maybeFullMetadata,
CancellationToken cancellationToken)
{
if (maybeFullMetadata.IsNone)
{
@@ -416,7 +422,7 @@ public abstract class MediaServerOtherVideoLibraryScanner<TConnectionParameters,
{
// TODO: move some of this code into this scanner
// will have to merge JF, Emby, Plex logic
return await UpdateMetadata(result, fullMetadata);
return await UpdateMetadata(result, fullMetadata, cancellationToken);
}
return result;
@@ -458,7 +464,8 @@ public abstract class MediaServerOtherVideoLibraryScanner<TConnectionParameters,
private async Task<Either<BaseError, MediaItemScanResult<TOtherVideo>>> UpdateSubtitles(
MediaItemScanResult<TOtherVideo> existing)
MediaItemScanResult<TOtherVideo> existing,
CancellationToken cancellationToken)
{
try
{
@@ -471,7 +478,7 @@ public abstract class MediaServerOtherVideoLibraryScanner<TConnectionParameters,
.Map(Subtitle.FromMediaStream)
.ToList();
if (await _metadataRepository.UpdateSubtitles(metadata, subtitles))
if (await _metadataRepository.UpdateSubtitles(metadata, subtitles, cancellationToken))
{
return existing;
}
@@ -486,7 +493,8 @@ public abstract class MediaServerOtherVideoLibraryScanner<TConnectionParameters,
}
private async Task<Either<BaseError, MediaItemScanResult<TOtherVideo>>> UpdateChapters(
MediaItemScanResult<TOtherVideo> existing)
MediaItemScanResult<TOtherVideo> existing,
CancellationToken cancellationToken)
{
try
{
@@ -496,7 +504,7 @@ public abstract class MediaServerOtherVideoLibraryScanner<TConnectionParameters,
return existing;
}
if (await _localChaptersProvider.UpdateChapters(existing.Item, Some(existing.LocalPath)))
if (await _localChaptersProvider.UpdateChapters(existing.Item, Some(existing.LocalPath), cancellationToken))
{
existing.IsUpdated = true;
}
@@ -113,7 +113,7 @@ public abstract class MediaServerTelevisionLibraryScanner<TConnectionParameters,
cancellationToken);
Either<BaseError, MediaItemScanResult<TShow>> maybeShow = await televisionRepository
.GetOrAdd(library, incoming)
.GetOrAdd(library, incoming, cancellationToken)
.BindT(existing => UpdateMetadata(connectionParameters, library, existing, incoming, deepScan));
if (maybeShow.IsLeft)
@@ -281,7 +281,8 @@ public abstract class MediaServerTelevisionLibraryScanner<TConnectionParameters,
protected abstract Task<Either<BaseError, MediaItemScanResult<TEpisode>>> UpdateMetadata(
MediaItemScanResult<TEpisode> result,
EpisodeMetadata fullMetadata);
EpisodeMetadata fullMetadata,
CancellationToken cancellationToken);
private async Task<Either<BaseError, Unit>> ScanSeasons(
IMediaServerTelevisionRepository<TLibrary, TShow, TSeason, TEpisode, TEtag> televisionRepository,
@@ -309,7 +310,7 @@ public abstract class MediaServerTelevisionLibraryScanner<TConnectionParameters,
incomingItemIds.Add(MediaServerItemId(incoming));
Either<BaseError, MediaItemScanResult<TSeason>> maybeSeason = await televisionRepository
.GetOrAdd(library, incoming)
.GetOrAdd(library, incoming, cancellationToken)
.BindT(existing => UpdateMetadata(connectionParameters, library, existing, incoming, deepScan));
if (maybeSeason.IsLeft)
@@ -424,7 +425,7 @@ public abstract class MediaServerTelevisionLibraryScanner<TConnectionParameters,
if (ServerReturnsStatisticsWithMetadata)
{
maybeEpisode = await televisionRepository
.GetOrAdd(library, incoming, deepScan)
.GetOrAdd(library, incoming, deepScan, cancellationToken)
.MapT(result =>
{
result.LocalPath = localPath;
@@ -435,13 +436,14 @@ public abstract class MediaServerTelevisionLibraryScanner<TConnectionParameters,
library,
existing,
incoming,
deepScan))
.BindT(UpdateChapters);
deepScan,
cancellationToken))
.BindT(existing => UpdateChapters(existing, cancellationToken));
}
else
{
maybeEpisode = await televisionRepository
.GetOrAdd(library, incoming, deepScan)
.GetOrAdd(library, incoming, deepScan, cancellationToken)
.MapT(result =>
{
result.LocalPath = localPath;
@@ -453,7 +455,8 @@ public abstract class MediaServerTelevisionLibraryScanner<TConnectionParameters,
existing,
incoming,
deepScan,
None))
None,
cancellationToken))
.BindT(existing => UpdateStatistics(
connectionParameters,
library,
@@ -461,8 +464,8 @@ public abstract class MediaServerTelevisionLibraryScanner<TConnectionParameters,
incoming,
deepScan,
None))
.BindT(UpdateSubtitles)
.BindT(UpdateChapters);
.BindT(existing => UpdateSubtitles(existing, cancellationToken))
.BindT(existing => UpdateChapters(existing, cancellationToken));
}
if (maybeEpisode.IsLeft)
@@ -666,7 +669,8 @@ public abstract class MediaServerTelevisionLibraryScanner<TConnectionParameters,
TLibrary library,
MediaItemScanResult<TEpisode> result,
TEpisode incoming,
bool deepScan)
bool deepScan,
CancellationToken cancellationToken)
{
Option<Tuple<EpisodeMetadata, MediaVersion>> maybeMetadataAndStatistics = await GetFullMetadataAndStatistics(
connectionParameters,
@@ -682,7 +686,8 @@ public abstract class MediaServerTelevisionLibraryScanner<TConnectionParameters,
result,
incoming,
deepScan,
fullMetadata);
fullMetadata,
cancellationToken);
foreach (BaseError error in metadataResult.LeftToSeq())
{
@@ -722,7 +727,8 @@ public abstract class MediaServerTelevisionLibraryScanner<TConnectionParameters,
MediaItemScanResult<TEpisode> result,
TEpisode incoming,
bool deepScan,
Option<EpisodeMetadata> maybeFullMetadata)
Option<EpisodeMetadata> maybeFullMetadata,
CancellationToken cancellationToken)
{
if (maybeFullMetadata.IsNone)
{
@@ -733,7 +739,7 @@ public abstract class MediaServerTelevisionLibraryScanner<TConnectionParameters,
{
// TODO: move some of this code into this scanner
// will have to merge JF, Emby, Plex logic
return await UpdateMetadata(result, fullMetadata);
return await UpdateMetadata(result, fullMetadata, cancellationToken);
}
return result;
@@ -801,7 +807,8 @@ public abstract class MediaServerTelevisionLibraryScanner<TConnectionParameters,
}
private async Task<Either<BaseError, MediaItemScanResult<TEpisode>>> UpdateSubtitles(
MediaItemScanResult<TEpisode> existing)
MediaItemScanResult<TEpisode> existing,
CancellationToken cancellationToken)
{
try
{
@@ -814,7 +821,7 @@ public abstract class MediaServerTelevisionLibraryScanner<TConnectionParameters,
.Map(Subtitle.FromMediaStream)
.ToList();
if (await _metadataRepository.UpdateSubtitles(metadata, subtitles))
if (await _metadataRepository.UpdateSubtitles(metadata, subtitles, cancellationToken))
{
return existing;
}
@@ -829,7 +836,8 @@ public abstract class MediaServerTelevisionLibraryScanner<TConnectionParameters,
}
private async Task<Either<BaseError, MediaItemScanResult<TEpisode>>> UpdateChapters(
MediaItemScanResult<TEpisode> existing)
MediaItemScanResult<TEpisode> existing,
CancellationToken cancellationToken)
{
try
{
@@ -839,7 +847,7 @@ public abstract class MediaServerTelevisionLibraryScanner<TConnectionParameters,
return existing;
}
if (await _localChaptersProvider.UpdateChapters(existing.Item, Some(existing.LocalPath)))
if (await _localChaptersProvider.UpdateChapters(existing.Item, Some(existing.LocalPath), cancellationToken))
{
existing.IsUpdated = true;
}
@@ -120,7 +120,10 @@ public class MovieFolderScanner : LocalFolderScanner, IMovieFolderScanner
cancellationToken);
string movieFolder = folderQueue.Dequeue();
Option<int> maybeParentFolder = await _libraryRepository.GetParentFolderId(libraryPath, movieFolder);
Option<int> maybeParentFolder = await _libraryRepository.GetParentFolderId(
libraryPath,
movieFolder,
cancellationToken);
foldersCompleted++;
var filesForEtag = _localFileSystem.ListFiles(movieFolder).ToList();
@@ -176,14 +179,14 @@ public class MovieFolderScanner : LocalFolderScanner, IMovieFolderScanner
{
// TODO: figure out how to rebuild playlists
Either<BaseError, MediaItemScanResult<Movie>> maybeMovie = await _movieRepository
.GetOrAdd(libraryPath, knownFolder, file)
.GetOrAdd(libraryPath, knownFolder, file, cancellationToken)
.BindT(movie => UpdateStatistics(movie, ffmpegPath, ffprobePath))
.BindT(video => UpdateLibraryFolderId(video, knownFolder))
.BindT(UpdateMetadata)
.BindT(movie => UpdateArtwork(movie, ArtworkKind.Poster, cancellationToken))
.BindT(movie => UpdateArtwork(movie, ArtworkKind.FanArt, cancellationToken))
.BindT(UpdateSubtitles)
.BindT(UpdateChapters)
.BindT(movie => UpdateSubtitles(movie, cancellationToken))
.BindT(movie => UpdateChapters(movie, cancellationToken))
.BindT(FlagNormal);
foreach (BaseError error in maybeMovie.LeftToSeq())
@@ -324,11 +327,13 @@ public class MovieFolderScanner : LocalFolderScanner, IMovieFolderScanner
}
}
private async Task<Either<BaseError, MediaItemScanResult<Movie>>> UpdateSubtitles(MediaItemScanResult<Movie> result)
private async Task<Either<BaseError, MediaItemScanResult<Movie>>> UpdateSubtitles(
MediaItemScanResult<Movie> result,
CancellationToken cancellationToken)
{
try
{
await _localSubtitlesProvider.UpdateSubtitles(result.Item, None, true);
await _localSubtitlesProvider.UpdateSubtitles(result.Item, None, true, cancellationToken);
return result;
}
catch (Exception ex)
@@ -338,11 +343,13 @@ public class MovieFolderScanner : LocalFolderScanner, IMovieFolderScanner
}
}
private async Task<Either<BaseError, MediaItemScanResult<Movie>>> UpdateChapters(MediaItemScanResult<Movie> result)
private async Task<Either<BaseError, MediaItemScanResult<Movie>>> UpdateChapters(
MediaItemScanResult<Movie> result,
CancellationToken cancellationToken)
{
try
{
await _localChaptersProvider.UpdateChapters(result.Item, None);
await _localChaptersProvider.UpdateChapters(result.Item, None, cancellationToken);
return result;
}
catch (Exception ex)
@@ -332,7 +332,10 @@ public class MusicVideoFolderScanner : LocalFolderScanner, IMusicVideoFolderScan
}
string musicVideoFolder = folderQueue.Dequeue();
Option<int> maybeParentFolder = await _libraryRepository.GetParentFolderId(libraryPath, musicVideoFolder);
Option<int> maybeParentFolder = await _libraryRepository.GetParentFolderId(
libraryPath,
musicVideoFolder,
cancellationToken);
// _logger.LogDebug("Scanning music video folder {Folder}", musicVideoFolder);
@@ -382,8 +385,8 @@ public class MusicVideoFolderScanner : LocalFolderScanner, IMusicVideoFolderScan
.BindT(video => UpdateLibraryFolderId(video, knownFolder))
.BindT(UpdateMetadata)
.BindT(result => UpdateThumbnail(result, cancellationToken))
.BindT(UpdateSubtitles)
.BindT(UpdateChapters)
.BindT(result => UpdateSubtitles(result, cancellationToken))
.BindT(result => UpdateChapters(result, cancellationToken))
.BindT(FlagNormal);
foreach (BaseError error in maybeMusicVideo.LeftToSeq())
@@ -533,11 +536,12 @@ public class MusicVideoFolderScanner : LocalFolderScanner, IMusicVideoFolderScan
}
private async Task<Either<BaseError, MediaItemScanResult<MusicVideo>>> UpdateSubtitles(
MediaItemScanResult<MusicVideo> result)
MediaItemScanResult<MusicVideo> result,
CancellationToken cancellationToken)
{
try
{
await _localSubtitlesProvider.UpdateSubtitles(result.Item, None, true);
await _localSubtitlesProvider.UpdateSubtitles(result.Item, None, true, cancellationToken);
return result;
}
catch (Exception ex)
@@ -548,11 +552,12 @@ public class MusicVideoFolderScanner : LocalFolderScanner, IMusicVideoFolderScan
}
private async Task<Either<BaseError, MediaItemScanResult<MusicVideo>>> UpdateChapters(
MediaItemScanResult<MusicVideo> result)
MediaItemScanResult<MusicVideo> result,
CancellationToken cancellationToken)
{
try
{
await _localChaptersProvider.UpdateChapters(result.Item, None);
await _localChaptersProvider.UpdateChapters(result.Item, None, cancellationToken);
return result;
}
catch (Exception ex)
@@ -132,7 +132,7 @@ public class OtherVideoFolderScanner : LocalFolderScanner, IOtherVideoFolderScan
string otherVideoFolder = folderQueue.Dequeue();
Option<int> maybeParentFolder =
await _libraryRepository.GetParentFolderId(libraryPath, otherVideoFolder);
await _libraryRepository.GetParentFolderId(libraryPath, otherVideoFolder, cancellationToken);
foldersCompleted++;
@@ -186,13 +186,13 @@ public class OtherVideoFolderScanner : LocalFolderScanner, IOtherVideoFolderScan
_logger.LogDebug("Processing other video file {File}", file);
Either<BaseError, MediaItemScanResult<OtherVideo>> maybeVideo = await _otherVideoRepository
.GetOrAdd(libraryPath, knownFolder, file)
.GetOrAdd(libraryPath, knownFolder, file, cancellationToken)
.BindT(video => UpdateStatistics(video, ffmpegPath, ffprobePath))
.BindT(video => UpdateLibraryFolderId(video, knownFolder))
.BindT(UpdateMetadata)
.BindT(video => UpdateThumbnail(video, cancellationToken))
.BindT(UpdateSubtitles)
.BindT(UpdateChapters)
.BindT(result => UpdateSubtitles(result, cancellationToken))
.BindT(result => UpdateChapters(result, cancellationToken))
.BindT(FlagNormal);
foreach (BaseError error in maybeVideo.LeftToSeq())
@@ -329,11 +329,12 @@ public class OtherVideoFolderScanner : LocalFolderScanner, IOtherVideoFolderScan
}
private async Task<Either<BaseError, MediaItemScanResult<OtherVideo>>> UpdateSubtitles(
MediaItemScanResult<OtherVideo> result)
MediaItemScanResult<OtherVideo> result,
CancellationToken cancellationToken)
{
try
{
await _localSubtitlesProvider.UpdateSubtitles(result.Item, None, true);
await _localSubtitlesProvider.UpdateSubtitles(result.Item, None, true, cancellationToken);
return result;
}
catch (Exception ex)
@@ -344,11 +345,12 @@ public class OtherVideoFolderScanner : LocalFolderScanner, IOtherVideoFolderScan
}
private async Task<Either<BaseError, MediaItemScanResult<OtherVideo>>> UpdateChapters(
MediaItemScanResult<OtherVideo> result)
MediaItemScanResult<OtherVideo> result,
CancellationToken cancellationToken)
{
try
{
await _localChaptersProvider.UpdateChapters(result.Item, None);
await _localChaptersProvider.UpdateChapters(result.Item, None, cancellationToken);
return result;
}
catch (Exception ex)
@@ -127,7 +127,7 @@ public class RemoteStreamFolderScanner : LocalFolderScanner, IRemoteStreamFolder
string remoteStreamFolder = folderQueue.Dequeue();
Option<int> maybeParentFolder =
await _libraryRepository.GetParentFolderId(libraryPath, remoteStreamFolder);
await _libraryRepository.GetParentFolderId(libraryPath, remoteStreamFolder, cancellationToken);
foldersCompleted++;
@@ -178,7 +178,7 @@ public class RemoteStreamFolderScanner : LocalFolderScanner, IRemoteStreamFolder
foreach (string file in allFiles.OrderBy(identity))
{
Either<BaseError, MediaItemScanResult<RemoteStream>> maybeVideo = await _remoteStreamRepository
.GetOrAdd(libraryPath, knownFolder, file)
.GetOrAdd(libraryPath, knownFolder, file, cancellationToken)
.BindT(video => ParseRemoteStreamDefinition(video, deserializer, cancellationToken))
.BindT(video => UpdateStatistics(video, ffmpegPath, ffprobePath))
.BindT(video => UpdateLibraryFolderId(video, knownFolder))
@@ -117,7 +117,8 @@ public class SongFolderScanner : LocalFolderScanner, ISongFolderScanner
cancellationToken);
string songFolder = folderQueue.Dequeue();
Option<int> maybeParentFolder = await _libraryRepository.GetParentFolderId(libraryPath, songFolder);
Option<int> maybeParentFolder =
await _libraryRepository.GetParentFolderId(libraryPath, songFolder, cancellationToken);
foldersCompleted++;
@@ -116,7 +116,8 @@ public class TelevisionFolderScanner : LocalFolderScanner, ITelevisionFolderScan
Array.Empty<int>()),
cancellationToken);
Option<int> maybeParentFolder = await _libraryRepository.GetParentFolderId(libraryPath, showFolder);
Option<int> maybeParentFolder =
await _libraryRepository.GetParentFolderId(libraryPath, showFolder, cancellationToken);
// this folder is unused by the show, but will be used as parents of season folders
LibraryFolder _ = await _libraryRepository.GetOrAddFolder(
@@ -250,7 +251,8 @@ public class TelevisionFolderScanner : LocalFolderScanner, ITelevisionFolderScan
return new ScanCanceled();
}
Option<int> maybeParentFolder = await _libraryRepository.GetParentFolderId(libraryPath, seasonFolder);
Option<int> maybeParentFolder =
await _libraryRepository.GetParentFolderId(libraryPath, seasonFolder, cancellationToken);
string etag = FolderEtag.CalculateWithSubfolders(seasonFolder, _localFileSystem);
LibraryFolder knownFolder = await _libraryRepository.GetOrAddFolder(
@@ -345,14 +347,14 @@ public class TelevisionFolderScanner : LocalFolderScanner, ITelevisionFolderScan
{
// TODO: figure out how to rebuild playlists
Either<BaseError, Episode> maybeEpisode = await _televisionRepository
.GetOrAddEpisode(season, libraryPath, seasonFolder, file)
.GetOrAddEpisode(season, libraryPath, seasonFolder, file, cancellationToken)
.BindT(episode => UpdateStatistics(new MediaItemScanResult<Episode>(episode), ffmpegPath, ffprobePath)
.MapT(_ => episode))
.BindT(video => UpdateLibraryFolderId(video, seasonFolder))
.BindT(UpdateMetadata)
.BindT(e => UpdateThumbnail(e, cancellationToken))
.BindT(UpdateSubtitles)
.BindT(UpdateChapters)
.BindT(e => UpdateSubtitles(e, cancellationToken))
.BindT(e => UpdateChapters(e, cancellationToken))
.BindT(e => FlagNormal(new MediaItemScanResult<Episode>(e)))
.MapT(r => r.Item);
@@ -577,11 +579,11 @@ public class TelevisionFolderScanner : LocalFolderScanner, ITelevisionFolderScan
}
}
private async Task<Either<BaseError, Episode>> UpdateSubtitles(Episode episode)
private async Task<Either<BaseError, Episode>> UpdateSubtitles(Episode episode, CancellationToken cancellationToken)
{
try
{
await _localSubtitlesProvider.UpdateSubtitles(episode, None, true);
await _localSubtitlesProvider.UpdateSubtitles(episode, None, true, cancellationToken);
return episode;
}
catch (Exception ex)
@@ -591,11 +593,11 @@ public class TelevisionFolderScanner : LocalFolderScanner, ITelevisionFolderScan
}
}
private async Task<Either<BaseError, Episode>> UpdateChapters(Episode episode)
private async Task<Either<BaseError, Episode>> UpdateChapters(Episode episode, CancellationToken cancellationToken)
{
try
{
await _localChaptersProvider.UpdateChapters(episode, None);
await _localChaptersProvider.UpdateChapters(episode, None, cancellationToken);
return episode;
}
catch (Exception ex)
@@ -159,7 +159,8 @@ public class PlexMovieLibraryScanner :
protected override async Task<Either<BaseError, MediaItemScanResult<PlexMovie>>> UpdateMetadata(
MediaItemScanResult<PlexMovie> result,
MovieMetadata fullMetadata)
MovieMetadata fullMetadata,
CancellationToken cancellationToken)
{
PlexMovie existing = result.Item;
MovieMetadata existingMetadata = existing.MovieMetadata.Head();
@@ -340,7 +341,7 @@ public class PlexMovieLibraryScanner :
}
}
if (await _metadataRepository.UpdateSubtitles(existingMetadata, fullMetadata.Subtitles))
if (await _metadataRepository.UpdateSubtitles(existingMetadata, fullMetadata.Subtitles, cancellationToken))
{
result.IsUpdated = true;
}
@@ -161,7 +161,8 @@ public class PlexOtherVideoLibraryScanner :
protected override async Task<Either<BaseError, MediaItemScanResult<PlexOtherVideo>>> UpdateMetadata(
MediaItemScanResult<PlexOtherVideo> result,
OtherVideoMetadata fullMetadata)
OtherVideoMetadata fullMetadata,
CancellationToken cancellationToken)
{
PlexOtherVideo existing = result.Item;
OtherVideoMetadata existingMetadata = existing.OtherVideoMetadata.Head();
@@ -342,7 +343,7 @@ public class PlexOtherVideoLibraryScanner :
}
}
if (await _metadataRepository.UpdateSubtitles(existingMetadata, fullMetadata.Subtitles))
if (await _metadataRepository.UpdateSubtitles(existingMetadata, fullMetadata.Subtitles, cancellationToken))
{
result.IsUpdated = true;
}
@@ -603,7 +603,8 @@ public partial class PlexTelevisionLibraryScanner :
protected override async Task<Either<BaseError, MediaItemScanResult<PlexEpisode>>> UpdateMetadata(
MediaItemScanResult<PlexEpisode> result,
EpisodeMetadata fullMetadata)
EpisodeMetadata fullMetadata,
CancellationToken cancellationToken)
{
PlexEpisode existing = result.Item;
EpisodeMetadata existingMetadata = existing.EpisodeMetadata.Head();
@@ -686,7 +687,7 @@ public partial class PlexTelevisionLibraryScanner :
result.IsUpdated = true;
}
if (await _metadataRepository.UpdateSubtitles(existingMetadata, fullMetadata.Subtitles))
if (await _metadataRepository.UpdateSubtitles(existingMetadata, fullMetadata.Subtitles, cancellationToken))
{
result.IsUpdated = true;
}