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
80 changes: 80 additions & 0 deletions GitLfsCache.Tests/Integration/ProxyFlowTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -2,11 +2,13 @@

namespace ktsu.GitLfsCache.Tests.Integration;

using System.Diagnostics.Metrics;
using System.Net;
using System.Net.Http.Json;
using System.Security.Cryptography;
using System.Text;
using System.Text.Json.Nodes;
using ktsu.GitLfsCache.Observability;
using Microsoft.VisualStudio.TestTools.UnitTesting;

[TestClass]
Expand Down Expand Up @@ -308,6 +310,84 @@
Assert.IsTrue(fixture.Store.Exists("github", oid), "A pushed object should be cached for the next fetch.");
}

/// <summary>
/// Puts a file where the upstream's staging directory belongs, so no staging file can be opened.
/// </summary>
/// <remarks>
/// Stands in for every way the store can stop accepting writes after the startup probe passed:
/// a permissions change, a read-only remount, inode exhaustion.
/// </remarks>
private static void BlockStaging(ProxyFixture fixture)
{
string upstreamDirectory = Path.Combine(ProxyFixture.StoreRoot, "github");
fixture.FileSystem.Directory.CreateDirectory(upstreamDirectory);
fixture.FileSystem.File.WriteAllText(Path.Combine(upstreamDirectory, "staging"), "not a directory");
}

/// <summary>Counts <c>gitlfscache.staging_failures</c> recorded while the listener is alive.</summary>
private sealed class StagingFailureCounter : IDisposable
{
private readonly MeterListener _listener = new();
private long _count;

public StagingFailureCounter()
{
_listener.InstrumentPublished = (instrument, listener) =>
{
if (instrument.Meter.Name == CacheMetrics.MeterName && instrument.Name == "gitlfscache.staging_failures")
{
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();
}

[TestMethod]
public async Task Download_StagingCannotBeOpened_IsServedFromUpstreamWithoutCaching()
{
await using ProxyFixture fixture = await ProxyFixture.StartAsync();
using StagingFailureCounter failures = new();
(byte[] content, string oid) = Object("pulled while the store is broken");
fixture.Upstream.AddObject(oid, content);
BlockStaging(fixture);

JsonNode batch = await PostBatchAsync(fixture, "download", oid, content.Length);
using HttpClient client = fixture.Client;
using HttpResponseMessage response = await client.GetAsync(Relative(HrefOf(batch, "download")));

Check warning on line 363 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=AaDnluXduxf7YbO2HskC&open=AaDnluXduxf7YbO2HskC&pullRequest=68

Assert.AreEqual(HttpStatusCode.OK, response.StatusCode);
CollectionAssert.AreEqual(content, await response.Content.ReadAsByteArrayAsync());

Check warning on line 366 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=AaDnluXduxf7YbO2HskB&open=AaDnluXduxf7YbO2HskB&pullRequest=68

Check warning on line 366 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=AaDnluXduxf7YbO2HskD&open=AaDnluXduxf7YbO2HskD&pullRequest=68
Assert.AreEqual(1, fixture.Upstream.FetchCount(oid));
Assert.IsFalse(fixture.Store.Exists("github", oid));
Assert.AreEqual(1, failures.Count);
}

[TestMethod]
public async Task Upload_StagingCannotBeOpened_IsStillRelayedUpstream()
{
await using ProxyFixture fixture = await ProxyFixture.StartAsync();
using StagingFailureCounter failures = new();
(byte[] content, string oid) = Object("pushed while the store is broken");
BlockStaging(fixture);

JsonNode batch = await PostBatchAsync(fixture, "upload", oid, content.Length);
using HttpClient client = fixture.Client;
using ByteArrayContent body = new(content);
using HttpResponseMessage response = await client.PutAsync(Relative(HrefOf(batch, "upload")), body);

Check warning on line 383 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=AaDnluXduxf7YbO2HskF&open=AaDnluXduxf7YbO2HskF&pullRequest=68

Assert.AreEqual(HttpStatusCode.OK, response.StatusCode);
CollectionAssert.AreEqual(content, fixture.Upstream.Uploaded[oid]);

Check warning on line 386 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=AaDnluXduxf7YbO2HskE&open=AaDnluXduxf7YbO2HskE&pullRequest=68
Assert.IsFalse(fixture.Store.Exists("github", oid));
Assert.AreEqual(1, failures.Count);
}

[TestMethod]
public async Task Upload_ThenDownload_IsServedFromTheStoreWithoutFetchingUpstream()
{
Expand Down
6 changes: 6 additions & 0 deletions GitLfsCache/Endpoints/EndpointLog.cs
Original file line number Diff line number Diff line change
Expand Up @@ -80,4 +80,10 @@ internal static partial class EndpointLog
Level = LogLevel.Warning,
Message = "Request for '{Path}' refused: no pattern in upstream '{Upstream}' allows it.")]
public static partial void RepositoryNotAllowed(ILogger logger, string path, string upstream);

[LoggerMessage(
EventId = 2012,
Level = LogLevel.Warning,
Message = "Could not open a staging file for {Oid} under upstream '{Upstream}'; the transfer is relayed without caching.")]
public static partial void StagingUnavailable(ILogger logger, Exception exception, string oid, string upstream);
}
52 changes: 45 additions & 7 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 / 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 / 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 / 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 / Analyze & Release

Constructor has 9 parameters, which is greater than the 7 authorized.
IUpstreamClient upstreamClient,
IHrefTokenCodec codec,
BatchRewriter rewriter,
Expand Down Expand Up @@ -223,13 +223,16 @@
return;
}

StagingHandle staging = store.OpenStaging(route.Upstream);
// Without a staging file the upload is still relayed, through a tee into nothing so the relayed
// byte count is kept. Failing the push because the cache cannot take a copy would make the
// cache the reason a push failed.
StagingHandle? staging = TryOpenStaging(route.Upstream, token.Oid);

await using (staging.ConfigureAwait(false))
try
{
ReadTeeStream teed = new(
context.Request.Body,
staging.Stream,
staging?.Stream ?? Stream.Null,
failure => EndpointLog.StoreSinkFailed(logger, failure, token.Oid));

await using ConfiguredAsyncDisposable teedDisposal = teed.ConfigureAwait(false);
Expand All @@ -247,7 +250,7 @@

// The object is published only after upstream accepts it. Caching an upload upstream
// rejected would serve bytes no one can verify against the real remote.
if (response.IsSuccessStatusCode && teed.SinkIsLive)
if (staging is not null && response.IsSuccessStatusCode && teed.SinkIsLive)
{
if (await store.PublishAsync(staging, route.Upstream, token.Oid, cancellationToken)
.ConfigureAwait(false))
Expand All @@ -262,6 +265,13 @@

await UpstreamRelay.CopyResponseAsync(response, context, cancellationToken).ConfigureAwait(false);
}
finally
{
if (staging is not null)
{
await staging.DisposeAsync().ConfigureAwait(false);
}
}
}

/// <summary>
Expand Down Expand Up @@ -319,7 +329,12 @@

await using ConfiguredAsyncDisposable upstreamBodyDisposal = upstreamBody.ConfigureAwait(false);

if (!storeLocally)
// By now upstream's status and headers are already on the response, so a staging file that
// cannot be opened has to fall back to relaying: throwing here would turn an object upstream
// served into a 500.
StagingHandle? staging = storeLocally ? TryOpenStaging(route.Upstream, token.Oid) : null;

if (staging is null)
{
long streamed = await StreamTee
.CopyAsync(upstreamBody, context.Response.Body, null, null, cancellationToken)
Expand All @@ -329,8 +344,6 @@
return false;
}

StagingHandle staging = store.OpenStaging(route.Upstream);

await using (staging.ConfigureAwait(false))
{
long streamed = await StreamTee.CopyAsync(
Expand Down Expand Up @@ -359,6 +372,31 @@
}
}

/// <summary>
/// Opens a staging file, or reports why the transfer has to go uncached.
/// </summary>
/// <remarks>
/// The startup probe only proves the store was writable when the process started. A permissions
/// change, a read-only remount, inode exhaustion or a stray file where the staging directory
/// belongs all appear later, and each one should cost a cold cache rather than a failed transfer.
/// </remarks>
/// <param name="upstream">The upstream key the object belongs to.</param>
/// <param name="oid">The object id, for the log.</param>
/// <returns>The open staging file, or null when one could not be opened.</returns>
private StagingHandle? TryOpenStaging(string upstream, string oid)
{
try
{
return store.OpenStaging(upstream);
}
catch (Exception exception) when (exception is IOException or UnauthorizedAccessException or ArgumentException)
{
EndpointLog.StagingUnavailable(logger, exception, oid, upstream);
metrics.RecordStagingFailure(upstream);
return null;
}
}

private bool TryGetToken(
HttpContext context,
LfsRoute route,
Expand Down
12 changes: 12 additions & 0 deletions GitLfsCache/Observability/CacheMetrics.cs
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,7 @@ public sealed class CacheMetrics : IDisposable
private readonly Counter<long> _bytesUploaded;
private readonly Counter<long> _objectsStored;
private readonly Counter<long> _verificationFailures;
private readonly Counter<long> _stagingFailures;
private readonly Counter<long> _coalescedWaits;
private readonly Counter<long> _rejectedTokens;
private readonly Counter<long> _lockListHits;
Expand Down Expand Up @@ -62,6 +63,7 @@ public CacheMetrics(IMeterFactory meterFactory)
_bytesUploaded = _meter.CreateCounter<long>("gitlfscache.upload_bytes_relayed", unit: "By", description: "Bytes relayed to upstream on upload.");
_objectsStored = _meter.CreateCounter<long>("gitlfscache.objects_stored", unit: ObjectUnit, description: "Objects verified and published to the store.");
_verificationFailures = _meter.CreateCounter<long>("gitlfscache.verification_failures", unit: ObjectUnit, description: "Transfers whose content did not hash to the expected object id.");
_stagingFailures = _meter.CreateCounter<long>("gitlfscache.staging_failures", unit: ObjectUnit, description: "Transfers relayed without caching because the store could not open a staging file.");
_coalescedWaits = _meter.CreateCounter<long>("gitlfscache.coalesced_waits", unit: RequestUnit, description: "Requests that waited for another request's fetch instead of fetching themselves.");
_rejectedTokens = _meter.CreateCounter<long>("gitlfscache.rejected_tokens", unit: RequestUnit, description: "Requests refused because their transfer token was invalid or expired.");
_lockListHits = _meter.CreateCounter<long>("gitlfscache.lock_list_hits", unit: RequestUnit, description: "Lock listings answered from a snapshot without reaching upstream.");
Expand Down Expand Up @@ -111,6 +113,16 @@ public void RecordStored(string upstream) =>
public void RecordVerificationFailure(string upstream) =>
_verificationFailures.Add(1, new KeyValuePair<string, object?>("upstream", upstream));

/// <summary>Records a transfer relayed without caching because no staging file could be opened.</summary>
/// <remarks>
/// Anything above zero means the store has stopped accepting writes (permissions, a read-only
/// remount, inode exhaustion) while clients carry on being served. The cache only goes cold, so
/// nothing else makes the condition visible.
/// </remarks>
/// <param name="upstream">The upstream key, recorded as a tag.</param>
public void RecordStagingFailure(string upstream) =>
_stagingFailures.Add(1, new KeyValuePair<string, object?>("upstream", upstream));

/// <summary>Records a request that waited for another request's fetch.</summary>
/// <param name="upstream">The upstream key, recorded as a tag.</param>
public void RecordCoalescedWait(string upstream) =>
Expand Down
Loading