Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
50 changes: 40 additions & 10 deletions Jellyfin/backend/Api/CustomRowController.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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<List<CustomRowItem>> Resolving = new();

public CustomRowController(
CustomRowCacheService cacheService,
CustomRowFetchService fetchService,
Expand Down Expand Up @@ -95,23 +102,23 @@ public async Task<ActionResult<CustomRowResponse>> 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
{
Success = true,
Items = items
});
}
catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested)
{
return new EmptyResult();
}
catch (NotSupportedException ex)
{
return BadRequest(new { Error = ex.Message });
Expand All @@ -126,6 +133,29 @@ public async Task<ActionResult<CustomRowResponse>> GetCustomRowItems(
});
}
}

private async Task<List<CustomRowItem>> FetchAndCacheAsync(
string cacheKey,
string source,
string type,
Dictionary<string, string> 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
Expand Down
30 changes: 30 additions & 0 deletions Jellyfin/backend/Helpers/InFlightTasks.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,30 @@
using System.Collections.Concurrent;

namespace Moonfin.Server.Helpers;

/// <summary>
/// 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.
/// </summary>
internal sealed class InFlightTasks<T>
{
private readonly ConcurrentDictionary<string, Lazy<Task<T>>> _running = new();

public Task<T> GetOrStart(string key, Func<Task<T>> 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<Task<T>>(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;
}
}
98 changes: 98 additions & 0 deletions Jellyfin/tests/Moonfin.Server.Tests/InFlightTasksTests.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,98 @@
using Moonfin.Server.Helpers;
using Xunit;

namespace Moonfin.Server.Tests;

/// <summary>
/// Custom rows resolve through <see cref="InFlightTasks{T}"/>, so a slow list keeps resolving
/// after the client that asked for it gives up, and a retry joins it instead of starting over.
/// </summary>
public sealed class InFlightTasksTests
{
[Fact]
public async Task ASecondCallerJoinsTheRunningTask()
{
var tasks = new InFlightTasks<int>();
var gate = new TaskCompletionSource<int>(TaskCreationOptions.RunContinuationsAsynchronously);
var starts = 0;
Task<int> 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<int>();
var gate = new TaskCompletionSource<int>(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<int>();
var starts = 0;
Task<int> 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<int>();

await Assert.ThrowsAsync<InvalidOperationException>(
() => tasks.GetOrStart("row", () => Task.FromException<int>(new InvalidOperationException())));

Assert.Equal(3, await tasks.GetOrStart("row", () => Task.FromResult(3)));
}

[Fact]
public async Task AStartThatThrowsFailsTheTaskAndFreesItsKey()
{
var tasks = new InFlightTasks<int>();

await Assert.ThrowsAsync<InvalidOperationException>(
() => 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<int>();
var gate = new TaskCompletionSource<int>(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<OperationCanceledException>(() => wait);
Assert.False(run.IsCompleted);

gate.SetResult(5);
Assert.Equal(5, await run);
}
}
Loading