diff --git a/GitLfsCache.Tests/Integration/ProxyFlowTests.cs b/GitLfsCache.Tests/Integration/ProxyFlowTests.cs index 3955c9f..dcc58e5 100644 --- a/GitLfsCache.Tests/Integration/ProxyFlowTests.cs +++ b/GitLfsCache.Tests/Integration/ProxyFlowTests.cs @@ -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] @@ -308,6 +310,84 @@ public async Task Upload_IsRelayedUpstreamAndStoredOnTheWayThrough() Assert.IsTrue(fixture.Store.Exists("github", oid), "A pushed object should be cached for the next fetch."); } + /// + /// Puts a file where the upstream's staging directory belongs, so no staging file can be opened. + /// + /// + /// 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. + /// + 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"); + } + + /// Counts gitlfscache.staging_failures recorded while the listener is alive. + 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((_, 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"))); + + Assert.AreEqual(HttpStatusCode.OK, response.StatusCode); + CollectionAssert.AreEqual(content, await response.Content.ReadAsByteArrayAsync()); + 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); + + Assert.AreEqual(HttpStatusCode.OK, response.StatusCode); + CollectionAssert.AreEqual(content, fixture.Upstream.Uploaded[oid]); + Assert.IsFalse(fixture.Store.Exists("github", oid)); + Assert.AreEqual(1, failures.Count); + } + [TestMethod] public async Task Upload_ThenDownload_IsServedFromTheStoreWithoutFetchingUpstream() { diff --git a/GitLfsCache/Endpoints/EndpointLog.cs b/GitLfsCache/Endpoints/EndpointLog.cs index 93c4bbe..cc27b15 100644 --- a/GitLfsCache/Endpoints/EndpointLog.cs +++ b/GitLfsCache/Endpoints/EndpointLog.cs @@ -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); } diff --git a/GitLfsCache/Endpoints/ObjectRouteHandler.cs b/GitLfsCache/Endpoints/ObjectRouteHandler.cs index 106b6d0..873d389 100644 --- a/GitLfsCache/Endpoints/ObjectRouteHandler.cs +++ b/GitLfsCache/Endpoints/ObjectRouteHandler.cs @@ -223,13 +223,16 @@ public async Task UploadAsync(HttpContext context, LfsRoute route, CancellationT 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); @@ -247,7 +250,7 @@ public async Task UploadAsync(HttpContext context, LfsRoute route, CancellationT // 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)) @@ -262,6 +265,13 @@ public async Task UploadAsync(HttpContext context, LfsRoute route, CancellationT await UpstreamRelay.CopyResponseAsync(response, context, cancellationToken).ConfigureAwait(false); } + finally + { + if (staging is not null) + { + await staging.DisposeAsync().ConfigureAwait(false); + } + } } /// @@ -319,7 +329,12 @@ private async Task StreamFromUpstreamAsync( 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) @@ -329,8 +344,6 @@ private async Task StreamFromUpstreamAsync( return false; } - StagingHandle staging = store.OpenStaging(route.Upstream); - await using (staging.ConfigureAwait(false)) { long streamed = await StreamTee.CopyAsync( @@ -359,6 +372,31 @@ private async Task StreamFromUpstreamAsync( } } + /// + /// Opens a staging file, or reports why the transfer has to go uncached. + /// + /// + /// 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. + /// + /// The upstream key the object belongs to. + /// The object id, for the log. + /// The open staging file, or null when one could not be opened. + 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, diff --git a/GitLfsCache/Observability/CacheMetrics.cs b/GitLfsCache/Observability/CacheMetrics.cs index 18e2ed1..5c43a0f 100644 --- a/GitLfsCache/Observability/CacheMetrics.cs +++ b/GitLfsCache/Observability/CacheMetrics.cs @@ -34,6 +34,7 @@ public sealed class CacheMetrics : IDisposable private readonly Counter _bytesUploaded; private readonly Counter _objectsStored; private readonly Counter _verificationFailures; + private readonly Counter _stagingFailures; private readonly Counter _coalescedWaits; private readonly Counter _rejectedTokens; private readonly Counter _lockListHits; @@ -62,6 +63,7 @@ public CacheMetrics(IMeterFactory meterFactory) _bytesUploaded = _meter.CreateCounter("gitlfscache.upload_bytes_relayed", unit: "By", description: "Bytes relayed to upstream on upload."); _objectsStored = _meter.CreateCounter("gitlfscache.objects_stored", unit: ObjectUnit, description: "Objects verified and published to the store."); _verificationFailures = _meter.CreateCounter("gitlfscache.verification_failures", unit: ObjectUnit, description: "Transfers whose content did not hash to the expected object id."); + _stagingFailures = _meter.CreateCounter("gitlfscache.staging_failures", unit: ObjectUnit, description: "Transfers relayed without caching because the store could not open a staging file."); _coalescedWaits = _meter.CreateCounter("gitlfscache.coalesced_waits", unit: RequestUnit, description: "Requests that waited for another request's fetch instead of fetching themselves."); _rejectedTokens = _meter.CreateCounter("gitlfscache.rejected_tokens", unit: RequestUnit, description: "Requests refused because their transfer token was invalid or expired."); _lockListHits = _meter.CreateCounter("gitlfscache.lock_list_hits", unit: RequestUnit, description: "Lock listings answered from a snapshot without reaching upstream."); @@ -111,6 +113,16 @@ public void RecordStored(string upstream) => public void RecordVerificationFailure(string upstream) => _verificationFailures.Add(1, new KeyValuePair("upstream", upstream)); + /// Records a transfer relayed without caching because no staging file could be opened. + /// + /// 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. + /// + /// The upstream key, recorded as a tag. + public void RecordStagingFailure(string upstream) => + _stagingFailures.Add(1, new KeyValuePair("upstream", upstream)); + /// Records a request that waited for another request's fetch. /// The upstream key, recorded as a tag. public void RecordCoalescedWait(string upstream) =>