aboutsummaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorShadowghost <Ghost_of_Stone@web.de>2026-09-01 21:17:07 +0200
committerShadowghost <Ghost_of_Stone@web.de>2026-09-01 22:03:59 +0200
commit0e6c52f4311f334a98a61191fe850f226b37d79b (patch)
tree521bd4a85d8e70f550f2bae6c44bd093b5a89c98
parent76418ec530a141d5e3d82c2ca92f05e469a04813 (diff)
Guard the refresh queue so concurrent callers cannot lose the items they queue
-rw-r--r--MediaBrowser.Providers/Manager/ProviderManager.cs50
-rw-r--r--tests/Jellyfin.Providers.Tests/Manager/ProviderManagerTests.cs132
2 files changed, 165 insertions, 17 deletions
diff --git a/MediaBrowser.Providers/Manager/ProviderManager.cs b/MediaBrowser.Providers/Manager/ProviderManager.cs
index fbd9e5435e..4ff9ab1a52 100644
--- a/MediaBrowser.Providers/Manager/ProviderManager.cs
+++ b/MediaBrowser.Providers/Manager/ProviderManager.cs
@@ -1143,16 +1143,21 @@ namespace MediaBrowser.Providers.Manager
return;
}
- _refreshQueue.Enqueue((itemId, options), priority);
-
+ // PriorityQueue is not thread safe, and this runs on whichever thread queued the refresh
+ // while the processor dequeues on its own, so every touch of the queue takes the lock.
lock (_refreshQueueLock)
{
- if (!_isProcessingRefreshQueue)
+ _refreshQueue.Enqueue((itemId, options), priority);
+
+ if (_isProcessingRefreshQueue)
{
- _isProcessingRefreshQueue = true;
- Task.Run(StartProcessingRefreshQueue);
+ return;
}
+
+ _isProcessingRefreshQueue = true;
}
+
+ Task.Run(StartProcessingRefreshQueue);
}
private async Task StartProcessingRefreshQueue()
@@ -1161,17 +1166,35 @@ namespace MediaBrowser.Providers.Manager
if (_disposed)
{
+ lock (_refreshQueueLock)
+ {
+ _isProcessingRefreshQueue = false;
+ }
+
return;
}
var cancellationToken = _disposeCancellationTokenSource.Token;
libraryManager.ClearIgnoreRuleCache();
- while (_refreshQueue.TryDequeue(out var refreshItem, out _))
+
+ while (true)
{
- if (_disposed)
+ (Guid ItemId, MetadataRefreshOptions RefreshOptions) refreshItem;
+
+ // Standing down and taking the next entry happen under one lock, and this is the
+ // only place the flag is handed back. Releasing it anywhere else would leave a gap
+ // in which a refresh queued just after the queue ran dry sees a processor that has
+ // already stopped, and waits forever.
+ lock (_refreshQueueLock)
{
- return;
+ if (_disposed
+ || cancellationToken.IsCancellationRequested
+ || !_refreshQueue.TryDequeue(out refreshItem, out _))
+ {
+ _isProcessingRefreshQueue = false;
+ break;
+ }
}
try
@@ -1188,19 +1211,22 @@ namespace MediaBrowser.Providers.Manager
await task.ConfigureAwait(false);
}
- catch (OperationCanceledException)
+ catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested)
{
- break;
+ // Shutting down. Whatever is still queued keeps its place; the next pass round
+ // stands the processor down, so the next refresh queued starts one of its own.
+ continue;
}
catch (Exception ex)
{
+ // A provider that cancelled for its own reasons lands here too, an HTTP timeout
+ // above all. One unreachable metadata server must not stop the queue draining.
_logger.LogError(ex, "Error refreshing item");
}
}
- lock (_refreshQueueLock)
+ if (!_disposed)
{
- _isProcessingRefreshQueue = false;
libraryManager.ClearIgnoreRuleCache();
}
}
diff --git a/tests/Jellyfin.Providers.Tests/Manager/ProviderManagerTests.cs b/tests/Jellyfin.Providers.Tests/Manager/ProviderManagerTests.cs
index 5749944fcd..b3f5af76b7 100644
--- a/tests/Jellyfin.Providers.Tests/Manager/ProviderManagerTests.cs
+++ b/tests/Jellyfin.Providers.Tests/Manager/ProviderManagerTests.cs
@@ -1,4 +1,5 @@
using System;
+using System.Collections.Concurrent;
using System.Collections.Generic;
using System.Linq;
using System.Net.Http;
@@ -377,6 +378,122 @@ namespace Jellyfin.Providers.Tests.Manager
GetMetadataProviders_CanRefreshMetadata_Tester(providerType, expected, ownedItem: true);
}
+ [Fact]
+ public async Task QueueRefresh_ManyItemsQueuedFromManyThreads_ProcessesEveryOne()
+ {
+ // The queue is filled from whichever thread wants a refresh and drained by a processor of
+ // its own, so an unsynchronised PriorityQueue can lose entries outright, and a processor
+ // that stands down before releasing its flag leaves whatever was queued in that gap with
+ // nobody to drain it. Either way an item silently never gets refreshed.
+ const int ItemCount = 2000;
+
+ var queued = Enumerable.Range(0, ItemCount).Select(_ => Guid.NewGuid()).ToArray();
+ var processed = new ConcurrentBag<Guid>();
+ var allProcessed = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
+
+ var libraryManager = new Mock<ILibraryManager>();
+ libraryManager.Setup(i => i.GetItemById(It.IsAny<Guid>()))
+ .Returns((Guid id) =>
+ {
+ // Returning null drains the entry without needing the whole refresh machinery.
+ processed.Add(id);
+ if (processed.Count == ItemCount)
+ {
+ allProcessed.TrySetResult();
+ }
+
+ return null;
+ });
+
+ using var providerManager = GetProviderManager(libraryManager: libraryManager.Object);
+
+ await Parallel.ForEachAsync(
+ queued,
+ TestContext.Current.CancellationToken,
+ (id, _) =>
+ {
+ providerManager.QueueRefresh(id, new MetadataRefreshOptions(Mock.Of<IDirectoryService>()), RefreshPriority.Normal);
+ return ValueTask.CompletedTask;
+ });
+
+ using var timeout = CancellationTokenSource.CreateLinkedTokenSource(TestContext.Current.CancellationToken);
+ timeout.CancelAfter(TimeSpan.FromSeconds(30));
+
+ try
+ {
+ await allProcessed.Task.WaitAsync(timeout.Token);
+ }
+ catch (OperationCanceledException)
+ {
+ // Fall through, so the assertions below name what was lost rather than the wait.
+ }
+
+ Assert.Empty(providerManager.GetRefreshQueue());
+ Assert.Equal(queued.Order().ToArray(), processed.Order().ToArray());
+ }
+
+ [Fact]
+ public async Task QueueRefresh_RefreshCancelsForItsOwnReasons_KeepsDrainingTheQueue()
+ {
+ // MetadataService rethrows OperationCanceledException out of a provider, so an HTTP
+ // timeout against an unreachable metadata server arrives here looking exactly like a
+ // shutdown. Treating it as one stops the processor and strands the rest of the queue.
+ const int ItemCount = 200;
+
+ var queued = Enumerable.Range(0, ItemCount).Select(_ => Guid.NewGuid()).ToArray();
+ var processed = new ConcurrentBag<Guid>();
+ var allProcessed = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
+ using var allQueued = new ManualResetEventSlim(false);
+ var cancelledOnce = false;
+
+ var libraryManager = new Mock<ILibraryManager>();
+ libraryManager.Setup(i => i.GetItemById(It.IsAny<Guid>()))
+ .Returns((Guid id) =>
+ {
+ if (!cancelledOnce)
+ {
+ cancelledOnce = true;
+
+ // Hold the first entry until the whole batch is queued, so everything that
+ // follows it is already waiting when the cancellation lands.
+ allQueued.Wait(TimeSpan.FromSeconds(30));
+ throw new OperationCanceledException("provider timed out");
+ }
+
+ processed.Add(id);
+ if (processed.Count == ItemCount - 1)
+ {
+ allProcessed.TrySetResult();
+ }
+
+ return null;
+ });
+
+ using var providerManager = GetProviderManager(libraryManager: libraryManager.Object);
+
+ foreach (var id in queued)
+ {
+ providerManager.QueueRefresh(id, new MetadataRefreshOptions(Mock.Of<IDirectoryService>()), RefreshPriority.Normal);
+ }
+
+ allQueued.Set();
+
+ using var timeout = CancellationTokenSource.CreateLinkedTokenSource(TestContext.Current.CancellationToken);
+ timeout.CancelAfter(TimeSpan.FromSeconds(30));
+
+ try
+ {
+ await allProcessed.Task.WaitAsync(timeout.Token);
+ }
+ catch (OperationCanceledException)
+ {
+ // Fall through, so the assertions below name what was left stranded.
+ }
+
+ Assert.Empty(providerManager.GetRefreshQueue());
+ Assert.Equal(ItemCount - 1, processed.Count);
+ }
+
private static void GetMetadataProviders_CanRefreshMetadata_Tester(
string providerType,
bool expected,
@@ -554,15 +671,20 @@ namespace Jellyfin.Providers.Tests.Manager
private static ProviderManager GetProviderManager(
ServerConfiguration? serverConfiguration = null,
LibraryOptions? libraryOptions = null,
- IBaseItemManager? baseItemManager = null)
+ IBaseItemManager? baseItemManager = null,
+ ILibraryManager? libraryManager = null)
{
var serverConfigurationManager = new Mock<IServerConfigurationManager>(MockBehavior.Strict);
serverConfigurationManager.Setup(i => i.Configuration)
.Returns(serverConfiguration ?? new ServerConfiguration());
- var libraryManager = new Mock<ILibraryManager>(MockBehavior.Strict);
- libraryManager.Setup(i => i.GetLibraryOptions(It.IsAny<BaseItem>()))
- .Returns(libraryOptions ?? new LibraryOptions());
+ if (libraryManager is null)
+ {
+ var libraryManagerMock = new Mock<ILibraryManager>(MockBehavior.Strict);
+ libraryManagerMock.Setup(i => i.GetLibraryOptions(It.IsAny<BaseItem>()))
+ .Returns(libraryOptions ?? new LibraryOptions());
+ libraryManager = libraryManagerMock.Object;
+ }
var providerManager = new ProviderManager(
Mock.Of<IHttpClientFactory>(),
@@ -572,7 +694,7 @@ namespace Jellyfin.Providers.Tests.Manager
_logger,
Mock.Of<IFileSystem>(),
Mock.Of<IServerApplicationPaths>(),
- libraryManager.Object,
+ libraryManager,
baseItemManager!,
Mock.Of<ILyricManager>(),
Mock.Of<IMemoryCache>(),