Improve proxy cache behavior and diagnostics

This commit is contained in:
ajp_anton
2026-07-20 02:17:33 +00:00
parent fd1cf12b6b
commit 25088ce91d
10 changed files with 504 additions and 65 deletions
@@ -33,6 +33,8 @@ public sealed class ItemsProxyCache
public required bool CacheStored { get; init; }
public required bool InFlightCoalesced { get; init; }
public required int StatusCode { get; init; }
public required int ItemCount { get; init; }
@@ -229,7 +231,8 @@ public sealed class ItemsProxyCache
long totalMs,
long upstreamMs,
long transformMs,
long sizeBytes)
long sizeBytes,
bool inFlightCoalesced = false)
{
lock (_lock)
{
@@ -239,6 +242,7 @@ public sealed class ItemsProxyCache
CreatedUtc = DateTimeOffset.UtcNow,
CacheHit = cacheHit,
CacheStored = cacheStored,
InFlightCoalesced = inFlightCoalesced,
StatusCode = statusCode,
ItemCount = itemCount,
TotalMs = totalMs,
@@ -265,6 +269,7 @@ public sealed class ItemsProxyCache
AgeSeconds = (long)(now - e.CreatedUtc).TotalSeconds,
e.CacheHit,
e.CacheStored,
e.InFlightCoalesced,
e.StatusCode,
e.ItemCount,
e.TotalMs,
@@ -0,0 +1,110 @@
using System.Text.Json.Nodes;
using Microsoft.Extensions.Logging;
namespace Jellyfin.Plugin.Multilang.Services;
public sealed class ItemsProxyPrecacheService
{
private static readonly TimeSpan IdleThreshold = TimeSpan.FromMinutes(30);
private readonly IHttpClientFactory _httpClientFactory;
private readonly ILogger<ItemsProxyPrecacheService> _logger;
private readonly object _lock = new();
private readonly Dictionary<string, DateTimeOffset> _lastSeen = new(StringComparer.Ordinal);
private readonly HashSet<string> _pendingUsers = new(StringComparer.Ordinal);
public ItemsProxyPrecacheService(IHttpClientFactory httpClientFactory, ILogger<ItemsProxyPrecacheService> logger)
{
_httpClientFactory = httpClientFactory;
_logger = logger;
}
public void ObserveUserActivity(
Uri baseUri,
string userId,
string token,
string clientLocale,
bool precacheMovieLibraries,
bool precacheTvShowLibraries)
{
if (!precacheMovieLibraries && !precacheTvShowLibraries)
return;
var now = DateTimeOffset.UtcNow;
lock (_lock)
{
var wasRecentlyActive = _lastSeen.TryGetValue(userId, out var lastSeen) && now - lastSeen < IdleThreshold;
_lastSeen[userId] = now;
if (wasRecentlyActive || !_pendingUsers.Add(userId))
return;
}
_ = PrecacheAsync(baseUri, userId, token, clientLocale, precacheMovieLibraries, precacheTvShowLibraries);
}
private async Task PrecacheAsync(
Uri baseUri,
string userId,
string token,
string clientLocale,
bool precacheMovieLibraries,
bool precacheTvShowLibraries)
{
try
{
var views = await GetJsonAsync(new Uri(baseUri, $"Users/{userId}/Views"), token).ConfigureAwait(false);
var items = views["Items"]?.AsArray() ?? [];
var targets = items
.Select(view => new
{
Id = view?["Id"]?.GetValue<string>() ?? string.Empty,
CollectionType = view?["CollectionType"]?.GetValue<string>() ?? string.Empty
})
.Where(view => !string.IsNullOrWhiteSpace(view.Id))
.Where(view =>
(precacheMovieLibraries && view.CollectionType.Equals("movies", StringComparison.OrdinalIgnoreCase)) ||
(precacheTvShowLibraries && view.CollectionType.Equals("tvshows", StringComparison.OrdinalIgnoreCase)))
.ToArray();
foreach (var target in targets)
{
var itemType = target.CollectionType.Equals("movies", StringComparison.OrdinalIgnoreCase) ? "Movie" : "Series";
var itemsUrl = $"/Users/{userId}/Items?ParentId={Uri.EscapeDataString(target.Id)}&IncludeItemTypes={itemType}&Recursive=true&SortBy=SortName&SortOrder=Ascending&StartIndex=0&Limit=100";
var endpoint = "Multilang/ItemsProxy?url=" + Uri.EscapeDataString(itemsUrl) + "&precache=true";
if (!string.IsNullOrWhiteSpace(clientLocale))
endpoint += "&mlLocale=" + Uri.EscapeDataString(clientLocale);
using var request = new HttpRequestMessage(HttpMethod.Get, new Uri(baseUri, endpoint));
AddTokenHeaders(request, token);
using var response = await _httpClientFactory.CreateClient().SendAsync(request, CancellationToken.None).ConfigureAwait(false);
if (!response.IsSuccessStatusCode)
_logger.LogWarning("Multilang pre-cache failed user={UserId} library={LibraryId} status={StatusCode}", userId, target.Id, (int)response.StatusCode);
}
}
catch (Exception ex)
{
_logger.LogWarning(ex, "Multilang pre-cache failed for user={UserId}", userId);
}
finally
{
lock (_lock)
_pendingUsers.Remove(userId);
}
}
private async Task<JsonObject> GetJsonAsync(Uri uri, string token)
{
using var request = new HttpRequestMessage(HttpMethod.Get, uri);
AddTokenHeaders(request, token);
using var response = await _httpClientFactory.CreateClient().SendAsync(request, CancellationToken.None).ConfigureAwait(false);
response.EnsureSuccessStatusCode();
var body = await response.Content.ReadAsStringAsync(CancellationToken.None).ConfigureAwait(false);
return JsonNode.Parse(body)?.AsObject() ?? throw new InvalidOperationException("Jellyfin views response was not a JSON object.");
}
private static void AddTokenHeaders(HttpRequestMessage request, string token)
{
request.Headers.TryAddWithoutValidation("Authorization", $"MediaBrowser Token=\"{token}\"");
request.Headers.TryAddWithoutValidation("X-Emby-Token", token);
request.Headers.TryAddWithoutValidation("X-MediaBrowser-Token", token);
}
}
@@ -0,0 +1,68 @@
namespace Jellyfin.Plugin.Multilang.Services;
public readonly record struct ItemsProxyUpstreamResponse(int StatusCode, string Body, string ContentType)
{
public bool IsSuccessStatusCode => StatusCode is >= 200 and < 300;
}
public readonly record struct CoalescedItemsProxyResponse(ItemsProxyUpstreamResponse Response, bool JoinedExistingRequest);
public sealed class ItemsProxyRequestCoalescer
{
private sealed class InFlightRequest
{
public TaskCompletionSource<ItemsProxyUpstreamResponse> Completion { get; } = new(TaskCreationOptions.RunContinuationsAsynchronously);
}
private readonly object _lock = new();
private readonly Dictionary<string, InFlightRequest> _requests = new(StringComparer.Ordinal);
public async Task<CoalescedItemsProxyResponse> GetOrFetchAsync(
string key,
Func<CancellationToken, Task<ItemsProxyUpstreamResponse>> fetch,
CancellationToken cancellationToken)
{
InFlightRequest request;
var joinedExistingRequest = true;
lock (_lock)
{
if (_requests.TryGetValue(key, out var existing) && existing is not null)
{
request = existing;
}
else
{
request = new InFlightRequest();
_requests.Add(key, request);
joinedExistingRequest = false;
_ = CompleteAsync(key, request, fetch);
}
}
var response = await request.Completion.Task.WaitAsync(cancellationToken).ConfigureAwait(false);
return new CoalescedItemsProxyResponse(response, joinedExistingRequest);
}
private async Task CompleteAsync(
string key,
InFlightRequest request,
Func<CancellationToken, Task<ItemsProxyUpstreamResponse>> fetch)
{
try
{
request.Completion.TrySetResult(await fetch(CancellationToken.None).ConfigureAwait(false));
}
catch (Exception ex)
{
request.Completion.TrySetException(ex);
}
finally
{
lock (_lock)
{
if (_requests.TryGetValue(key, out var current) && ReferenceEquals(current, request))
_requests.Remove(key);
}
}
}
}