Merge pull request #17763 from Shadowghost/fix-refresh-queue-and-directory-cache

Fix items being lost from the refresh queue and bound the directory service caches
This commit is contained in:
Cody Robibero
2026-09-06 07:36:13 -04:00
committed by GitHub
14 changed files with 879 additions and 85 deletions
@@ -17,7 +17,8 @@ namespace MediaBrowser.Controller.LibraryTaskScheduler;
/// </summary>
public sealed class LimitedConcurrencyLibraryScheduler : ILimitedConcurrencyLibraryScheduler, IAsyncDisposable
{
private const int CleanupGracePeriod = 60;
private static readonly TimeSpan _cleanupGracePeriod = TimeSpan.FromSeconds(60);
private readonly IHostApplicationLifetime _hostApplicationLifetime;
private readonly ILogger<LimitedConcurrencyLibraryScheduler> _logger;
private readonly IServerConfigurationManager _serverConfigurationManager;
@@ -31,6 +32,8 @@ public sealed class LimitedConcurrencyLibraryScheduler : ILimitedConcurrencyLibr
private readonly Lock _taskLock = new();
private readonly Channel<TaskQueueItem> _tasks = Channel.CreateUnbounded<TaskQueueItem>();
private readonly CancellationTokenSource _disposeTokenSource = new();
private readonly TimeSpan _gracePeriod;
private volatile int _workCounter;
private Task? _cleanupTask;
@@ -46,10 +49,34 @@ public sealed class LimitedConcurrencyLibraryScheduler : ILimitedConcurrencyLibr
IHostApplicationLifetime hostApplicationLifetime,
ILogger<LimitedConcurrencyLibraryScheduler> logger,
IServerConfigurationManager serverConfigurationManager)
: this(hostApplicationLifetime, logger, serverConfigurationManager, _cleanupGracePeriod)
{
}
internal LimitedConcurrencyLibraryScheduler(
IHostApplicationLifetime hostApplicationLifetime,
ILogger<LimitedConcurrencyLibraryScheduler> logger,
IServerConfigurationManager serverConfigurationManager,
TimeSpan gracePeriod)
{
_hostApplicationLifetime = hostApplicationLifetime;
_logger = logger;
_serverConfigurationManager = serverConfigurationManager;
_gracePeriod = gracePeriod;
}
/// <summary>
/// Gets the number of runners the scheduler currently keeps alive.
/// </summary>
internal int ActiveRunnerCount
{
get
{
lock (_taskLock)
{
return _taskRunners.Count;
}
}
}
private void ScheduleTaskCleanup()
@@ -68,31 +95,65 @@ public sealed class LimitedConcurrencyLibraryScheduler : ILimitedConcurrencyLibr
async Task RunCleanupTask()
{
_logger.LogDebug("Schedule cleanup task in {CleanupGracePerioid} sec.", CleanupGracePeriod);
await Task.Delay(TimeSpan.FromSeconds(CleanupGracePeriod)).ConfigureAwait(false);
if (_disposed)
while (true)
{
_logger.LogDebug("Abort cleaning up, already disposed.");
return;
}
lock (_taskLock)
{
if (_tasks.Reader.Count > 0 || _workCounter > 0)
_logger.LogDebug("Schedule cleanup task in {CleanupGracePeriod}.", _gracePeriod);
try
{
_logger.LogDebug("Delay cleanup task, operations still running.");
// tasks are still there so its still in use. Reschedule cleanup task.
// we cannot just exit here and rely on the other invoker because there is a considerable timeframe where it could have already ended.
_cleanupTask = RunCleanupTask();
await Task.Delay(_gracePeriod, _disposeTokenSource.Token).ConfigureAwait(false);
}
catch (OperationCanceledException)
{
_logger.LogDebug("Abort cleaning up, already disposed.");
return;
}
}
_logger.LogDebug("Cleanup runners.");
foreach (var item in _taskRunners.ToArray())
if (_disposed)
{
_logger.LogDebug("Abort cleaning up, already disposed.");
return;
}
CancellationTokenSource[] runners;
lock (_taskLock)
{
if (_tasks.Reader.Count > 0 || _workCounter > 0)
{
_logger.LogDebug("Delay cleanup task, operations still running.");
// tasks are still there so its still in use. Wait another grace period.
// we cannot just exit here and rely on the other invoker because there is a considerable timeframe where it could have already ended.
continue;
}
runners = [.. _taskRunners.Keys];
// Retire the runners before they are told to stop: an operation starting while
// they wind down must spawn its own instead of counting these towards the fanout.
_taskRunners.Clear();
// Hand the next operation the ability to schedule a cleanup again. Without this
// the very first cleanup would be the only one that ever runs.
_cleanupTask = null;
}
_logger.LogDebug("Cleanup runners.");
await StopRunners(runners).ConfigureAwait(false);
return;
}
}
}
private static async Task StopRunners(CancellationTokenSource[] runners)
{
foreach (var runner in runners)
{
try
{
await item.Key.CancelAsync().ConfigureAwait(false);
_taskRunners.Remove(item.Key);
await runner.CancelAsync().ConfigureAwait(false);
}
catch (ObjectDisposedException)
{
// The runner already stopped on its own and disposed its stop source.
}
}
}
@@ -127,12 +188,17 @@ public sealed class LimitedConcurrencyLibraryScheduler : ILimitedConcurrencyLibr
{
var stopToken = new CancellationTokenSource();
var combinedSource = CancellationTokenSource.CreateLinkedTokenSource(stopToken.Token, _hostApplicationLifetime.ApplicationStopping);
// Keyed on its own stop source, because cancelling that is what reaches the linked
// source the runner waits on. Cancellation does not travel the other way.
// Started without the runner's own token: a task cancelled before it is scheduled
// never runs its body, so it would never take itself out of _taskRunners again.
_taskRunners.Add(
combinedSource,
stopToken,
Task.Factory.StartNew(
ItemWorker,
(combinedSource, stopToken),
combinedSource.Token,
(stopToken, combinedSource),
CancellationToken.None,
TaskCreationOptions.PreferFairness,
TaskScheduler.Default));
}
@@ -145,7 +211,7 @@ public sealed class LimitedConcurrencyLibraryScheduler : ILimitedConcurrencyLibr
_deadlockDetector.Value = stopToken.TaskStop;
try
{
while (!stopToken.GlobalStop.Token.IsCancellationRequested)
while (!stopToken.GlobalStop.IsCancellationRequested)
{
var item = await _tasks.Reader.ReadAsync(stopToken.GlobalStop.Token).ConfigureAwait(false);
try
@@ -162,15 +228,24 @@ public sealed class LimitedConcurrencyLibraryScheduler : ILimitedConcurrencyLibr
}
}
}
catch (OperationCanceledException) when (stopToken.TaskStop.IsCancellationRequested)
catch (OperationCanceledException) when (stopToken.GlobalStop.IsCancellationRequested)
{
// thats how you do it, interupt the waiter thread. There is nothing to do here when it was on purpose.
}
catch (ChannelClosedException)
{
// the scheduler was disposed and will not hand out any more work.
}
finally
{
_logger.LogDebug("Cleanup Runner'.");
_deadlockDetector.Value = default!;
_taskRunners.Remove(stopToken.TaskStop);
lock (_taskLock)
{
_taskRunners.Remove(stopToken.TaskStop);
}
stopToken.GlobalStop.Dispose();
stopToken.TaskStop.Dispose();
}
@@ -195,7 +270,7 @@ public sealed class LimitedConcurrencyLibraryScheduler : ILimitedConcurrencyLibr
finally
{
item.Progress.Report(100);
item.Done.SetResult();
item.Done.TrySetResult();
}
}
@@ -285,16 +360,33 @@ public sealed class LimitedConcurrencyLibraryScheduler : ILimitedConcurrencyLibr
_disposed = true;
_tasks.Writer.Complete();
foreach (var item in _taskRunners)
// Nobody is left to run these, so release whoever is waiting on them.
while (_tasks.Reader.TryRead(out var item))
{
await item.Key.CancelAsync().ConfigureAwait(false);
item.Done.TrySetResult();
}
if (_cleanupTask is not null)
CancellationTokenSource[] runners;
Task? cleanupTask;
lock (_taskLock)
{
await _cleanupTask.ConfigureAwait(false);
_cleanupTask?.Dispose();
runners = [.. _taskRunners.Keys];
_taskRunners.Clear();
cleanupTask = _cleanupTask;
}
await StopRunners(runners).ConfigureAwait(false);
// Cuts the grace period short instead of holding up shutdown for the rest of it.
await _disposeTokenSource.CancelAsync().ConfigureAwait(false);
if (cleanupTask is not null)
{
await cleanupTask.ConfigureAwait(false);
}
_disposeTokenSource.Dispose();
}
private class TaskQueueItem
@@ -18,7 +18,6 @@
</PropertyGroup>
<ItemGroup>
<PackageReference Include="BitFaster.Caching" />
<PackageReference Include="Microsoft.Extensions.Configuration.Binder" />
</ItemGroup>
@@ -5,13 +5,19 @@ using System.Collections.Concurrent;
using System.Collections.Generic;
using System.IO;
using System.Linq;
using System.Threading;
using MediaBrowser.Model.IO;
namespace MediaBrowser.Controller.Providers
{
public class DirectoryService : IDirectoryService
{
// TODO make static and switch to FastConcurrentLru.
// TODO replace with one shared bounded cache.
private const int MaxCachedRecords = 100_000;
private const int AccessIntervalMs = 1_000;
// Timeout cache if no access for 5 minutes.
private const int IdleTimeoutMs = 5 * 60 * 1_000;
private readonly ConcurrentDictionary<string, FileSystemMetadata[]> _cache = new(StringComparer.Ordinal);
private readonly ConcurrentDictionary<string, FileSystemMetadata> _fileCache = new(StringComparer.Ordinal);
@@ -20,6 +26,12 @@ namespace MediaBrowser.Controller.Providers
private readonly IFileSystem _fileSystem;
// ConcurrentDictionary.Count locks the dictionary, so keep an estimated counter.
// Concurrent factory runs can overcount and a clear racing an add can undercount,
// it only has to be roughly right.
private int _recordCount;
private long _lastAccess = Environment.TickCount64;
public DirectoryService(IFileSystem fileSystem)
{
_fileSystem = fileSystem;
@@ -27,20 +39,26 @@ namespace MediaBrowser.Controller.Providers
public FileSystemMetadata[] GetFileSystemEntries(string path)
{
DropCacheIfIdleOrFull();
return _cache.GetOrAdd(
path,
static (p, fileSystem) =>
static (p, state) =>
{
FileSystemMetadata[] entries;
try
{
return fileSystem.GetFileSystemEntries(p).ToArray();
entries = state.FileSystem.GetFileSystemEntries(p).ToArray();
}
catch (DirectoryNotFoundException)
{
return [];
entries = [];
}
Interlocked.Add(ref state.Service._recordCount, entries.Length + 1);
return entries;
},
_fileSystem);
(FileSystem: _fileSystem, Service: this));
}
public List<FileSystemMetadata> GetDirectories(string path)
@@ -89,13 +107,18 @@ namespace MediaBrowser.Controller.Providers
public FileSystemMetadata? GetFileSystemEntry(string path)
{
DropCacheIfIdleOrFull();
if (!_fileCache.TryGetValue(path, out var result))
{
var file = _fileSystem.GetFileSystemInfo(path);
if (file?.Exists ?? false)
{
result = file;
_fileCache.TryAdd(path, result);
if (_fileCache.TryAdd(path, result))
{
Interlocked.Increment(ref _recordCount);
}
}
}
@@ -107,32 +130,96 @@ namespace MediaBrowser.Controller.Providers
public IReadOnlyList<string> GetFilePaths(string path, bool clearCache)
{
if (clearCache)
if (clearCache && _filePathCache.TryRemove(path, out var cached))
{
_filePathCache.TryRemove(path, out _);
Interlocked.Add(ref _recordCount, -(cached.Count + 1));
}
DropCacheIfIdleOrFull();
var filePaths = _filePathCache.GetOrAdd(
path,
static (p, fileSystem) =>
static (p, state) =>
{
List<string> filePaths;
try
{
return fileSystem.GetFilePaths(p).OrderBy(x => x).ToList();
filePaths = state.FileSystem.GetFilePaths(p).OrderBy(x => x).ToList();
}
catch (DirectoryNotFoundException)
{
return [];
filePaths = [];
}
Interlocked.Add(ref state.Service._recordCount, filePaths.Count + 1);
return filePaths;
},
_fileSystem);
(FileSystem: _fileSystem, Service: this));
return filePaths;
}
public void Invalidate(string path)
{
Forget(path);
var parent = Path.GetDirectoryName(path);
if (!string.IsNullOrEmpty(parent))
{
Forget(parent);
}
}
public void Move(string source, string destination)
{
Directory.Move(source, destination);
Invalidate(source);
Invalidate(destination);
}
public bool IsAccessible(string path)
{
return _fileSystem.GetFileSystemEntryPaths(path).Any();
}
private void DropCacheIfIdleOrFull()
{
var nowMs = Environment.TickCount64;
var idleMs = nowMs - _lastAccess;
if (idleMs >= IdleTimeoutMs || _recordCount >= MaxCachedRecords)
{
_cache.Clear();
_fileCache.Clear();
_filePathCache.Clear();
_recordCount = 0;
_lastAccess = nowMs;
return;
}
if (idleMs >= AccessIntervalMs)
{
_lastAccess = nowMs;
}
}
private void Forget(string path)
{
if (_cache.TryRemove(path, out var entries))
{
Interlocked.Add(ref _recordCount, -(entries.Length + 1));
}
if (_fileCache.TryRemove(path, out _))
{
Interlocked.Decrement(ref _recordCount);
}
if (_filePathCache.TryRemove(path, out var filePaths))
{
Interlocked.Add(ref _recordCount, -(filePaths.Count + 1));
}
}
}
}
@@ -23,6 +23,19 @@ namespace MediaBrowser.Controller.Providers
IReadOnlyList<string> GetFilePaths(string path, bool clearCache);
/// <summary>
/// Forgets what is cached about a path and about the directory containing it.
/// </summary>
/// <param name="path">The file or directory path that changed.</param>
void Invalidate(string path);
/// <summary>
/// Moves a directory and forgets what is cached about both paths.
/// </summary>
/// <param name="source">The directory to move.</param>
/// <param name="destination">The path to move the directory to.</param>
void Move(string source, string destination);
bool IsAccessible(string path);
}
}