From 95ba8440c6a413249930cce349e4ef2c3435d5bb Mon Sep 17 00:00:00 2001 From: Claude Date: Wed, 7 Oct 2026 23:27:18 +0000 Subject: [PATCH] Keep a coalesced clone or fetch running when the request that started it disconnects [patch] The leader ran clone and fetch under its own request's RequestAborted token, so a client hanging up killed git and released every follower as failed. A clone longer than the client timeout therefore never finished, and followers got 503 no-mirror. The work now runs under IHostApplicationLifetime.ApplicationStopping and owns the leader's ticket, while each caller awaits it with its own token. FetchTimeout still bounds it through GitRunner, and host shutdown still cancels it. Fixes ktsu-dev/GitBranchStateCache#44 Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01FM7p8BsZu36YCjPsutDXQN --- .../Fakes/FakeHostApplicationLifetime.cs | 35 ++++ .../Mirrors/MirrorFetcherCoalescingTests.cs | 153 ++++++++++++++++++ .../Mirrors/MirrorFetcherRealGitTests.cs | 2 + .../Mirrors/MirrorMaintenanceServiceTests.cs | 1 + GitBranchStateCache/Mirrors/MirrorFetcher.cs | 64 ++++++-- 5 files changed, 243 insertions(+), 12 deletions(-) create mode 100644 GitBranchStateCache.Tests/Fakes/FakeHostApplicationLifetime.cs create mode 100644 GitBranchStateCache.Tests/Mirrors/MirrorFetcherCoalescingTests.cs diff --git a/GitBranchStateCache.Tests/Fakes/FakeHostApplicationLifetime.cs b/GitBranchStateCache.Tests/Fakes/FakeHostApplicationLifetime.cs new file mode 100644 index 0000000..57494c1 --- /dev/null +++ b/GitBranchStateCache.Tests/Fakes/FakeHostApplicationLifetime.cs @@ -0,0 +1,35 @@ +// Copyright (c) 2023-2026 ktsu-dev contributors + +namespace ktsu.GitBranchStateCache.Tests.Fakes; + +using Microsoft.Extensions.Hosting; + +/// +/// A host lifetime whose shutdown a test triggers by hand. +/// +internal sealed class FakeHostApplicationLifetime : IHostApplicationLifetime, IDisposable +{ + private readonly CancellationTokenSource _started = new(); + private readonly CancellationTokenSource _stopping = new(); + private readonly CancellationTokenSource _stopped = new(); + + /// + public CancellationToken ApplicationStarted => _started.Token; + + /// + public CancellationToken ApplicationStopping => _stopping.Token; + + /// + public CancellationToken ApplicationStopped => _stopped.Token; + + /// + public void StopApplication() => _stopping.Cancel(); + + /// + public void Dispose() + { + _started.Dispose(); + _stopping.Dispose(); + _stopped.Dispose(); + } +} diff --git a/GitBranchStateCache.Tests/Mirrors/MirrorFetcherCoalescingTests.cs b/GitBranchStateCache.Tests/Mirrors/MirrorFetcherCoalescingTests.cs new file mode 100644 index 0000000..a347ed4 --- /dev/null +++ b/GitBranchStateCache.Tests/Mirrors/MirrorFetcherCoalescingTests.cs @@ -0,0 +1,153 @@ +// Copyright (c) 2023-2026 ktsu-dev contributors + +namespace ktsu.GitBranchStateCache.Tests.Mirrors; + +using System.Diagnostics.Metrics; +using ktsu.GitBranchStateCache.Coalescing; +using ktsu.GitBranchStateCache.Configuration; +using ktsu.GitBranchStateCache.Mirrors; +using ktsu.GitBranchStateCache.Observability; +using ktsu.GitBranchStateCache.Tests.Fakes; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Logging.Abstractions; +using Microsoft.Extensions.Options; +using Microsoft.Extensions.Time.Testing; +using Microsoft.VisualStudio.TestTools.UnitTesting; +using Testably.Abstractions.Testing; + +/// +/// The clone or fetch one request starts is shared by every request coalesced onto it, so it must not +/// belong to whichever request happened to arrive first. +/// +[TestClass] +public class MirrorFetcherCoalescingTests +{ + /// Gets or sets the context MSTest supplies, used for the run's cancellation token. + public TestContext TestContext { get; set; } = null!; + + private static readonly string Root = Path.Combine( + Path.GetPathRoot(Path.GetTempPath()) ?? Path.DirectorySeparatorChar.ToString(), + "gitbranchstatecache-coalescing"); + + private static readonly MirrorKey Key = new("github", "studio/game.git"); + + private sealed record Harness( + MirrorFetcher Fetcher, + FakeGitRunner Runner, + FakeHostApplicationLifetime Lifetime, + string Directory, + TaskCompletionSource CloneStarted, + TaskCompletionSource CloneMayFinish); + + /// + /// Builds a fetcher whose clone blocks until the test releases it, and which leaves a mirror on the + /// mock filesystem as a real clone would. + /// + private static Harness Build() + { + MockFileSystem fileSystem = new(); + fileSystem.Directory.CreateDirectory(Root); + + FakeTimeProvider time = new(new DateTimeOffset(2026, 8, 19, 9, 47, 0, TimeSpan.Zero)); + IOptions options = Options.Create(new GitBranchStateCacheOptions + { + MirrorRoot = Root, + }); + + MirrorStore store = new(fileSystem, options, time); + Assert.IsTrue(store.TryResolve(Key, out string? directory)); + + TaskCompletionSource cloneStarted = new(TaskCreationOptions.RunContinuationsAsynchronously); + TaskCompletionSource cloneMayFinish = new(TaskCreationOptions.RunContinuationsAsynchronously); + FakeGitRunner runner = BlockingCloneRunner(fileSystem, cloneStarted, cloneMayFinish); + + FakeHostApplicationLifetime lifetime = new(); + + MirrorFetcher fetcher = new( + runner, + store, + fileSystem, + new SingleFlight(), + new BranchStateMetrics(MeterFactory()), + options, + time, + lifetime, + NullLogger.Instance); + + return new Harness(fetcher, runner, lifetime, directory!, cloneStarted, cloneMayFinish); + } + + private static FakeGitRunner BlockingCloneRunner( + MockFileSystem fileSystem, + TaskCompletionSource cloneStarted, + TaskCompletionSource cloneMayFinish) => + new() + { + Before = async invocation => + { + if (invocation.Arguments[0] != "clone") + { + return; + } + + cloneStarted.TrySetResult(); + await cloneMayFinish.Task.ConfigureAwait(false); + fileSystem.Directory.CreateDirectory(fileSystem.Path.Combine(invocation.Arguments[^1], "objects")); + }, + }; + + private static IMeterFactory MeterFactory() + { + ServiceCollection services = new(); + services.AddMetrics(); + return services.BuildServiceProvider().GetRequiredService(); + } + + private static Task EnsureAsync(Harness harness, CancellationToken cancellationToken) => + harness.Fetcher.EnsureCurrentAsync( + Key, + harness.Directory, + new Uri("https://forge.example/studio/game.git"), + new Uri("https://forge.example"), + authorization: null, + cancellationToken); + + [TestMethod] + public async Task EnsureCurrent_WhenTheLeaderDisconnectsMidClone_FinishesTheCloneForTheFollowerAsync() + { + Harness harness = Build(); + using FakeHostApplicationLifetime lifetime = harness.Lifetime; + + using CancellationTokenSource leaderDisconnect = new(); + Task leader = EnsureAsync(harness, leaderDisconnect.Token); + await harness.CloneStarted.Task.WaitAsync(TestContext.CancellationTokenSource.Token).ConfigureAwait(false); + + Task follower = EnsureAsync(harness, TestContext.CancellationTokenSource.Token); + + await leaderDisconnect.CancelAsync().ConfigureAwait(false); + await Assert.ThrowsAsync(() => leader).ConfigureAwait(false); + + harness.CloneMayFinish.TrySetResult(); + MirrorFetchResult result = await follower.ConfigureAwait(false); + + Assert.AreEqual(MirrorFetchStatus.Current, result.Status, result.Failure); + Assert.AreEqual(1, harness.Runner.CountOf("clone")); + } + + [TestMethod] + public async Task EnsureCurrent_WhenTheHostStopsMidClone_CancelsTheCloneAsync() + { + Harness harness = Build(); + using FakeHostApplicationLifetime lifetime = harness.Lifetime; + + Task leader = EnsureAsync(harness, TestContext.CancellationTokenSource.Token); + await harness.CloneStarted.Task.WaitAsync(TestContext.CancellationTokenSource.Token).ConfigureAwait(false); + + lifetime.StopApplication(); + harness.CloneMayFinish.TrySetResult(); + + // The fake runner throws once its token is cancelled, which is where GitRunner kills the process. + await Assert.ThrowsAsync(() => leader).ConfigureAwait(false); + Assert.AreEqual(0, harness.Runner.CountOf("config")); + } +} diff --git a/GitBranchStateCache.Tests/Mirrors/MirrorFetcherRealGitTests.cs b/GitBranchStateCache.Tests/Mirrors/MirrorFetcherRealGitTests.cs index 32a24d0..be03f18 100644 --- a/GitBranchStateCache.Tests/Mirrors/MirrorFetcherRealGitTests.cs +++ b/GitBranchStateCache.Tests/Mirrors/MirrorFetcherRealGitTests.cs @@ -13,6 +13,7 @@ namespace ktsu.GitBranchStateCache.Tests.Mirrors; using ktsu.GitBranchStateCache.Mirrors; using ktsu.GitBranchStateCache.Observability; using ktsu.GitBranchStateCache.Refs; +using ktsu.GitBranchStateCache.Tests.Fakes; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Logging.Abstractions; using Microsoft.Extensions.Options; @@ -100,6 +101,7 @@ public void RemoveFixture() new BranchStateMetrics(meterFactory), options, time, + new FakeHostApplicationLifetime(), NullLogger.Instance); Assert.IsTrue(store.TryResolve(new MirrorKey("github", "studio/game.git"), out string? directory)); diff --git a/GitBranchStateCache.Tests/Mirrors/MirrorMaintenanceServiceTests.cs b/GitBranchStateCache.Tests/Mirrors/MirrorMaintenanceServiceTests.cs index 282dd8a..53d144e 100644 --- a/GitBranchStateCache.Tests/Mirrors/MirrorMaintenanceServiceTests.cs +++ b/GitBranchStateCache.Tests/Mirrors/MirrorMaintenanceServiceTests.cs @@ -81,6 +81,7 @@ private static MirrorFetcher BuildFetcher( new BranchStateMetrics(meterFactory), options, time, + new FakeHostApplicationLifetime(), NullLogger.Instance); } diff --git a/GitBranchStateCache/Mirrors/MirrorFetcher.cs b/GitBranchStateCache/Mirrors/MirrorFetcher.cs index eba5b69..59babe6 100644 --- a/GitBranchStateCache/Mirrors/MirrorFetcher.cs +++ b/GitBranchStateCache/Mirrors/MirrorFetcher.cs @@ -7,6 +7,7 @@ namespace ktsu.GitBranchStateCache.Mirrors; using ktsu.GitBranchStateCache.Configuration; using ktsu.GitBranchStateCache.Git; using ktsu.GitBranchStateCache.Observability; +using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Logging; using Microsoft.Extensions.Options; @@ -20,6 +21,13 @@ namespace ktsu.GitBranchStateCache.Mirrors; /// LFS the large assets are pointer files of a hundred or so bytes each, so the git object store was /// never carrying the bulk anyway. /// +/// A clone or fetch is shared by every request coalesced onto it, so it runs under the service's own +/// lifetime rather than under the request that happened to start it. Each request waits on it with its +/// own disconnect token: a client that hangs up stops waiting, but the work carries on for everyone +/// else. Host shutdown still cancels it, and +/// still bounds it, because enforces the invocation's timeout itself. +/// +/// /// A clone lands in a temporary directory and is moved into place only once it has finished. A crash /// part way through a clone of a large repository would otherwise leave a directory that looks like a /// mirror, and every later request would be answered from a repository missing most of its history. @@ -32,6 +40,7 @@ namespace ktsu.GitBranchStateCache.Mirrors; /// Service counters. /// The configured options. /// Clock, injected so freshness is testable. +/// The host's lifetime, whose shutdown cancels a clone or fetch in progress. /// Logger. public sealed class MirrorFetcher( IGitRunner runner, @@ -41,6 +50,7 @@ public sealed class MirrorFetcher( BranchStateMetrics metrics, IOptions options, TimeProvider timeProvider, + IHostApplicationLifetime lifetime, ILogger logger) : IMirrorFetcher { /// @@ -65,25 +75,51 @@ public async Task EnsureCurrentAsync( return MirrorFetchResult.Current(current); } - using IWorkTicket ticket = flights.Acquire(key.ToFlightKey()); + IWorkTicket ticket = flights.Acquire(key.ToFlightKey()); if (!ticket.IsLeader) { return await FollowAsync(key, directory, ticket, cancellationToken).ConfigureAwait(false); } - MirrorFetchResult result = exists - ? await FetchAsync(key, directory, upstreamBase, authorization, fetchedAt, cancellationToken) - .ConfigureAwait(false) - : await CloneAsync(key, directory, repositoryUrl, upstreamBase, authorization, cancellationToken) - .ConfigureAwait(false); + // The work owns the leader's ticket from here on, so a leader whose client disconnects stops waiting + // without abandoning the clone or fetch every follower is waiting on. + Task work = LeadAsync( + key, directory, repositoryUrl, upstreamBase, authorization, exists, fetchedAt, ticket); - ticket.Complete(result.Status == MirrorFetchStatus.Current); - return result; + return await work.WaitAsync(cancellationToken).ConfigureAwait(false); } /// - /// Waits for whichever request is already working on this repository. + /// Does the clone or fetch on behalf of every request coalesced onto it, and reports the outcome. + /// + private async Task LeadAsync( + MirrorKey key, + string directory, + Uri repositoryUrl, + Uri upstreamBase, + string? authorization, + bool exists, + DateTimeOffset? fetchedAt, + IWorkTicket ticket) + { + using (ticket) + { + CancellationToken stopping = lifetime.ApplicationStopping; + + MirrorFetchResult result = exists + ? await FetchAsync(key, directory, upstreamBase, authorization, fetchedAt, stopping) + .ConfigureAwait(false) + : await CloneAsync(key, directory, repositoryUrl, upstreamBase, authorization, stopping) + .ConfigureAwait(false); + + ticket.Complete(result.Status == MirrorFetchStatus.Current); + return result; + } + } + + /// + /// Waits for whichever request is already working on this repository, and releases its ticket. /// /// /// A follower whose leader succeeded re-reads the marker rather than trusting the leader's answer, @@ -99,9 +135,13 @@ private async Task FollowAsync( { metrics.RecordFetchWait(key.Upstream); - bool succeeded = await ticket - .WaitForLeaderAsync(options.Value.FetchTimeout, cancellationToken) - .ConfigureAwait(false); + bool succeeded; + using (ticket) + { + succeeded = await ticket + .WaitForLeaderAsync(options.Value.FetchTimeout, cancellationToken) + .ConfigureAwait(false); + } if (succeeded && mirrors.RefsFetchedAt(directory) is DateTimeOffset refreshed) {