fix search index threading (#141)

* fix search index threading

* code cleanup
This commit is contained in:
Jason Dove
2021-04-05 05:41:29 -05:00
committed by GitHub
parent 24cdf6295f
commit 3fb6da0754
29 changed files with 280 additions and 122 deletions
@@ -1,4 +1,5 @@
using System;
using System.Collections.Generic;
using System.Threading.Tasks;
using ErsatzTV.Core.Domain;
using LanguageExt;
@@ -12,6 +13,8 @@ namespace ErsatzTV.Core.Interfaces.Metadata
string ffprobePath,
DateTimeOffset lastScan,
decimal progressMin,
decimal progressMax);
decimal progressMax,
Func<List<MediaItem>, ValueTask> addToSearchIndex,
Func<List<int>, ValueTask> removeFromSearchIndex);
}
}
@@ -1,4 +1,5 @@
using System;
using System.Collections.Generic;
using System.Threading.Tasks;
using ErsatzTV.Core.Domain;
using LanguageExt;
@@ -12,6 +13,8 @@ namespace ErsatzTV.Core.Interfaces.Metadata
string ffprobePath,
DateTimeOffset lastScan,
decimal progressMin,
decimal progressMax);
decimal progressMax,
Func<List<MediaItem>, ValueTask> addToSearchIndex,
Func<List<int>, ValueTask> removeFromSearchIndex);
}
}
@@ -1,4 +1,5 @@
using System;
using System.Collections.Generic;
using System.Threading.Tasks;
using ErsatzTV.Core.Domain;
using LanguageExt;
@@ -12,6 +13,8 @@ namespace ErsatzTV.Core.Interfaces.Metadata
string ffprobePath,
DateTimeOffset lastScan,
decimal progressMin,
decimal progressMax);
decimal progressMax,
Func<List<MediaItem>, ValueTask> addToSearchIndex,
Func<List<int>, ValueTask> removeFromSearchIndex);
}
}
@@ -1,4 +1,6 @@
using System.Threading.Tasks;
using System;
using System.Collections.Generic;
using System.Threading.Tasks;
using ErsatzTV.Core.Domain;
using ErsatzTV.Core.Plex;
using LanguageExt;
@@ -10,6 +12,8 @@ namespace ErsatzTV.Core.Interfaces.Plex
Task<Either<BaseError, Unit>> ScanLibrary(
PlexConnection connection,
PlexServerAuthToken token,
PlexLibrary plexMediaSourceLibrary);
PlexLibrary plexMediaSourceLibrary,
Func<List<MediaItem>, ValueTask> addToSearchIndex,
Func<List<int>, ValueTask> removeFromSearchIndex);
}
}
@@ -1,4 +1,6 @@
using System.Threading.Tasks;
using System;
using System.Collections.Generic;
using System.Threading.Tasks;
using ErsatzTV.Core.Domain;
using ErsatzTV.Core.Plex;
using LanguageExt;
@@ -10,6 +12,8 @@ namespace ErsatzTV.Core.Interfaces.Plex
Task<Either<BaseError, Unit>> ScanLibrary(
PlexConnection connection,
PlexServerAuthToken token,
PlexLibrary plexMediaSourceLibrary);
PlexLibrary plexMediaSourceLibrary,
Func<List<MediaItem>, ValueTask> addToSearchIndex,
Func<List<int>, ValueTask> removeFromSearchIndex);
}
}
+5 -17
View File
@@ -8,7 +8,6 @@ using ErsatzTV.Core.Errors;
using ErsatzTV.Core.Interfaces.Images;
using ErsatzTV.Core.Interfaces.Metadata;
using ErsatzTV.Core.Interfaces.Repositories;
using ErsatzTV.Core.Interfaces.Search;
using LanguageExt;
using MediatR;
using Microsoft.Extensions.Logging;
@@ -25,7 +24,6 @@ namespace ErsatzTV.Core.Metadata
private readonly ILogger<MovieFolderScanner> _logger;
private readonly IMediator _mediator;
private readonly IMovieRepository _movieRepository;
private readonly ISearchIndex _searchIndex;
public MovieFolderScanner(
ILocalFileSystem localFileSystem,
@@ -34,7 +32,6 @@ namespace ErsatzTV.Core.Metadata
ILocalMetadataProvider localMetadataProvider,
IMetadataRepository metadataRepository,
IImageCache imageCache,
ISearchIndex searchIndex,
IMediator mediator,
ILogger<MovieFolderScanner> logger)
: base(localFileSystem, localStatisticsProvider, metadataRepository, imageCache, logger)
@@ -42,7 +39,6 @@ namespace ErsatzTV.Core.Metadata
_localFileSystem = localFileSystem;
_movieRepository = movieRepository;
_localMetadataProvider = localMetadataProvider;
_searchIndex = searchIndex;
_mediator = mediator;
_logger = logger;
}
@@ -52,7 +48,9 @@ namespace ErsatzTV.Core.Metadata
string ffprobePath,
DateTimeOffset lastScan,
decimal progressMin,
decimal progressMax)
decimal progressMax,
Func<List<MediaItem>, ValueTask> addToSearchIndex,
Func<List<int>, ValueTask> removeFromSearchIndex)
{
decimal progressSpread = progressMax - progressMin;
@@ -111,17 +109,7 @@ namespace ErsatzTV.Core.Metadata
.BindT(movie => UpdateArtwork(movie, ArtworkKind.FanArt));
await maybeMovie.Match(
async result =>
{
if (result.IsAdded)
{
await _searchIndex.AddItems(new List<MediaItem> { result.Item });
}
else if (result.IsUpdated)
{
await _searchIndex.UpdateItems(new List<MediaItem> { result.Item });
}
},
async result => await addToSearchIndex(new List<MediaItem> { result.Item }),
error =>
{
_logger.LogWarning("Error processing movie at {Path}: {Error}", file, error.Value);
@@ -136,7 +124,7 @@ namespace ErsatzTV.Core.Metadata
{
_logger.LogInformation("Removing missing movie at {Path}", path);
List<int> ids = await _movieRepository.DeleteByPath(libraryPath, path);
await _searchIndex.RemoveItems(ids);
await removeFromSearchIndex(ids);
}
}
@@ -8,7 +8,6 @@ using ErsatzTV.Core.Errors;
using ErsatzTV.Core.Interfaces.Images;
using ErsatzTV.Core.Interfaces.Metadata;
using ErsatzTV.Core.Interfaces.Repositories;
using ErsatzTV.Core.Interfaces.Search;
using LanguageExt;
using MediatR;
using Microsoft.Extensions.Logging;
@@ -24,7 +23,6 @@ namespace ErsatzTV.Core.Metadata
private readonly ILogger<MusicVideoFolderScanner> _logger;
private readonly IMediator _mediator;
private readonly IMusicVideoRepository _musicVideoRepository;
private readonly ISearchIndex _searchIndex;
public MusicVideoFolderScanner(
ILocalFileSystem localFileSystem,
@@ -32,7 +30,6 @@ namespace ErsatzTV.Core.Metadata
ILocalMetadataProvider localMetadataProvider,
IMetadataRepository metadataRepository,
IImageCache imageCache,
ISearchIndex searchIndex,
IMusicVideoRepository musicVideoRepository,
IMediator mediator,
ILogger<MusicVideoFolderScanner> logger) : base(
@@ -44,7 +41,6 @@ namespace ErsatzTV.Core.Metadata
{
_localFileSystem = localFileSystem;
_localMetadataProvider = localMetadataProvider;
_searchIndex = searchIndex;
_musicVideoRepository = musicVideoRepository;
_mediator = mediator;
_logger = logger;
@@ -55,7 +51,9 @@ namespace ErsatzTV.Core.Metadata
string ffprobePath,
DateTimeOffset lastScan,
decimal progressMin,
decimal progressMax)
decimal progressMax,
Func<List<MediaItem>, ValueTask> addToSearchIndex,
Func<List<int>, ValueTask> removeFromSearchIndex)
{
decimal progressSpread = progressMax - progressMin;
@@ -104,17 +102,7 @@ namespace ErsatzTV.Core.Metadata
.BindT(UpdateThumbnail);
await maybeMusicVideo.Match(
async result =>
{
if (result.IsAdded)
{
await _searchIndex.AddItems(new List<MediaItem> { result.Item });
}
else if (result.IsUpdated)
{
await _searchIndex.UpdateItems(new List<MediaItem> { result.Item });
}
},
async result => await addToSearchIndex(new List<MediaItem> { result.Item }),
error =>
{
_logger.LogWarning("Error processing music video at {Path}: {Error}", file, error.Value);
@@ -129,7 +117,7 @@ namespace ErsatzTV.Core.Metadata
{
_logger.LogInformation("Removing missing music video at {Path}", path);
List<int> ids = await _musicVideoRepository.DeleteByPath(libraryPath, path);
await _searchIndex.RemoveItems(ids);
await removeFromSearchIndex(ids);
}
}
@@ -8,7 +8,6 @@ using ErsatzTV.Core.Errors;
using ErsatzTV.Core.Interfaces.Images;
using ErsatzTV.Core.Interfaces.Metadata;
using ErsatzTV.Core.Interfaces.Repositories;
using ErsatzTV.Core.Interfaces.Search;
using LanguageExt;
using MediatR;
using Microsoft.Extensions.Logging;
@@ -23,7 +22,6 @@ namespace ErsatzTV.Core.Metadata
private readonly ILocalMetadataProvider _localMetadataProvider;
private readonly ILogger<TelevisionFolderScanner> _logger;
private readonly IMediator _mediator;
private readonly ISearchIndex _searchIndex;
private readonly ITelevisionRepository _televisionRepository;
public TelevisionFolderScanner(
@@ -33,7 +31,6 @@ namespace ErsatzTV.Core.Metadata
ILocalMetadataProvider localMetadataProvider,
IMetadataRepository metadataRepository,
IImageCache imageCache,
ISearchIndex searchIndex,
IMediator mediator,
ILogger<TelevisionFolderScanner> logger) : base(
localFileSystem,
@@ -45,7 +42,6 @@ namespace ErsatzTV.Core.Metadata
_localFileSystem = localFileSystem;
_televisionRepository = televisionRepository;
_localMetadataProvider = localMetadataProvider;
_searchIndex = searchIndex;
_mediator = mediator;
_logger = logger;
}
@@ -55,7 +51,9 @@ namespace ErsatzTV.Core.Metadata
string ffprobePath,
DateTimeOffset lastScan,
decimal progressMin,
decimal progressMax)
decimal progressMax,
Func<List<MediaItem>, ValueTask> addToSearchIndex,
Func<List<int>, ValueTask> removeFromSearchIndex)
{
decimal progressSpread = progressMax - progressMin;
@@ -84,15 +82,7 @@ namespace ErsatzTV.Core.Metadata
await maybeShow.Match(
async result =>
{
if (result.IsAdded)
{
await _searchIndex.AddItems(new List<MediaItem> { result.Item });
}
else if (result.IsUpdated)
{
await _searchIndex.UpdateItems(new List<MediaItem> { result.Item });
}
await addToSearchIndex(new List<MediaItem> { result.Item });
await ScanSeasons(
libraryPath,
ffprobePath,
@@ -122,7 +112,7 @@ namespace ErsatzTV.Core.Metadata
await _televisionRepository.DeleteEmptySeasons(libraryPath);
List<int> ids = await _televisionRepository.DeleteEmptyShows(libraryPath);
await _searchIndex.RemoveItems(ids);
await removeFromSearchIndex(ids);
return Unit.Default;
}
+7 -18
View File
@@ -1,10 +1,10 @@
using System.Collections.Generic;
using System;
using System.Collections.Generic;
using System.Linq;
using System.Threading.Tasks;
using ErsatzTV.Core.Domain;
using ErsatzTV.Core.Interfaces.Plex;
using ErsatzTV.Core.Interfaces.Repositories;
using ErsatzTV.Core.Interfaces.Search;
using ErsatzTV.Core.Metadata;
using LanguageExt;
using MediatR;
@@ -20,13 +20,11 @@ namespace ErsatzTV.Core.Plex
private readonly IMetadataRepository _metadataRepository;
private readonly IMovieRepository _movieRepository;
private readonly IPlexServerApiClient _plexServerApiClient;
private readonly ISearchIndex _searchIndex;
public PlexMovieLibraryScanner(
IPlexServerApiClient plexServerApiClient,
IMovieRepository movieRepository,
IMetadataRepository metadataRepository,
ISearchIndex searchIndex,
IMediator mediator,
ILogger<PlexMovieLibraryScanner> logger)
: base(metadataRepository, logger)
@@ -34,7 +32,6 @@ namespace ErsatzTV.Core.Plex
_plexServerApiClient = plexServerApiClient;
_movieRepository = movieRepository;
_metadataRepository = metadataRepository;
_searchIndex = searchIndex;
_mediator = mediator;
_logger = logger;
}
@@ -42,7 +39,9 @@ namespace ErsatzTV.Core.Plex
public async Task<Either<BaseError, Unit>> ScanLibrary(
PlexConnection connection,
PlexServerAuthToken token,
PlexLibrary plexMediaSourceLibrary)
PlexLibrary plexMediaSourceLibrary,
Func<List<MediaItem>, ValueTask> addToSearchIndex,
Func<List<int>, ValueTask> removeFromSearchIndex)
{
Either<BaseError, List<PlexMovie>> entries = await _plexServerApiClient.GetMovieLibraryContents(
plexMediaSourceLibrary,
@@ -65,17 +64,7 @@ namespace ErsatzTV.Core.Plex
.BindT(existing => UpdateArtwork(existing, incoming));
await maybeMovie.Match(
async result =>
{
if (result.IsAdded)
{
await _searchIndex.AddItems(new List<MediaItem> { result.Item });
}
else if (result.IsUpdated)
{
await _searchIndex.UpdateItems(new List<MediaItem> { result.Item });
}
},
async result => await addToSearchIndex(new List<MediaItem> { result.Item }),
error =>
{
_logger.LogWarning(
@@ -88,7 +77,7 @@ namespace ErsatzTV.Core.Plex
var movieKeys = movieEntries.Map(s => s.Key).ToList();
List<int> ids = await _movieRepository.RemoveMissingPlexMovies(plexMediaSourceLibrary, movieKeys);
await _searchIndex.RemoveItems(ids);
await removeFromSearchIndex(ids);
await _mediator.Publish(new LibraryScanProgress(plexMediaSourceLibrary.Id, 0));
},
@@ -1,10 +1,10 @@
using System.Collections.Generic;
using System;
using System.Collections.Generic;
using System.Linq;
using System.Threading.Tasks;
using ErsatzTV.Core.Domain;
using ErsatzTV.Core.Interfaces.Plex;
using ErsatzTV.Core.Interfaces.Repositories;
using ErsatzTV.Core.Interfaces.Search;
using ErsatzTV.Core.Metadata;
using LanguageExt;
using MediatR;
@@ -20,14 +20,12 @@ namespace ErsatzTV.Core.Plex
private readonly IMediator _mediator;
private readonly IMetadataRepository _metadataRepository;
private readonly IPlexServerApiClient _plexServerApiClient;
private readonly ISearchIndex _searchIndex;
private readonly ITelevisionRepository _televisionRepository;
public PlexTelevisionLibraryScanner(
IPlexServerApiClient plexServerApiClient,
ITelevisionRepository televisionRepository,
IMetadataRepository metadataRepository,
ISearchIndex searchIndex,
IMediator mediator,
ILogger<PlexTelevisionLibraryScanner> logger)
: base(metadataRepository, logger)
@@ -35,7 +33,6 @@ namespace ErsatzTV.Core.Plex
_plexServerApiClient = plexServerApiClient;
_televisionRepository = televisionRepository;
_metadataRepository = metadataRepository;
_searchIndex = searchIndex;
_mediator = mediator;
_logger = logger;
}
@@ -43,7 +40,9 @@ namespace ErsatzTV.Core.Plex
public async Task<Either<BaseError, Unit>> ScanLibrary(
PlexConnection connection,
PlexServerAuthToken token,
PlexLibrary plexMediaSourceLibrary)
PlexLibrary plexMediaSourceLibrary,
Func<List<MediaItem>, ValueTask> addToSearchIndex,
Func<List<int>, ValueTask> removeFromSearchIndex)
{
Either<BaseError, List<PlexShow>> entries = await _plexServerApiClient.GetShowLibraryContents(
plexMediaSourceLibrary,
@@ -67,15 +66,7 @@ namespace ErsatzTV.Core.Plex
await maybeShow.Match(
async result =>
{
if (result.IsAdded)
{
await _searchIndex.AddItems(new List<MediaItem> { result.Item });
}
else if (result.IsUpdated)
{
await _searchIndex.UpdateItems(new List<MediaItem> { result.Item });
}
await addToSearchIndex(new List<MediaItem> { result.Item });
await ScanSeasons(plexMediaSourceLibrary, result.Item, connection, token);
},
error =>
@@ -91,7 +82,7 @@ namespace ErsatzTV.Core.Plex
var showKeys = showEntries.Map(s => s.Key).ToList();
List<int> ids =
await _televisionRepository.RemoveMissingPlexShows(plexMediaSourceLibrary, showKeys);
await _searchIndex.RemoveItems(ids);
await removeFromSearchIndex(ids);
await _mediator.Publish(new LibraryScanProgress(plexMediaSourceLibrary.Id, 0));