From 17cbf7d7b2559f5ff98167dde9c10d428d625aaf Mon Sep 17 00:00:00 2001 From: "Emil H. Clausen" Date: Wed, 29 Jul 2026 16:09:58 +0200 Subject: [PATCH] refactor(catalog): extract fetcher and caching into dedicated classes with stale-while-revalidate strategy --- .../TestCatalogApplicationService.cs | 139 +++++++++++- .../Application/CatalogApplicationService.cs | 197 +----------------- src/SelfService/Application/CatalogFetcher.cs | 178 ++++++++++++++++ .../Application/CatalogSnapshotCache.cs | 112 ++++++++++ src/SelfService/Configuration/Domain.cs | 2 + 5 files changed, 430 insertions(+), 198 deletions(-) create mode 100644 src/SelfService/Application/CatalogFetcher.cs create mode 100644 src/SelfService/Application/CatalogSnapshotCache.cs diff --git a/src/SelfService.Tests/Application/TestCatalogApplicationService.cs b/src/SelfService.Tests/Application/TestCatalogApplicationService.cs index 9038403f..40a7c5aa 100644 --- a/src/SelfService.Tests/Application/TestCatalogApplicationService.cs +++ b/src/SelfService.Tests/Application/TestCatalogApplicationService.cs @@ -1,4 +1,5 @@ -using Microsoft.Extensions.Caching.Memory; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Logging; using Microsoft.Extensions.Logging.Abstractions; using Moq; using SelfService.Application; @@ -13,15 +14,31 @@ public class TestCatalogApplicationService private static CatalogApplicationService BuildService( CatalogConfig config, ICatalogClient catalogClient, - ICapabilityRepository capabilityRepository + ICapabilityRepository capabilityRepository, + TimeSpan? ttl = null + ) => new(BuildCache(config, catalogClient, capabilityRepository, ttl)); + + /// The cache resolves the fetcher from a scope of its own, so tests wire a small + /// real container around the mocks rather than handing it a fetcher directly. + private static CatalogSnapshotCache BuildCache( + CatalogConfig config, + ICatalogClient catalogClient, + ICapabilityRepository capabilityRepository, + TimeSpan? ttl = null ) { - return new CatalogApplicationService( - config, - catalogClient, - capabilityRepository, - new MemoryCache(new MemoryCacheOptions()), - NullLogger.Instance + var services = new ServiceCollection(); + services.AddSingleton(config); + services.AddSingleton(catalogClient); + services.AddSingleton(capabilityRepository); + services.AddSingleton>(NullLogger.Instance); + services.AddTransient(); + var provider = services.BuildServiceProvider(); + + return new CatalogSnapshotCache( + provider.GetRequiredService(), + NullLogger.Instance, + ttl ?? CatalogSnapshotCache.DefaultTtl ); } @@ -267,6 +284,112 @@ public async Task All_clusters_fail_reports_unavailable_with_no_items() Assert.Equal(1, result.Availability.ClustersFailed); } + [Fact] + public async Task Expired_snapshot_is_served_without_waiting_for_the_refresh() + { + var capability = A.Capability.WithId(CapabilityId.Parse("team-alpha-abcde")).WithName("Team Alpha").Build(); + var refreshReached = new TaskCompletionSource(); + var releaseRefresh = new TaskCompletionSource(); + var calls = 0; + + var catalogClient = new Mock(); + catalogClient + .Setup(x => x.GetCatalog(It.IsAny(), It.IsAny())) + .Returns( + async (Uri _, CancellationToken _) => + { + // Every fetch after the first one hangs, standing in for a slow fan-out. + if (Interlocked.Increment(ref calls) > 1) + { + refreshReached.TrySetResult(); + await releaseRefresh.Task; + } + return SingleAppSnapshot(); + } + ); + + var capabilityRepository = new Mock(); + capabilityRepository + .Setup(x => x.GetByIds(It.IsAny>())) + .ReturnsAsync(new[] { capability }); + + // Zero TTL: the snapshot is stale the moment it lands. + var service = BuildService( + SingleCluster(), + catalogClient.Object, + capabilityRepository.Object, + ttl: TimeSpan.Zero + ); + + await service.ListApplications(new ApplicationFilters()); // cold start — this one does wait + + // The snapshot is now expired, so this call kicks off a refresh that will not + // complete. It must still return, from the stale data, rather than block on it. + var stale = await service.ListApplications(new ApplicationFilters()).WaitAsync(TimeSpan.FromSeconds(10)); + + Assert.Equal("api", Assert.Single(stale.Items).Name); + await refreshReached.Task.WaitAsync(TimeSpan.FromSeconds(10)); // the refresh did start + releaseRefresh.SetResult(); + } + + [Fact] + public async Task Cold_callers_share_a_single_fetch() + { + var capability = A.Capability.WithId(CapabilityId.Parse("team-alpha-abcde")).WithName("Team Alpha").Build(); + var firstCallReached = new TaskCompletionSource(); + var releaseFirstCall = new TaskCompletionSource(); + var calls = 0; + + var catalogClient = new Mock(); + catalogClient + .Setup(x => x.GetCatalog(It.IsAny(), It.IsAny())) + .Returns( + async (Uri _, CancellationToken _) => + { + if (Interlocked.Increment(ref calls) == 1) + { + firstCallReached.TrySetResult(); + await releaseFirstCall.Task; + } + return SingleAppSnapshot(); + } + ); + + var capabilityRepository = new Mock(); + capabilityRepository + .Setup(x => x.GetByIds(It.IsAny>())) + .ReturnsAsync(new[] { capability }); + + var service = BuildService(SingleCluster(), catalogClient.Object, capabilityRepository.Object); + + var first = service.ListApplications(new ApplicationFilters()); + await firstCallReached.Task.WaitAsync(TimeSpan.FromSeconds(10)); // a fetch is in flight + var second = service.ListApplications(new ApplicationFilters()); // arrives mid-fetch + releaseFirstCall.SetResult(); + + await Task.WhenAll(first, second).WaitAsync(TimeSpan.FromSeconds(10)); + + // The second caller queued on the refresh gate and read the snapshot the first + // one produced — an expiry must not fan out once per concurrent request. + Assert.Equal(1, calls); + Assert.Equal("api", Assert.Single((await second).Items).Name); + } + + private static CatalogSnapshotDto SingleAppSnapshot() => + new() + { + Applications = + { + new ApplicationEntryDto + { + Namespace = "team-alpha-abcde", + Name = "api", + Kind = "Deployment", + CapabilityId = "team-alpha-abcde", + }, + }, + }; + [Fact] public async Task TokenProvider_unconfigured_scope_returns_null_without_acquiring() { diff --git a/src/SelfService/Application/CatalogApplicationService.cs b/src/SelfService/Application/CatalogApplicationService.cs index 5fa63e7e..39e0ca17 100644 --- a/src/SelfService/Application/CatalogApplicationService.cs +++ b/src/SelfService/Application/CatalogApplicationService.cs @@ -1,4 +1,3 @@ -using Microsoft.Extensions.Caching.Memory; using SelfService.Domain.Models; using SelfService.Infrastructure.Catalog; @@ -41,34 +40,13 @@ Task> GetDependencies( ); } -/// Caching proxy over the per-cluster ssu-catalog services. On cache miss it does a -/// full-snapshot fetch per cluster, concatenates, then once joins on capabilityId against Capability data (filtering to capability-owned apps and attaching the name). All -/// query methods read from the single cached merged structure. Only stored in-memory public class CatalogApplicationService : ICatalogApplicationService { - private const string CacheKey = "catalog:merged"; - private static readonly TimeSpan CacheTtl = TimeSpan.FromSeconds(45); - private static readonly TimeSpan FetchTimeout = TimeSpan.FromSeconds(20); + private readonly CatalogSnapshotCache _snapshots; - private readonly CatalogConfig _config; - private readonly ICatalogClient _catalogClient; - private readonly ICapabilityRepository _capabilityRepository; - private readonly IMemoryCache _cache; - private readonly ILogger _logger; - - public CatalogApplicationService( - CatalogConfig config, - ICatalogClient catalogClient, - ICapabilityRepository capabilityRepository, - IMemoryCache cache, - ILogger logger - ) + public CatalogApplicationService(CatalogSnapshotCache snapshots) { - _config = config; - _catalogClient = catalogClient; - _capabilityRepository = capabilityRepository; - _cache = cache; - _logger = logger; + _snapshots = snapshots; } public async Task> GetDeploymentsForCapability( @@ -76,7 +54,7 @@ public async Task> GetDeploymentsForCapabilit CancellationToken cancellationToken = default ) { - var merged = await GetMerged(); + var merged = await _snapshots.Get(cancellationToken); var id = capabilityId.ToString(); var items = merged .Applications.Where(a => string.Equals(a.CapabilityId, id, StringComparison.OrdinalIgnoreCase)) @@ -89,7 +67,7 @@ public async Task> ListApplications( CancellationToken cancellationToken = default ) { - var merged = await GetMerged(); + var merged = await _snapshots.Get(cancellationToken); IEnumerable apps = merged.Applications; if (!string.IsNullOrWhiteSpace(filters.CapabilityId)) @@ -120,7 +98,7 @@ public async Task> ListApplications( public async Task> ListNamespaces(CancellationToken cancellationToken = default) { - var merged = await GetMerged(); + var merged = await _snapshots.Get(cancellationToken); return new CatalogResult(merged.Namespaces, merged.Availability); } @@ -129,7 +107,7 @@ public async Task> GetDependencies( CancellationToken cancellationToken = default ) { - var merged = await GetMerged(); + var merged = await _snapshots.Get(cancellationToken); IEnumerable deps = merged.Dependencies; if (!string.IsNullOrWhiteSpace(filters.Namespace)) @@ -148,165 +126,4 @@ public async Task> GetDependencies( } private static bool HasDocs(ApplicationEntryDto app) => app.Services.Any(s => s.ApiDocs.Count > 0); - - private Task GetMerged() - { - return _cache.GetOrCreateAsync( - CacheKey, - async entry => - { - entry.AbsoluteExpirationRelativeToNow = CacheTtl; - using var cts = new CancellationTokenSource(FetchTimeout); - return await FetchAndMerge(cts.Token); - } - )!; - } - - private async Task FetchAndMerge(CancellationToken cancellationToken) - { - var registry = _config.Clusters; - - var fetches = registry - .Select(async endpoint => - { - var snapshot = await _catalogClient.GetCatalog(endpoint.Url, cancellationToken); - return (endpoint, snapshot); - }) - .ToList(); - - var results = await Task.WhenAll(fetches); - - var apps = new List(); - var namespaces = new List(); - var dependencies = new List(); - var clustersFailed = 0; - var collectedTimes = new List(); - var publishedTimes = new List(); - - foreach (var (endpoint, snapshot) in results) - { - if (snapshot is null) - { - clustersFailed++; - continue; - } - - foreach (var app in snapshot.Applications) - { - app.Cluster = endpoint.Cluster; - } - foreach (var ns in snapshot.Namespaces) - { - ns.Cluster = endpoint.Cluster; - } - - apps.AddRange(snapshot.Applications); - namespaces.AddRange(snapshot.Namespaces); - dependencies.AddRange(snapshot.Dependencies); - if (snapshot.CollectedAt.Year > 1) - { - collectedTimes.Add(snapshot.CollectedAt); - } - if (snapshot.PublishedAt.Year > 1) - { - publishedTimes.Add(snapshot.PublishedAt); - } - } - - var availability = new CatalogAvailability( - CatalogAvailable: registry.Count > 0 && clustersFailed < registry.Count, - ClustersQueried: registry.Count, - ClustersFailed: clustersFailed, - CollectedAt: collectedTimes.Count > 0 ? collectedTimes.Min() : null, - PublishedAt: publishedTimes.Count > 0 ? publishedTimes.Min() : null - ); - - var capabilityNames = await ResolveCapabilityNames(apps, namespaces); - - var ownedApps = apps.Where(a => - TryAttachCapability(a.CapabilityId, capabilityNames, name => a.CapabilityName = name) - ) - .ToList(); - var ownedNamespaces = namespaces - .Where(n => TryAttachCapability(n.CapabilityId, capabilityNames, name => n.CapabilityName = name)) - .ToList(); - - var ownedNamespaceKeys = ownedApps - .Select(a => (a.Cluster, a.Namespace)) - .Concat(ownedNamespaces.Select(n => (n.Cluster, n.Name))) - .ToHashSet(); - var ownedDependencies = dependencies - .Where(d => - ownedNamespaceKeys.Contains((d.Source.Cluster, d.Source.Namespace)) - || ownedNamespaceKeys.Contains((d.Target.Cluster, d.Target.Namespace)) - ) - .ToList(); - - _logger.LogDebug( - "Catalog merged: {Apps} owned apps, {Namespaces} owned namespaces, {Deps} dependencies across {Queried} clusters ({Failed} failed)", - ownedApps.Count, - ownedNamespaces.Count, - ownedDependencies.Count, - availability.ClustersQueried, - availability.ClustersFailed - ); - - return new MergedCatalog(ownedApps, ownedNamespaces, ownedDependencies, availability); - } - - private async Task> ResolveCapabilityNames( - IEnumerable apps, - IEnumerable namespaces - ) - { - var candidateIds = apps.Select(a => a.CapabilityId) - .Concat(namespaces.Select(n => n.CapabilityId)) - .Where(id => !string.IsNullOrWhiteSpace(id)) - .Distinct(StringComparer.OrdinalIgnoreCase) - .ToList(); - - var parsed = new List(); - foreach (var id in candidateIds) - { - if (CapabilityId.TryParse(id, out var capabilityId)) - { - parsed.Add(capabilityId); - } - } - - if (parsed.Count == 0) - { - return new Dictionary(); - } - - var capabilities = await _capabilityRepository.GetByIds(parsed); - var names = new Dictionary(StringComparer.OrdinalIgnoreCase); - foreach (var capability in capabilities) - { - names[capability.Id.ToString()] = capability.Name; - } - return names; - } - - private static bool TryAttachCapability( - string capabilityId, - IReadOnlyDictionary capabilityNames, - Action attachName - ) - { - if (string.IsNullOrWhiteSpace(capabilityId) || !capabilityNames.TryGetValue(capabilityId, out var name)) - { - return false; - } - attachName(name); - return true; - } - - /// The cached, merged + capability-joined catalog. All query methods read from this. - private sealed record MergedCatalog( - IReadOnlyList Applications, - IReadOnlyList Namespaces, - IReadOnlyList Dependencies, - CatalogAvailability Availability - ); } diff --git a/src/SelfService/Application/CatalogFetcher.cs b/src/SelfService/Application/CatalogFetcher.cs new file mode 100644 index 00000000..0ee2edbe --- /dev/null +++ b/src/SelfService/Application/CatalogFetcher.cs @@ -0,0 +1,178 @@ +using SelfService.Domain.Models; +using SelfService.Infrastructure.Catalog; + +namespace SelfService.Application; + +/// The merged, capability-joined catalog. All catalog query methods read from this. +public sealed record MergedCatalog( + IReadOnlyList Applications, + IReadOnlyList Namespaces, + IReadOnlyList Dependencies, + CatalogAvailability Availability +); + +public interface ICatalogFetcher +{ + Task FetchAndMerge(CancellationToken cancellationToken); +} + +public class CatalogFetcher : ICatalogFetcher +{ + private readonly CatalogConfig _config; + private readonly ICatalogClient _catalogClient; + private readonly ICapabilityRepository _capabilityRepository; + private readonly ILogger _logger; + + public CatalogFetcher( + CatalogConfig config, + ICatalogClient catalogClient, + ICapabilityRepository capabilityRepository, + ILogger logger + ) + { + _config = config; + _catalogClient = catalogClient; + _capabilityRepository = capabilityRepository; + _logger = logger; + } + + public async Task FetchAndMerge(CancellationToken cancellationToken) + { + var registry = _config.Clusters; + + var fetches = registry + .Select(async endpoint => + { + var snapshot = await _catalogClient.GetCatalog(endpoint.Url, cancellationToken); + return (endpoint, snapshot); + }) + .ToList(); + + var results = await Task.WhenAll(fetches); + + var apps = new List(); + var namespaces = new List(); + var dependencies = new List(); + var clustersFailed = 0; + var collectedTimes = new List(); + var publishedTimes = new List(); + + foreach (var (endpoint, snapshot) in results) + { + if (snapshot is null) + { + clustersFailed++; + continue; + } + + foreach (var app in snapshot.Applications) + { + app.Cluster = endpoint.Cluster; + } + foreach (var ns in snapshot.Namespaces) + { + ns.Cluster = endpoint.Cluster; + } + + apps.AddRange(snapshot.Applications); + namespaces.AddRange(snapshot.Namespaces); + dependencies.AddRange(snapshot.Dependencies); + if (snapshot.CollectedAt.Year > 1) + { + collectedTimes.Add(snapshot.CollectedAt); + } + if (snapshot.PublishedAt.Year > 1) + { + publishedTimes.Add(snapshot.PublishedAt); + } + } + + var availability = new CatalogAvailability( + CatalogAvailable: registry.Count > 0 && clustersFailed < registry.Count, + ClustersQueried: registry.Count, + ClustersFailed: clustersFailed, + CollectedAt: collectedTimes.Count > 0 ? collectedTimes.Min() : null, + PublishedAt: publishedTimes.Count > 0 ? publishedTimes.Min() : null + ); + + var capabilityNames = await ResolveCapabilityNames(apps, namespaces); + + var ownedApps = apps.Where(a => + TryAttachCapability(a.CapabilityId, capabilityNames, name => a.CapabilityName = name) + ) + .ToList(); + var ownedNamespaces = namespaces + .Where(n => TryAttachCapability(n.CapabilityId, capabilityNames, name => n.CapabilityName = name)) + .ToList(); + + var ownedNamespaceKeys = ownedApps + .Select(a => (a.Cluster, a.Namespace)) + .Concat(ownedNamespaces.Select(n => (n.Cluster, n.Name))) + .ToHashSet(); + var ownedDependencies = dependencies + .Where(d => + ownedNamespaceKeys.Contains((d.Source.Cluster, d.Source.Namespace)) + || ownedNamespaceKeys.Contains((d.Target.Cluster, d.Target.Namespace)) + ) + .ToList(); + + _logger.LogDebug( + "Catalog merged: {Apps} owned apps, {Namespaces} owned namespaces, {Deps} dependencies across {Queried} clusters ({Failed} failed)", + ownedApps.Count, + ownedNamespaces.Count, + ownedDependencies.Count, + availability.ClustersQueried, + availability.ClustersFailed + ); + + return new MergedCatalog(ownedApps, ownedNamespaces, ownedDependencies, availability); + } + + private async Task> ResolveCapabilityNames( + IEnumerable apps, + IEnumerable namespaces + ) + { + var candidateIds = apps.Select(a => a.CapabilityId) + .Concat(namespaces.Select(n => n.CapabilityId)) + .Where(id => !string.IsNullOrWhiteSpace(id)) + .Distinct(StringComparer.OrdinalIgnoreCase) + .ToList(); + + var parsed = new List(); + foreach (var id in candidateIds) + { + if (CapabilityId.TryParse(id, out var capabilityId)) + { + parsed.Add(capabilityId); + } + } + + if (parsed.Count == 0) + { + return new Dictionary(); + } + + var capabilities = await _capabilityRepository.GetByIds(parsed); + var names = new Dictionary(StringComparer.OrdinalIgnoreCase); + foreach (var capability in capabilities) + { + names[capability.Id.ToString()] = capability.Name; + } + return names; + } + + private static bool TryAttachCapability( + string capabilityId, + IReadOnlyDictionary capabilityNames, + Action attachName + ) + { + if (string.IsNullOrWhiteSpace(capabilityId) || !capabilityNames.TryGetValue(capabilityId, out var name)) + { + return false; + } + attachName(name); + return true; + } +} diff --git a/src/SelfService/Application/CatalogSnapshotCache.cs b/src/SelfService/Application/CatalogSnapshotCache.cs new file mode 100644 index 00000000..7fe35629 --- /dev/null +++ b/src/SelfService/Application/CatalogSnapshotCache.cs @@ -0,0 +1,112 @@ +namespace SelfService.Application; + +public sealed class CatalogSnapshotCache +{ + public static readonly TimeSpan DefaultTtl = TimeSpan.FromSeconds(45); + private static readonly TimeSpan FetchTimeout = TimeSpan.FromSeconds(20); + + private readonly IServiceScopeFactory _scopeFactory; + private readonly ILogger _logger; + private readonly TimeSpan _ttl; + private readonly SemaphoreSlim _refreshGate = new(1, 1); + + private volatile Entry? _entry; + + public CatalogSnapshotCache(IServiceScopeFactory scopeFactory, ILogger logger) + : this(scopeFactory, logger, DefaultTtl) { } + + /// TTL overload — for tests that need to age a snapshot out without waiting. + public CatalogSnapshotCache(IServiceScopeFactory scopeFactory, ILogger logger, TimeSpan ttl) + { + _scopeFactory = scopeFactory; + _logger = logger; + _ttl = ttl; + } + + public async Task Get(CancellationToken cancellationToken = default) + { + var entry = _entry; + + // Cold start — nothing to serve, so this one caller has to wait. + if (entry is null) + { + return (await RefreshAndWait(cancellationToken)).Catalog; + } + + if (!IsFresh(entry)) + { + RefreshInBackground(); + } + + return entry.Catalog; + } + + private async Task RefreshAndWait(CancellationToken cancellationToken) + { + await _refreshGate.WaitAsync(cancellationToken); + try + { + // Another caller may have populated the cache while we queued on the gate. + var existing = _entry; + if (existing is not null && IsFresh(existing)) + { + return existing; + } + return await Fetch(); + } + finally + { + _refreshGate.Release(); + } + } + + private void RefreshInBackground() + { + // A refresh is already in flight — the stale snapshot stays in place until + // it finishes. Nothing to do but let this request read the old data. + if (!_refreshGate.Wait(0)) + { + return; + } + + _ = Task.Run(async () => + { + try + { + await Fetch(); + } + catch (Exception e) + { + // The previous snapshot is still there and still being served, so a + // failed refresh degrades to stale data rather than a failed request. + _logger.LogWarning( + e, + "Background refresh of the catalog snapshot failed; continuing to serve stale data" + ); + } + finally + { + _refreshGate.Release(); + } + }); + } + + /// Always called with held. + private async Task Fetch() + { + using var scope = _scopeFactory.CreateScope(); + var fetcher = scope.ServiceProvider.GetRequiredService(); + + // Not tied to the triggering request: a background refresh outlives it, and a + // cold-start caller that gives up should not abort the fetch for everyone else. + using var cts = new CancellationTokenSource(FetchTimeout); + + var entry = new Entry(await fetcher.FetchAndMerge(cts.Token), DateTime.UtcNow); + _entry = entry; + return entry; + } + + private bool IsFresh(Entry entry) => DateTime.UtcNow - entry.RefreshedAt < _ttl; + + private sealed record Entry(MergedCatalog Catalog, DateTime RefreshedAt); +} diff --git a/src/SelfService/Configuration/Domain.cs b/src/SelfService/Configuration/Domain.cs index c26e57d2..2e495d01 100644 --- a/src/SelfService/Configuration/Domain.cs +++ b/src/SelfService/Configuration/Domain.cs @@ -197,6 +197,8 @@ public static void AddDomain(this WebApplicationBuilder builder) builder.Services.AddMemoryCache(); builder.Services.AddScoped(); builder.Services.AddHttpClient(); + builder.Services.AddTransient(); + builder.Services.AddSingleton(); builder.Services.AddTransient(); } }