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) {