diff --git a/Jellyfin/backend/Api/CustomRowController.cs b/Jellyfin/backend/Api/CustomRowController.cs index 5b09d3a..380fc28 100644 --- a/Jellyfin/backend/Api/CustomRowController.cs +++ b/Jellyfin/backend/Api/CustomRowController.cs @@ -9,6 +9,7 @@ using Microsoft.AspNetCore.Http; using Microsoft.AspNetCore.Mvc; using Microsoft.Extensions.Logging; +using Moonfin.Server.Helpers; using Moonfin.Server.Services; namespace Moonfin.Server.Api; @@ -23,6 +24,12 @@ public class CustomRowController : ControllerBase private static readonly TimeSpan CacheTtl = TimeSpan.FromHours(24); + // A large list can take longer to resolve than a client waits, so a resolution isn't tied + // to the request that started it. It still caches when that client gives up, and a retry + // while it's running waits on it instead of starting over. Keyed per user as well, since + // each user resolves with their own API keys. + private static readonly InFlightTasks> Resolving = new(); + public CustomRowController( CustomRowCacheService cacheService, CustomRowFetchService fetchService, @@ -95,16 +102,12 @@ public async Task> GetCustomRowItems( try { - var items = await _fetchService.FetchCustomRowAsync(source, type, parsedParams, userId.Value, cancellationToken); - - // Don't cache empty results. An empty row cached for 24h looks like a - // broken list when the real cause was an upstream hiccup. - if (items.Count > 0) - { - _cacheService.Set(cacheKey, items); - _cacheService.PruneOlderThan(TimeSpan.FromDays(7)); - await _cacheService.FlushAsync(); - } + var items = await Resolving + .GetOrStart( + $"{cacheKey}:{userId.Value}", + () => FetchAndCacheAsync(cacheKey, source, type, parsedParams, userId.Value)) + .WaitAsync(cancellationToken) + .ConfigureAwait(false); return Ok(new CustomRowResponse { @@ -112,6 +115,10 @@ public async Task> GetCustomRowItems( Items = items }); } + catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested) + { + return new EmptyResult(); + } catch (NotSupportedException ex) { return BadRequest(new { Error = ex.Message }); @@ -126,6 +133,29 @@ public async Task> GetCustomRowItems( }); } } + + private async Task> FetchAndCacheAsync( + string cacheKey, + string source, + string type, + Dictionary parsedParams, + Guid userId) + { + var items = await _fetchService + .FetchCustomRowAsync(source, type, parsedParams, userId, CancellationToken.None) + .ConfigureAwait(false); + + // Don't cache empty results. An empty row cached for 24h looks like a + // broken list when the real cause was an upstream hiccup. + if (items.Count > 0) + { + _cacheService.Set(cacheKey, items); + _cacheService.PruneOlderThan(TimeSpan.FromDays(7)); + await _cacheService.FlushAsync().ConfigureAwait(false); + } + + return items; + } } public class CustomRowResponse diff --git a/Jellyfin/backend/Helpers/InFlightTasks.cs b/Jellyfin/backend/Helpers/InFlightTasks.cs new file mode 100644 index 0000000..15942e9 --- /dev/null +++ b/Jellyfin/backend/Helpers/InFlightTasks.cs @@ -0,0 +1,30 @@ +using System.Collections.Concurrent; + +namespace Moonfin.Server.Helpers; + +/// +/// Runs one task per key at a time. Asking for a key that is still running hands back the +/// running task instead of starting another, and the key frees up once it finishes. +/// +internal sealed class InFlightTasks +{ + private readonly ConcurrentDictionary>> _running = new(); + + public Task GetOrStart(string key, Func> start) + { + // The async wrapper turns a start that throws into a failed task, so a bad start + // still frees its key. + var entry = new Lazy>(async () => await start().ConfigureAwait(false)); + var current = _running.GetOrAdd(key, entry); + if (ReferenceEquals(current, entry)) + { + _ = entry.Value.ContinueWith( + _ => _running.TryRemove(KeyValuePair.Create(key, entry)), + CancellationToken.None, + TaskContinuationOptions.ExecuteSynchronously, + TaskScheduler.Default); + } + + return current.Value; + } +} diff --git a/Jellyfin/tests/Moonfin.Server.Tests/InFlightTasksTests.cs b/Jellyfin/tests/Moonfin.Server.Tests/InFlightTasksTests.cs new file mode 100644 index 0000000..7b069ee --- /dev/null +++ b/Jellyfin/tests/Moonfin.Server.Tests/InFlightTasksTests.cs @@ -0,0 +1,98 @@ +using Moonfin.Server.Helpers; +using Xunit; + +namespace Moonfin.Server.Tests; + +/// +/// Custom rows resolve through , so a slow list keeps resolving +/// after the client that asked for it gives up, and a retry joins it instead of starting over. +/// +public sealed class InFlightTasksTests +{ + [Fact] + public async Task ASecondCallerJoinsTheRunningTask() + { + var tasks = new InFlightTasks(); + var gate = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var starts = 0; + Task Start() + { + starts++; + return gate.Task; + } + + var first = tasks.GetOrStart("row", Start); + var second = tasks.GetOrStart("row", Start); + gate.SetResult(7); + + Assert.Same(first, second); + Assert.Equal(7, await second); + Assert.Equal(1, starts); + } + + [Fact] + public async Task DifferentKeysRunSeparately() + { + var tasks = new InFlightTasks(); + var gate = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + + var first = tasks.GetOrStart("row:userA", () => gate.Task); + var second = tasks.GetOrStart("row:userB", () => Task.FromResult(2)); + + Assert.NotSame(first, second); + Assert.Equal(2, await second); + gate.SetResult(1); + Assert.Equal(1, await first); + } + + [Fact] + public async Task AFinishedKeyStartsAgain() + { + var tasks = new InFlightTasks(); + var starts = 0; + Task Start() => Task.FromResult(++starts); + + Assert.Equal(1, await tasks.GetOrStart("row", Start)); + Assert.Equal(2, await tasks.GetOrStart("row", Start)); + } + + [Fact] + public async Task AFailedRunFreesItsKey() + { + var tasks = new InFlightTasks(); + + await Assert.ThrowsAsync( + () => tasks.GetOrStart("row", () => Task.FromException(new InvalidOperationException()))); + + Assert.Equal(3, await tasks.GetOrStart("row", () => Task.FromResult(3))); + } + + [Fact] + public async Task AStartThatThrowsFailsTheTaskAndFreesItsKey() + { + var tasks = new InFlightTasks(); + + await Assert.ThrowsAsync( + () => tasks.GetOrStart("row", () => throw new InvalidOperationException())); + + Assert.Equal(3, await tasks.GetOrStart("row", () => Task.FromResult(3))); + } + + [Fact] + public async Task ACallerThatStopsWaitingLeavesTheRunGoing() + { + var tasks = new InFlightTasks(); + var gate = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + using var client = new CancellationTokenSource(); + + var run = tasks.GetOrStart("row", () => gate.Task); + var wait = run.WaitAsync(client.Token); + client.Cancel(); + + await Assert.ThrowsAnyAsync(() => wait); + Assert.False(run.IsCompleted); + + gate.SetResult(5); + Assert.Equal(5, await run); + } +}