From 926dc21fecb36148dfa41f37594aa8084edf6623 Mon Sep 17 00:00:00 2001 From: Claude Date: Thu, 8 Oct 2026 18:30:49 +0000 Subject: [PATCH] Hand an aborted leader's fetch to one follower instead of all of them [patch] When a leader's client disconnected mid-transfer, its followers were released with a failure and each fetched the object from upstream itself: the aborted fetch plus one per follower. A follower released without the object now acquires again, so one of them becomes the new leader and the rest wait for it, bounded at three leaders before a follower fetches for itself. The single flight retires an entry before releasing its followers, so a follower that acquires again gets a fresh entry rather than the finished one. Fixes ktsu-dev/GitLfsCache#58 Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_017ft9gJVQx6JFNQ4VRoAwd9 --- .../Integration/ProxyFlowTests.cs | 69 ++++++++++++++ GitLfsCache.Tests/Integration/StubUpstream.cs | 11 +++ GitLfsCache/Endpoints/ObjectRouteHandler.cs | 90 ++++++++++++------- GitLfsCache/Fetching/SingleFlight.cs | 9 +- 4 files changed, 145 insertions(+), 34 deletions(-) diff --git a/GitLfsCache.Tests/Integration/ProxyFlowTests.cs b/GitLfsCache.Tests/Integration/ProxyFlowTests.cs index 76c30d8..c4db63c 100644 --- a/GitLfsCache.Tests/Integration/ProxyFlowTests.cs +++ b/GitLfsCache.Tests/Integration/ProxyFlowTests.cs @@ -558,4 +558,73 @@ public async Task ConcurrentMissesForOneObject_ProduceASingleUpstreamFetch() // This is the property the whole cache exists for: eight clients, one upstream transfer. Assert.AreEqual(1, fixture.Upstream.FetchCount(oid)); } + + [TestMethod] + public async Task AbortedLeader_HandsTheFetchToOneFollowerInsteadOfEveryFollowerFetching() + { + const int Followers = 5; + await using ProxyFixture fixture = await ProxyFixture.StartAsync(); + using CoalescedWaitCounter waits = new(); + (byte[] content, string oid) = Object(new string('y', 200_000)); + fixture.Upstream.AddObject(oid, content); + fixture.Upstream.ObjectGate = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + + JsonNode batch = await PostBatchAsync(fixture, "download", oid, content.Length); + string href = Relative(HrefOf(batch, "download")); + using HttpClient client = fixture.Client; + + // The leader's fetch is held at upstream, and the followers queue behind it. + using CancellationTokenSource leaderAbort = new(); + Task leader = client.GetAsync(href, leaderAbort.Token); + await WaitUntilAsync(() => fixture.Upstream.FetchCount(oid) == 1); + Task[] followers = [.. Enumerable.Range(0, Followers).Select(_ => client.GetByteArrayAsync(href))]; + await WaitUntilAsync(() => waits.Count == Followers); + + // The leader's client goes away mid-transfer, then upstream is allowed to answer. + await leaderAbort.CancelAsync(); + await Assert.ThrowsAsync(() => leader); + fixture.Upstream.ObjectGate.SetResult(); + + foreach (byte[] body in await Task.WhenAll(followers)) + { + CollectionAssert.AreEqual(content, body); + } + + // The aborted fetch plus a single replacement, not one fetch per follower. + Assert.AreEqual(2, fixture.Upstream.FetchCount(oid)); + } + + private static async Task WaitUntilAsync(Func condition) + { + using CancellationTokenSource timeout = new(TimeSpan.FromSeconds(10)); + while (!condition()) + { + await Task.Delay(10, timeout.Token); + } + } + + /// Counts gitlfscache.coalesced_waits recorded while the listener is alive. + private sealed class CoalescedWaitCounter : IDisposable + { + private readonly MeterListener _listener = new(); + private long _count; + + public CoalescedWaitCounter() + { + _listener.InstrumentPublished = (instrument, listener) => + { + if (instrument.Meter.Name == CacheMetrics.MeterName && instrument.Name == "gitlfscache.coalesced_waits") + { + listener.EnableMeasurementEvents(instrument); + } + }; + + _listener.SetMeasurementEventCallback((_, value, _, _) => Interlocked.Add(ref _count, value)); + _listener.Start(); + } + + public long Count => Interlocked.Read(ref _count); + + public void Dispose() => _listener.Dispose(); + } } diff --git a/GitLfsCache.Tests/Integration/StubUpstream.cs b/GitLfsCache.Tests/Integration/StubUpstream.cs index 4bcb8e6..ccd37e0 100644 --- a/GitLfsCache.Tests/Integration/StubUpstream.cs +++ b/GitLfsCache.Tests/Integration/StubUpstream.cs @@ -53,6 +53,12 @@ public IReadOnlyList Requests /// public byte[]? CorruptedContent { get; set; } + /// + /// Gets or sets a gate that object fetches wait on before answering, so a test can hold a transfer + /// in flight. A fetch whose request is cancelled while it waits is abandoned. + /// + public TaskCompletionSource? ObjectGate { get; set; } + /// Gets the bytes uploads delivered, keyed by object id. public Dictionary Uploaded { get; } = new(StringComparer.Ordinal); @@ -117,6 +123,11 @@ protected override async Task SendAsync( if (path.Contains("/storage/", StringComparison.Ordinal)) { + if (ObjectGate is TaskCompletionSource gate) + { + await gate.Task.WaitAsync(cancellationToken).ConfigureAwait(false); + } + return BuildObjectResponse(path); } diff --git a/GitLfsCache/Endpoints/ObjectRouteHandler.cs b/GitLfsCache/Endpoints/ObjectRouteHandler.cs index 873d389..d78eb66 100644 --- a/GitLfsCache/Endpoints/ObjectRouteHandler.cs +++ b/GitLfsCache/Endpoints/ObjectRouteHandler.cs @@ -48,6 +48,12 @@ internal sealed class ObjectRouteHandler( private const string OctetStream = "application/octet-stream"; private const string TokenQueryParameter = "t"; + /// + /// How many leaders a follower waits for before it fetches for itself, so a run of failing or + /// stalled leaders cannot keep a request waiting indefinitely. + /// + private const int MaxFollowerAttempts = 3; + /// /// Answers a Batch API call, rewriting the hrefs upstream returns to point back here. /// @@ -160,49 +166,71 @@ await StreamFromUpstreamAsync(context, route, token, range, storeLocally: false, return; } - using IFetchTicket ticket = coalescer.Acquire(route.Upstream, token.Oid); + IFetchTicket ticket = coalescer.Acquire(route.Upstream, token.Oid); - if (!ticket.IsLeader) + try { - metrics.RecordCoalescedWait(route.Upstream); - EndpointLog.WaitingForLeader(logger, token.Oid, route.Upstream); + if (!ticket.IsLeader) + { + metrics.RecordCoalescedWait(route.Upstream); + } - bool published = await ticket - .WaitForLeaderAsync(options.Value.Fetch.FollowerTimeout, cancellationToken) - .ConfigureAwait(false); + // A follower released without the object, because the leader's client went away or the + // leader stalled, queues again: one of the released followers becomes the new leader and + // the rest wait for it, rather than every one of them fetching the same object at once. + for (int attempt = 1; !ticket.IsLeader; attempt++) + { + EndpointLog.WaitingForLeader(logger, token.Oid, route.Upstream); - long nowLength = 0; - Stream? nowCached = published - ? store.OpenRead(route.Upstream, token.Oid, out nowLength) - : null; + bool published = await ticket + .WaitForLeaderAsync(options.Value.Fetch.FollowerTimeout, cancellationToken) + .ConfigureAwait(false); - if (nowCached is not null) - { - await using (nowCached.ConfigureAwait(false)) + long nowLength = 0; + Stream? nowCached = published + ? store.OpenRead(route.Upstream, token.Oid, out nowLength) + : null; + + if (nowCached is not null) { - store.Touch(route.Upstream, token.Oid); - metrics.RecordHit(route.Upstream, nowLength); - await ServeFromStoreAsync(context, nowCached, nowLength, cancellationToken) - .ConfigureAwait(false); + await using (nowCached.ConfigureAwait(false)) + { + store.Touch(route.Upstream, token.Oid); + metrics.RecordHit(route.Upstream, nowLength); + await ServeFromStoreAsync(context, nowCached, nowLength, cancellationToken) + .ConfigureAwait(false); + } + + return; } - return; - } + EndpointLog.LeaderDidNotFinish(logger, token.Oid); - EndpointLog.LeaderDidNotFinish(logger, token.Oid); - } + if (attempt >= MaxFollowerAttempts) + { + break; + } - bool stored = await StreamFromUpstreamAsync( - context, - route, - token, - range: null, - storeLocally: true, - cancellationToken).ConfigureAwait(false); + ticket.Dispose(); + ticket = coalescer.Acquire(route.Upstream, token.Oid); + } - if (ticket.IsLeader) + bool stored = await StreamFromUpstreamAsync( + context, + route, + token, + range: null, + storeLocally: true, + cancellationToken).ConfigureAwait(false); + + if (ticket.IsLeader) + { + ticket.Complete(stored); + } + } + finally { - ticket.Complete(stored); + ticket.Dispose(); } } diff --git a/GitLfsCache/Fetching/SingleFlight.cs b/GitLfsCache/Fetching/SingleFlight.cs index a0ed0f6..11a17a3 100644 --- a/GitLfsCache/Fetching/SingleFlight.cs +++ b/GitLfsCache/Fetching/SingleFlight.cs @@ -103,8 +103,10 @@ public void Complete(bool published) throw new InvalidOperationException("Only the leader reports the outcome of a fetch."); } - entry.Completion.TrySetResult(published); + // Retired first, so a follower that acquires again once released gets a fresh entry + // rather than this finished one. owner.Retire(key, entry); + entry.Completion.TrySetResult(published); } public void Dispose() @@ -122,9 +124,10 @@ public void Dispose() } // A leader that never reported an outcome releases its followers as a failure rather than - // leaving them waiting for a fetch that is no longer happening. - entry.Completion.TrySetResult(false); + // leaving them waiting for a fetch that is no longer happening. Retired first, so a released + // follower that acquires again can become the next leader. owner.Retire(key, entry); + entry.Completion.TrySetResult(false); } } }