Close a change batch on its own window so a library scan cannot grow it without bound
This commit is contained in:
@@ -27,6 +27,11 @@ namespace Emby.Server.Implementations.EntryPoints;
|
||||
/// </summary>
|
||||
public sealed class LibraryChangedNotifier : IHostedService, IDisposable
|
||||
{
|
||||
// A batch holds a live reference to every item it names, so it has to stay small enough that a
|
||||
// library scan - which changes items faster than any batch window closes - cannot grow it without
|
||||
// bound. Reached only by a scan; interactive use closes a batch on the window long before this.
|
||||
internal const int MaxBatchSize = 2000;
|
||||
|
||||
private readonly ILibraryManager _libraryManager;
|
||||
private readonly IServerConfigurationManager _configurationManager;
|
||||
private readonly IProviderManager _providerManager;
|
||||
@@ -35,11 +40,11 @@ public sealed class LibraryChangedNotifier : IHostedService, IDisposable
|
||||
private readonly ILogger<LibraryChangedNotifier> _logger;
|
||||
|
||||
private readonly Lock _libraryChangedSyncLock = new();
|
||||
private readonly List<Folder> _foldersAddedTo = new();
|
||||
private readonly List<Folder> _foldersRemovedFrom = new();
|
||||
private readonly List<BaseItem> _itemsAdded = new();
|
||||
private readonly List<BaseItem> _itemsRemoved = new();
|
||||
private readonly List<BaseItem> _itemsUpdated = new();
|
||||
private readonly Dictionary<Guid, Folder> _foldersAddedTo = [];
|
||||
private readonly Dictionary<Guid, Folder> _foldersRemovedFrom = [];
|
||||
private readonly Dictionary<Guid, BaseItem> _itemsAdded = [];
|
||||
private readonly Dictionary<Guid, BaseItem> _itemsRemoved = [];
|
||||
private readonly Dictionary<Guid, BaseItem> _itemsUpdated = [];
|
||||
private readonly ConcurrentDictionary<Guid, DateTime> _lastProgressMessageTimes = new();
|
||||
|
||||
private Timer? _libraryUpdateTimer;
|
||||
@@ -173,7 +178,7 @@ public sealed class LibraryChangedNotifier : IHostedService, IDisposable
|
||||
private void OnLibraryItemRemoved(object? sender, ItemChangeEventArgs e)
|
||||
=> OnLibraryChange(e.Item, e.Parent, _itemsRemoved, _foldersRemovedFrom);
|
||||
|
||||
private void OnLibraryChange(BaseItem item, BaseItem parent, List<BaseItem> itemsList, List<Folder>? foldersList)
|
||||
private void OnLibraryChange(BaseItem item, BaseItem parent, Dictionary<Guid, BaseItem> itemsList, Dictionary<Guid, Folder>? foldersList)
|
||||
{
|
||||
if (!FilterItem(item))
|
||||
{
|
||||
@@ -182,23 +187,28 @@ public sealed class LibraryChangedNotifier : IHostedService, IDisposable
|
||||
|
||||
lock (_libraryChangedSyncLock)
|
||||
{
|
||||
var updateDuration = TimeSpan.FromSeconds(_configurationManager.Configuration.LibraryUpdateDuration);
|
||||
|
||||
// The window runs from the first change of a batch and is never extended. Extending it on
|
||||
// every change would keep a library scan's batch open for the whole scan, and the batch
|
||||
// holds the items it names alive, so it would grow to the size of the library.
|
||||
if (_libraryUpdateTimer is null)
|
||||
{
|
||||
var updateDuration = TimeSpan.FromSeconds(_configurationManager.Configuration.LibraryUpdateDuration);
|
||||
_libraryUpdateTimer = new Timer(LibraryUpdateTimerCallback, null, updateDuration, Timeout.InfiniteTimeSpan);
|
||||
}
|
||||
else
|
||||
{
|
||||
_libraryUpdateTimer.Change(updateDuration, Timeout.InfiniteTimeSpan);
|
||||
}
|
||||
|
||||
if (foldersList is not null && parent is Folder folder)
|
||||
{
|
||||
foldersList.Add(folder);
|
||||
foldersList[folder.Id] = folder;
|
||||
}
|
||||
|
||||
itemsList.Add(item);
|
||||
itemsList[item.Id] = item;
|
||||
|
||||
// A window long enough to cover a burst still has to give way once the batch is large
|
||||
// enough to be worth sending on its own.
|
||||
if (_itemsAdded.Count + _itemsRemoved.Count + _itemsUpdated.Count >= MaxBatchSize)
|
||||
{
|
||||
_libraryUpdateTimer.Change(TimeSpan.Zero, Timeout.InfiniteTimeSpan);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -211,22 +221,16 @@ public sealed class LibraryChangedNotifier : IHostedService, IDisposable
|
||||
List<BaseItem> itemsRemoved;
|
||||
lock (_libraryChangedSyncLock)
|
||||
{
|
||||
// Remove dupes in case some were saved multiple times
|
||||
foldersAddedTo = _foldersAddedTo
|
||||
.DistinctBy(x => x.Id)
|
||||
.ToList();
|
||||
|
||||
foldersRemovedFrom = _foldersRemovedFrom
|
||||
.DistinctBy(x => x.Id)
|
||||
.ToList();
|
||||
foldersAddedTo = _foldersAddedTo.Values.ToList();
|
||||
foldersRemovedFrom = _foldersRemovedFrom.Values.ToList();
|
||||
|
||||
itemsUpdated = _itemsUpdated
|
||||
.Where(i => !_itemsAdded.Contains(i))
|
||||
.DistinctBy(x => x.Id)
|
||||
.Where(e => !_itemsAdded.ContainsKey(e.Key))
|
||||
.Select(e => e.Value)
|
||||
.ToList();
|
||||
|
||||
itemsAdded = _itemsAdded.ToList();
|
||||
itemsRemoved = _itemsRemoved.ToList();
|
||||
itemsAdded = _itemsAdded.Values.ToList();
|
||||
itemsRemoved = _itemsRemoved.Values.ToList();
|
||||
|
||||
if (_libraryUpdateTimer is not null)
|
||||
{
|
||||
@@ -241,6 +245,15 @@ public sealed class LibraryChangedNotifier : IHostedService, IDisposable
|
||||
_foldersRemovedFrom.Clear();
|
||||
}
|
||||
|
||||
if (itemsAdded.Count == 0
|
||||
&& itemsUpdated.Count == 0
|
||||
&& itemsRemoved.Count == 0
|
||||
&& foldersAddedTo.Count == 0
|
||||
&& foldersRemovedFrom.Count == 0)
|
||||
{
|
||||
return;
|
||||
}
|
||||
|
||||
await SendChangeNotifications(itemsAdded, itemsUpdated, itemsRemoved, foldersAddedTo, foldersRemovedFrom, CancellationToken.None).ConfigureAwait(false);
|
||||
}
|
||||
|
||||
|
||||
@@ -18,15 +18,17 @@ namespace Emby.Server.Implementations.EntryPoints
|
||||
public sealed class UserDataChangeNotifier : IHostedService, IDisposable
|
||||
{
|
||||
private const int UpdateDuration = 500;
|
||||
internal const int MaxBatchSize = 2000;
|
||||
|
||||
private readonly ISessionManager _sessionManager;
|
||||
private readonly IUserDataManager _userDataManager;
|
||||
private readonly IUserManager _userManager;
|
||||
|
||||
private readonly Dictionary<Guid, List<BaseItem>> _changedItems = new();
|
||||
private readonly Dictionary<Guid, Dictionary<Guid, BaseItem>> _changedItems = [];
|
||||
private readonly Lock _syncLock = new();
|
||||
|
||||
private Timer? _updateTimer;
|
||||
private int _changedItemCount;
|
||||
|
||||
/// <summary>
|
||||
/// Initializes a new instance of the <see cref="UserDataChangeNotifier"/> class.
|
||||
@@ -69,50 +71,64 @@ namespace Emby.Server.Implementations.EntryPoints
|
||||
|
||||
lock (_syncLock)
|
||||
{
|
||||
if (_updateTimer is null)
|
||||
{
|
||||
_updateTimer = new Timer(
|
||||
UpdateTimerCallback,
|
||||
null,
|
||||
UpdateDuration,
|
||||
Timeout.Infinite);
|
||||
}
|
||||
else
|
||||
{
|
||||
_updateTimer.Change(UpdateDuration, Timeout.Infinite);
|
||||
}
|
||||
// The window runs from the first change of a batch and is never extended, so a stream
|
||||
// of changes that never pauses - a library scan - still closes its batches instead of
|
||||
// holding every item it touched alive until the stream stops.
|
||||
_updateTimer ??= new Timer(
|
||||
UpdateTimerCallback,
|
||||
null,
|
||||
UpdateDuration,
|
||||
Timeout.Infinite);
|
||||
|
||||
if (!_changedItems.TryGetValue(e.UserId, out List<BaseItem>? keys))
|
||||
if (!_changedItems.TryGetValue(e.UserId, out Dictionary<Guid, BaseItem>? keys))
|
||||
{
|
||||
keys = new List<BaseItem>();
|
||||
keys = [];
|
||||
_changedItems[e.UserId] = keys;
|
||||
}
|
||||
|
||||
keys.Add(e.Item);
|
||||
|
||||
var baseItem = e.Item;
|
||||
|
||||
// Go up one level for indicators
|
||||
if (baseItem is not null)
|
||||
{
|
||||
Track(keys, baseItem);
|
||||
|
||||
var parent = baseItem.GetOwner() ?? baseItem.GetParent();
|
||||
|
||||
if (parent is not null)
|
||||
{
|
||||
keys.Add(parent);
|
||||
Track(keys, parent);
|
||||
}
|
||||
}
|
||||
|
||||
// A window long enough to cover a burst still has to give way once the batch is
|
||||
// large enough to be worth sending on its own.
|
||||
if (_changedItemCount >= MaxBatchSize)
|
||||
{
|
||||
_updateTimer.Change(0, Timeout.Infinite);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private void Track(Dictionary<Guid, BaseItem> keys, BaseItem item)
|
||||
{
|
||||
var before = keys.Count;
|
||||
keys[item.Id] = item;
|
||||
|
||||
if (keys.Count != before)
|
||||
{
|
||||
_changedItemCount++;
|
||||
}
|
||||
}
|
||||
|
||||
private async void UpdateTimerCallback(object? state)
|
||||
{
|
||||
List<KeyValuePair<Guid, List<BaseItem>>> changes;
|
||||
List<KeyValuePair<Guid, Dictionary<Guid, BaseItem>>> changes;
|
||||
lock (_syncLock)
|
||||
{
|
||||
// Remove dupes in case some were saved multiple times
|
||||
changes = _changedItems.ToList();
|
||||
_changedItems.Clear();
|
||||
_changedItemCount = 0;
|
||||
|
||||
if (_updateTimer is not null)
|
||||
{
|
||||
@@ -121,17 +137,22 @@ namespace Emby.Server.Implementations.EntryPoints
|
||||
}
|
||||
}
|
||||
|
||||
if (changes.Count == 0)
|
||||
{
|
||||
return;
|
||||
}
|
||||
|
||||
foreach (var (userId, changedItems) in changes)
|
||||
{
|
||||
await _sessionManager.SendMessageToUserSessions(
|
||||
[userId],
|
||||
SessionMessageType.UserDataChanged,
|
||||
() => GetUserDataChangeInfo(userId, changedItems),
|
||||
() => GetUserDataChangeInfo(userId, changedItems.Values),
|
||||
default).ConfigureAwait(false);
|
||||
}
|
||||
}
|
||||
|
||||
private UserDataChangeInfo GetUserDataChangeInfo(Guid userId, List<BaseItem> changedItems)
|
||||
private UserDataChangeInfo GetUserDataChangeInfo(Guid userId, IEnumerable<BaseItem> changedItems)
|
||||
{
|
||||
var user = _userManager.GetUserById(userId)
|
||||
?? throw new ArgumentException("Invalid user ID", nameof(userId));
|
||||
@@ -140,7 +161,6 @@ namespace Emby.Server.Implementations.EntryPoints
|
||||
{
|
||||
UserId = userId,
|
||||
UserDataList = changedItems
|
||||
.DistinctBy(x => x.Id)
|
||||
.Select(i =>
|
||||
{
|
||||
var dto = _userDataManager.GetUserDataDto(i, user);
|
||||
|
||||
+123
@@ -0,0 +1,123 @@
|
||||
using System;
|
||||
using System.Collections.Generic;
|
||||
using System.Diagnostics;
|
||||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
using Emby.Server.Implementations.EntryPoints;
|
||||
using MediaBrowser.Controller.Configuration;
|
||||
using MediaBrowser.Controller.Entities;
|
||||
using MediaBrowser.Controller.Library;
|
||||
using MediaBrowser.Controller.Providers;
|
||||
using MediaBrowser.Controller.Session;
|
||||
using MediaBrowser.Model.Configuration;
|
||||
using Microsoft.Extensions.Logging.Abstractions;
|
||||
using Moq;
|
||||
using Xunit;
|
||||
|
||||
namespace Jellyfin.Server.Implementations.Tests.EntryPoints;
|
||||
|
||||
public class LibraryChangedNotifierTests
|
||||
{
|
||||
// How long a test waits for the notifier's timer callback to run. Generous: the assertions are
|
||||
// about a batch being sent at all, not about how promptly.
|
||||
private static readonly TimeSpan _flushTimeout = TimeSpan.FromSeconds(15);
|
||||
|
||||
private readonly Mock<ILibraryManager> _libraryManager = new();
|
||||
private readonly Mock<IServerConfigurationManager> _configurationManager = new();
|
||||
private readonly Mock<ISessionManager> _sessionManager = new();
|
||||
private readonly Mock<IUserManager> _userManager = new();
|
||||
private readonly Mock<IProviderManager> _providerManager = new();
|
||||
private readonly ServerConfiguration _configuration = new();
|
||||
|
||||
private int _flushCount;
|
||||
|
||||
public LibraryChangedNotifierTests()
|
||||
{
|
||||
_configurationManager.SetupGet(e => e.Configuration).Returns(_configuration);
|
||||
|
||||
// Reading the session list is the first thing a flush does, so it stands in for "a batch was
|
||||
// sent" without having to mock a whole user library behind it.
|
||||
_sessionManager.SetupGet(e => e.Sessions)
|
||||
.Returns(() =>
|
||||
{
|
||||
Interlocked.Increment(ref _flushCount);
|
||||
return [];
|
||||
});
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public async Task OnLibraryItemUpdated_BatchSizeCapReached_SendsWithoutWaitingForWindow()
|
||||
{
|
||||
// Long enough that only the size cap can close the batch.
|
||||
_configuration.LibraryUpdateDuration = 3600;
|
||||
|
||||
var notifier = CreateNotifier();
|
||||
await notifier.StartAsync(TestContext.Current.CancellationToken);
|
||||
|
||||
for (var i = 0; i < LibraryChangedNotifier.MaxBatchSize; i++)
|
||||
{
|
||||
RaiseItemUpdated();
|
||||
}
|
||||
|
||||
Assert.True(await WaitForFlushAsync(1), "The batch was not sent once it hit the size cap.");
|
||||
|
||||
await notifier.StopAsync(TestContext.Current.CancellationToken);
|
||||
notifier.Dispose();
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public async Task OnLibraryItemUpdated_ChangesNeverPause_StillSendsOnTheWindow()
|
||||
{
|
||||
// A scan changes items continuously. The window must run from the first change of a batch, or
|
||||
// the batch never closes and holds every item it named alive for the length of the scan.
|
||||
_configuration.LibraryUpdateDuration = 1;
|
||||
|
||||
var notifier = CreateNotifier();
|
||||
await notifier.StartAsync(TestContext.Current.CancellationToken);
|
||||
|
||||
var stopwatch = Stopwatch.StartNew();
|
||||
while (stopwatch.Elapsed < _flushTimeout && Volatile.Read(ref _flushCount) == 0)
|
||||
{
|
||||
// Well below the window, and well below the size cap over the whole loop.
|
||||
RaiseItemUpdated();
|
||||
await Task.Delay(25, TestContext.Current.CancellationToken);
|
||||
}
|
||||
|
||||
Assert.True(Volatile.Read(ref _flushCount) > 0, "The batch was never sent while changes kept arriving.");
|
||||
|
||||
await notifier.StopAsync(TestContext.Current.CancellationToken);
|
||||
notifier.Dispose();
|
||||
}
|
||||
|
||||
private LibraryChangedNotifier CreateNotifier()
|
||||
=> new(
|
||||
_libraryManager.Object,
|
||||
_configurationManager.Object,
|
||||
_sessionManager.Object,
|
||||
_userManager.Object,
|
||||
NullLogger<LibraryChangedNotifier>.Instance,
|
||||
_providerManager.Object);
|
||||
|
||||
// A folder passes the notifier's item filter without needing any of BaseItem's static services.
|
||||
private void RaiseItemUpdated()
|
||||
=> _libraryManager.Raise(
|
||||
e => e.ItemUpdated += null,
|
||||
_libraryManager.Object,
|
||||
new ItemChangeEventArgs { Item = new Folder { Id = Guid.NewGuid() } });
|
||||
|
||||
private async Task<bool> WaitForFlushAsync(int expected)
|
||||
{
|
||||
var stopwatch = Stopwatch.StartNew();
|
||||
while (stopwatch.Elapsed < _flushTimeout)
|
||||
{
|
||||
if (Volatile.Read(ref _flushCount) >= expected)
|
||||
{
|
||||
return true;
|
||||
}
|
||||
|
||||
await Task.Delay(25, TestContext.Current.CancellationToken);
|
||||
}
|
||||
|
||||
return false;
|
||||
}
|
||||
}
|
||||
+78
@@ -0,0 +1,78 @@
|
||||
using System;
|
||||
using System.Diagnostics;
|
||||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
using Emby.Server.Implementations.EntryPoints;
|
||||
using MediaBrowser.Controller.Entities;
|
||||
using MediaBrowser.Controller.Library;
|
||||
using MediaBrowser.Controller.Session;
|
||||
using MediaBrowser.Model.Entities;
|
||||
using MediaBrowser.Model.Session;
|
||||
using Moq;
|
||||
using Xunit;
|
||||
|
||||
namespace Jellyfin.Server.Implementations.Tests.EntryPoints;
|
||||
|
||||
public class UserDataChangeNotifierTests
|
||||
{
|
||||
// How long a test waits for the notifier's timer callback to run. Generous: the assertions are
|
||||
// about a batch being sent at all, not about how promptly.
|
||||
private static readonly TimeSpan _flushTimeout = TimeSpan.FromSeconds(15);
|
||||
|
||||
private readonly Mock<IUserDataManager> _userDataManager = new();
|
||||
private readonly Mock<ISessionManager> _sessionManager = new();
|
||||
private readonly Mock<IUserManager> _userManager = new();
|
||||
|
||||
private int _flushCount;
|
||||
|
||||
public UserDataChangeNotifierTests()
|
||||
{
|
||||
_sessionManager
|
||||
.Setup(e => e.SendMessageToUserSessions(
|
||||
It.IsAny<System.Collections.Generic.List<Guid>>(),
|
||||
SessionMessageType.UserDataChanged,
|
||||
It.IsAny<Func<UserDataChangeInfo>>(),
|
||||
It.IsAny<CancellationToken>()))
|
||||
.Callback(() => Interlocked.Increment(ref _flushCount))
|
||||
.Returns(Task.CompletedTask);
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public async Task OnUserDataSaved_ChangesNeverPause_StillSendsOnTheWindow()
|
||||
{
|
||||
// A scan changes user data continuously. The window must run from the first change of a batch,
|
||||
// or the batch never closes and holds every item it named alive for the length of the scan.
|
||||
var notifier = CreateNotifier();
|
||||
await notifier.StartAsync(TestContext.Current.CancellationToken);
|
||||
|
||||
var userId = Guid.NewGuid();
|
||||
var stopwatch = Stopwatch.StartNew();
|
||||
while (stopwatch.Elapsed < _flushTimeout && Volatile.Read(ref _flushCount) == 0)
|
||||
{
|
||||
// Well below the window, and well below the size cap over the whole loop.
|
||||
RaiseUserDataSaved(userId);
|
||||
await Task.Delay(25, TestContext.Current.CancellationToken);
|
||||
}
|
||||
|
||||
Assert.True(Volatile.Read(ref _flushCount) > 0, "The batch was never sent while changes kept arriving.");
|
||||
|
||||
await notifier.StopAsync(TestContext.Current.CancellationToken);
|
||||
notifier.Dispose();
|
||||
}
|
||||
|
||||
private UserDataChangeNotifier CreateNotifier()
|
||||
=> new(_userDataManager.Object, _sessionManager.Object, _userManager.Object);
|
||||
|
||||
// A folder needs none of BaseItem's static services, and PlaybackProgress is the one reason the
|
||||
// notifier ignores outright.
|
||||
private void RaiseUserDataSaved(Guid userId)
|
||||
=> _userDataManager.Raise(
|
||||
e => e.UserDataSaved += null,
|
||||
_userDataManager.Object,
|
||||
new UserDataSaveEventArgs
|
||||
{
|
||||
UserId = userId,
|
||||
SaveReason = UserDataSaveReason.UpdateUserRating,
|
||||
Item = new Folder { Id = Guid.NewGuid() }
|
||||
});
|
||||
}
|
||||
Reference in New Issue
Block a user