Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
69 changes: 69 additions & 0 deletions GitLfsCache.Tests/Integration/ProxyFlowTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -558,4 +558,73 @@
// 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<HttpResponseMessage> leader = client.GetAsync(href, leaderAbort.Token);
await WaitUntilAsync(() => fixture.Upstream.FetchCount(oid) == 1);
Task<byte[]>[] followers = [.. Enumerable.Range(0, Followers).Select(_ => client.GetByteArrayAsync(href))];

Check warning on line 580 in GitLfsCache.Tests/Integration/ProxyFlowTests.cs

View check run for this annotation

SonarQubeCloud / SonarCloud Code Analysis

Consider using the overload that accepts a CancellationToken and pass 'TestContext.CancellationToken'

See more on https://sonarcloud.io/project/issues?id=ktsu-dev_GitLfsCache&issues=AaEcz7B80YFr0R4GWXVf&open=AaEcz7B80YFr0R4GWXVf&pullRequest=110
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<OperationCanceledException>(() => leader);
fixture.Upstream.ObjectGate.SetResult();

foreach (byte[] body in await Task.WhenAll(followers))
{
CollectionAssert.AreEqual(content, body);

Check warning on line 590 in GitLfsCache.Tests/Integration/ProxyFlowTests.cs

View check run for this annotation

SonarQubeCloud / SonarCloud Code Analysis

Use 'Assert.AreSequenceEqual' instead of 'CollectionAssert.AreEqual'

See more on https://sonarcloud.io/project/issues?id=ktsu-dev_GitLfsCache&issues=AaEcz7B80YFr0R4GWXVe&open=AaEcz7B80YFr0R4GWXVe&pullRequest=110
}

// 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<bool> condition)
{
using CancellationTokenSource timeout = new(TimeSpan.FromSeconds(10));
while (!condition())
{
await Task.Delay(10, timeout.Token);
}
}

/// <summary>Counts <c>gitlfscache.coalesced_waits</c> recorded while the listener is alive.</summary>
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<long>((_, value, _, _) => Interlocked.Add(ref _count, value));
_listener.Start();
}

public long Count => Interlocked.Read(ref _count);

public void Dispose() => _listener.Dispose();
}
}
11 changes: 11 additions & 0 deletions GitLfsCache.Tests/Integration/StubUpstream.cs
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,12 @@ public IReadOnlyList<RecordedRequest> Requests
/// </summary>
public byte[]? CorruptedContent { get; set; }

/// <summary>
/// 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.
/// </summary>
public TaskCompletionSource? ObjectGate { get; set; }

/// <summary>Gets the bytes uploads delivered, keyed by object id.</summary>
public Dictionary<string, byte[]> Uploaded { get; } = new(StringComparer.Ordinal);

Expand Down Expand Up @@ -117,6 +123,11 @@ protected override async Task<HttpResponseMessage> SendAsync(

if (path.Contains("/storage/", StringComparison.Ordinal))
{
if (ObjectGate is TaskCompletionSource gate)
{
await gate.Task.WaitAsync(cancellationToken).ConfigureAwait(false);
}

return BuildObjectResponse(path);
}

Expand Down
90 changes: 59 additions & 31 deletions GitLfsCache/Endpoints/ObjectRouteHandler.cs
Original file line number Diff line number Diff line change
Expand Up @@ -34,7 +34,7 @@
/// <param name="metrics">Cache counters.</param>
/// <param name="options">The configured options.</param>
/// <param name="logger">Logger.</param>
internal sealed class ObjectRouteHandler(

Check warning on line 37 in GitLfsCache/Endpoints/ObjectRouteHandler.cs

View workflow job for this annotation

GitHub Actions / ci / .NET / Analyze & Release

Constructor has 9 parameters, which is greater than the 7 authorized.

Check warning on line 37 in GitLfsCache/Endpoints/ObjectRouteHandler.cs

View workflow job for this annotation

GitHub Actions / ci / .NET / Analyze & Release

Constructor has 9 parameters, which is greater than the 7 authorized.
IUpstreamClient upstreamClient,
IHrefTokenCodec codec,
BatchRewriter rewriter,
Expand All @@ -48,6 +48,12 @@
private const string OctetStream = "application/octet-stream";
private const string TokenQueryParameter = "t";

/// <summary>
/// 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.
/// </summary>
private const int MaxFollowerAttempts = 3;

/// <summary>
/// Answers a Batch API call, rewriting the hrefs upstream returns to point back here.
/// </summary>
Expand Down Expand Up @@ -160,49 +166,71 @@
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++)

Check warning on line 181 in GitLfsCache/Endpoints/ObjectRouteHandler.cs

View workflow job for this annotation

GitHub Actions / ci / .NET / Analyze & Release

This loop's stop condition tests 'ticket' but the incrementer updates 'attempt'.

Check warning on line 181 in GitLfsCache/Endpoints/ObjectRouteHandler.cs

View workflow job for this annotation

GitHub Actions / ci / .NET / Analyze & Release

This loop's stop condition tests 'ticket' but the incrementer updates 'attempt'.

Check failure on line 181 in GitLfsCache/Endpoints/ObjectRouteHandler.cs

View check run for this annotation

SonarQubeCloud / SonarCloud Code Analysis

This loop's stop condition tests 'ticket' but the incrementer updates 'attempt'.

See more on https://sonarcloud.io/project/issues?id=ktsu-dev_GitLfsCache&issues=AaEcz67-0YFr0R4GWXVc&open=AaEcz67-0YFr0R4GWXVc&pullRequest=110
{
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();

Check warning on line 233 in GitLfsCache/Endpoints/ObjectRouteHandler.cs

View workflow job for this annotation

GitHub Actions / ci / .NET / Analyze & Release

Resource 'ticket' has already been disposed explicitly or through a using statement implicitly. Remove the redundant disposal.

Check warning on line 233 in GitLfsCache/Endpoints/ObjectRouteHandler.cs

View workflow job for this annotation

GitHub Actions / ci / .NET / Analyze & Release

Resource 'ticket' has already been disposed explicitly or through a using statement implicitly. Remove the redundant disposal.

Check warning on line 233 in GitLfsCache/Endpoints/ObjectRouteHandler.cs

View check run for this annotation

SonarQubeCloud / SonarCloud Code Analysis

Resource 'ticket' has already been disposed explicitly or through a using statement implicitly. Remove the redundant disposal.

See more on https://sonarcloud.io/project/issues?id=ktsu-dev_GitLfsCache&issues=AaEcz67-0YFr0R4GWXVd&open=AaEcz67-0YFr0R4GWXVd&pullRequest=110
}
}

Expand Down
9 changes: 6 additions & 3 deletions GitLfsCache/Fetching/SingleFlight.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Expand All @@ -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);
}
}
}
Loading