diff --git a/GitLfsCache.Tests/Endpoints/HandlerCompositionTests.cs b/GitLfsCache.Tests/Endpoints/HandlerCompositionTests.cs new file mode 100644 index 0000000..e2e7773 --- /dev/null +++ b/GitLfsCache.Tests/Endpoints/HandlerCompositionTests.cs @@ -0,0 +1,111 @@ +// Copyright (c) 2023-2026 ktsu-dev contributors + +namespace ktsu.GitLfsCache.Tests.Endpoints; + +using System.Reflection; +using ktsu.GitLfsCache.Endpoints; +using ktsu.GitLfsCache.Locks; +using ktsu.GitLfsCache.Storage; +using ktsu.GitLfsCache.Tokens; +using Microsoft.VisualStudio.TestTools.UnitTesting; + +/// +/// Guards the split between the dispatcher and the two route handlers. +/// +/// +/// Written against constructors rather than behaviour because what is being guarded is a dependency +/// boundary, and the behaviour either side of it is already covered by the integration suite. The +/// dispatcher acquired its lock dependencies one at a time, each reasonable on its own, until it held +/// fourteen; a test that fails the moment an unrelated dependency is added is the only thing that +/// notices that happening again. +/// +[TestClass] +public class HandlerCompositionTests +{ + /// + /// The constructor parameter budget SonarQube enforces (S107). + /// + private const int MaxConstructorParameters = 7; + + private static ParameterInfo[] ConstructorParametersOf() => + typeof(T).GetConstructors(BindingFlags.Public | BindingFlags.NonPublic | BindingFlags.Instance) + .Single() + .GetParameters(); + + [TestMethod] + public void Dispatcher_StaysWithinTheConstructorParameterBudget() + { + ParameterInfo[] parameters = ConstructorParametersOf(); + + Assert.IsLessThanOrEqualTo( + MaxConstructorParameters, + parameters.Length, + $"{nameof(GitLfsCacheHandler)} dispatches; it does not do the work. It now takes " + + $"{parameters.Length} dependencies ({string.Join(", ", parameters.Select(p => p.Name))}). " + + "A new dependency here almost always belongs to one of the route handlers instead."); + } + + [TestMethod] + public void LockHandler_StaysWithinTheConstructorParameterBudget() + { + ParameterInfo[] parameters = ConstructorParametersOf(); + + Assert.IsLessThanOrEqualTo(MaxConstructorParameters, parameters.Length); + } + + /// + /// The lock routes need the upstream and the allow-list check, both already done by the dispatcher + /// before either handler is called. Nothing about transfer tokens or the object store reaches them, + /// and taking either back would put the two concerns into one class again. + /// + [TestMethod] + public void LockHandler_TakesNoObjectSideDependencies() + { + IEnumerable dependencies = ConstructorParametersOf() + .Select(parameter => parameter.ParameterType); + + Assert.IsFalse( + dependencies.Any(type => type == typeof(IObjectStore) || type == typeof(IHrefTokenCodec)), + $"{nameof(LockRouteHandler)} took an object-store or transfer-token dependency."); + } + + /// + /// The three the issue named, and the reason the split was worth making. + /// + [TestMethod] + public void LockHandler_HoldsTheLockDependencies() + { + IEnumerable dependencies = ConstructorParametersOf() + .Select(parameter => parameter.ParameterType); + + CollectionAssert.IsSubsetOf( + new[] { typeof(LockListService), typeof(ILockSnapshotStore), typeof(LockFanOut) }, + dependencies.ToList()); + } + + /// + /// The dispatcher owning a lock dependency directly is how the previous shape started. + /// + [TestMethod] + public void Dispatcher_RoutesLocksThroughTheLockHandlerRatherThanHoldingItsDependencies() + { + List dependencies = [.. ConstructorParametersOf() + .Select(parameter => parameter.ParameterType)]; + + CollectionAssert.Contains(dependencies, typeof(LockRouteHandler)); + + foreach (Type lockDependency in new[] + { + typeof(LockListService), + typeof(ILockSnapshotStore), + typeof(LockFanOut), + }) + { + CollectionAssert.DoesNotContain( + dependencies, + lockDependency, + $"{nameof(GitLfsCacheHandler)} holds {lockDependency.Name} directly; it belongs to " + + $"{nameof(LockRouteHandler)}."); + } + } +} diff --git a/GitLfsCache/Endpoints/GitLfsCacheHandler.cs b/GitLfsCache/Endpoints/GitLfsCacheHandler.cs index 3f769b8..01f7067 100644 --- a/GitLfsCache/Endpoints/GitLfsCacheHandler.cs +++ b/GitLfsCache/Endpoints/GitLfsCacheHandler.cs @@ -2,63 +2,38 @@ namespace ktsu.GitLfsCache.Endpoints; -using System.Diagnostics.CodeAnalysis; -using System.Net; -using System.Runtime.CompilerServices; -using System.Text.Json.Nodes; -using ktsu.GitLfsCache.Batch; using ktsu.GitLfsCache.Configuration; -using ktsu.GitLfsCache.Fetching; -using ktsu.GitLfsCache.Locks; -using ktsu.GitLfsCache.Observability; -using ktsu.GitLfsCache.Storage; -using ktsu.GitLfsCache.Tokens; using ktsu.GitLfsCache.Upstreams; using Microsoft.AspNetCore.Http; using Microsoft.Extensions.Logging; using Microsoft.Extensions.Options; /// -/// Handles every request under an upstream prefix. +/// Dispatches every request under an upstream prefix to the handler its route selects. /// /// -/// One handler rather than several because the paths share their front half: resolve the upstream, -/// validate the token, and decide whether the bytes come from the store or from upstream. Splitting -/// that across five endpoint classes would mean five copies of the same preamble. +/// The preamble that every route genuinely shares lives here and runs once: parse the path, resolve +/// the upstream, and check the repository against the allow-list. Past that point the routes divide +/// into two groups that share nothing else, so each has its own handler — +/// for batch, transfer and verify, which need transfer tokens and +/// the object store, and for the lock routes, which need neither. /// /// Resolves upstream keys. /// Decides which repository paths an upstream may be used for. -/// Sends requests upstream. -/// Decodes transfer tokens. -/// Rewrites batch responses. -/// The local object store. -/// Keeps concurrent misses to one upstream fetch. -/// Answers lock listings from a snapshot. -/// Holds lock snapshots, so a relayed change can invalidate one. -/// Runs the individual calls of a batched lock request. -/// Resolves the base URL rewritten hrefs point at. -/// Cache counters. +/// Handles batch, transfer and verify. +/// Handles the lock routes. +/// Passes anything the proxy does not model upstream verbatim. /// The configured options. /// Logger. -public sealed class GitLfsCacheHandler( +internal sealed class GitLfsCacheHandler( IUpstreamRegistry registry, IRepositoryAllowList allowList, - IUpstreamClient upstreamClient, - IHrefTokenCodec codec, - BatchRewriter rewriter, - IObjectStore store, - IFetchCoalescer coalescer, - LockListService lockLists, - ILockSnapshotStore lockSnapshots, - LockFanOut lockFanOut, - PublicUrlResolver publicUrls, - CacheMetrics metrics, + ObjectRouteHandler objects, + LockRouteHandler locks, + UpstreamRelay relay, IOptions options, ILogger logger) { - private const string OctetStream = "application/octet-stream"; - private const string TokenQueryParameter = "t"; - /// /// Dispatches one request. /// @@ -105,620 +80,41 @@ public async Task HandleAsync(HttpContext context) switch (route.Kind) { case LfsRouteKind.Batch when caching && HttpMethods.IsPost(context.Request.Method): - await BatchAsync(context, route, upstreamBase, cancellationToken).ConfigureAwait(false); + await objects.BatchAsync(context, route, upstreamBase, cancellationToken).ConfigureAwait(false); return; case LfsRouteKind.Transfer when caching && HttpMethods.IsGet(context.Request.Method): - await DownloadAsync(context, route, cancellationToken).ConfigureAwait(false); + await objects.DownloadAsync(context, route, cancellationToken).ConfigureAwait(false); return; case LfsRouteKind.Transfer when caching && HttpMethods.IsPut(context.Request.Method): - await UploadAsync(context, route, cancellationToken).ConfigureAwait(false); + await objects.UploadAsync(context, route, cancellationToken).ConfigureAwait(false); return; case LfsRouteKind.Verify when caching && HttpMethods.IsPost(context.Request.Method): - await VerifyAsync(context, route, cancellationToken).ConfigureAwait(false); + await objects.VerifyAsync(context, route, cancellationToken).ConfigureAwait(false); return; case LfsRouteKind.Locks when HttpMethods.IsGet(context.Request.Method): - await LockListAsync(context, route, upstreamBase, cancellationToken).ConfigureAwait(false); + await locks.ListAsync(context, route, upstreamBase, cancellationToken).ConfigureAwait(false); return; case LfsRouteKind.LocksBatch when HttpMethods.IsPost(context.Request.Method): - await LockFanOutAsync(context, route, upstreamBase, cancellationToken).ConfigureAwait(false); + await locks.FanOutAsync(context, route, upstreamBase, cancellationToken).ConfigureAwait(false); return; // Creation and release are relayed, never terminated, because upstream is the only thing - // that may grant or release a lock. The snapshot is dropped afterwards so the change this - // client just made is visible to the next listing rather than waiting out the lifetime. + // that may grant or release a lock, and the snapshot is dropped afterwards. case LfsRouteKind.Locks when HttpMethods.IsPost(context.Request.Method): case LfsRouteKind.LocksUnlock when HttpMethods.IsPost(context.Request.Method): - await RelayAsync(context, route, upstreamBase, cancellationToken).ConfigureAwait(false); - InvalidateLocksIfChanged(context, route); + await locks.RelayChangeAsync(context, route, upstreamBase, cancellationToken).ConfigureAwait(false); return; default: // Includes a recognized path reached with an unexpected method. Relaying rather than // rejecting keeps the proxy transparent to anything it does not model. - await RelayAsync(context, route, upstreamBase, cancellationToken).ConfigureAwait(false); - return; - } - } - - private async Task BatchAsync( - HttpContext context, - LfsRoute route, - Uri upstreamBase, - CancellationToken cancellationToken) - { - using HttpRequestMessage request = UpstreamRequests.BuildBatchRequest( - upstreamBase, - route.RepositoryPath, - context.Request.Body, - context.Request.Headers.Authorization.ToString()); - - using HttpResponseMessage response = await upstreamClient - .SendAsync(request, cancellationToken) - .ConfigureAwait(false); - - // Upstream is the authority on access. A refusal is relayed exactly as it arrived, so the - // client sees upstream's real answer rather than a proxy interpretation of it. - if (!response.IsSuccessStatusCode) - { - await CopyResponseAsync(response, context, cancellationToken).ConfigureAwait(false); - return; - } - - JsonNode? upstreamBody; - - Stream batchBody = await response.Content - .ReadAsStreamAsync(cancellationToken) - .ConfigureAwait(false); - - await using (batchBody.ConfigureAwait(false)) - { - upstreamBody = await JsonNode.ParseAsync(batchBody, cancellationToken: cancellationToken) - .ConfigureAwait(false); - } - - if (upstreamBody is null) - { - context.Response.StatusCode = StatusCodes.Status502BadGateway; - return; - } - - JsonNode rewritten = rewriter.Rewrite(upstreamBody, new BatchRewriteContext - { - Upstream = route.Upstream, - RepositoryPath = route.RepositoryPath, - PublicBaseUrl = publicUrls.Resolve(context.Request), - }); - - context.Response.StatusCode = StatusCodes.Status200OK; - context.Response.ContentType = UpstreamRequests.LfsMediaType; - await context.Response - .WriteAsync(rewritten.ToJsonString(), cancellationToken) - .ConfigureAwait(false); - } - - private async Task DownloadAsync(HttpContext context, LfsRoute route, CancellationToken cancellationToken) - { - if (!TryGetToken(context, route, TokenAction.Download, out HrefToken? token)) - { - return; - } - - string range = context.Request.Headers.Range.ToString(); - - Stream? cached = store.OpenRead(route.Upstream, token.Oid, out long length); - - if (cached is not null) - { - await using (cached.ConfigureAwait(false)) - { - store.Touch(route.Upstream, token.Oid); - metrics.RecordHit(route.Upstream, length); - EndpointLog.ServedFromCache(logger, token.Oid, route.Upstream); - - await ServeFromStoreAsync(context, cached, length, cancellationToken).ConfigureAwait(false); - } - - return; - } - - if (!string.IsNullOrEmpty(range)) - { - // A partial response cannot be stored as a whole object, so the range is forwarded and the - // result streamed straight through. Rare enough not to be worth partial-object bookkeeping. - EndpointLog.RangeRequestNotStored(logger, token.Oid); - await StreamFromUpstreamAsync(context, route, token, range, storeLocally: false, cancellationToken) - .ConfigureAwait(false); - return; - } - - using IFetchTicket ticket = coalescer.Acquire(route.Upstream, token.Oid); - - if (!ticket.IsLeader) - { - metrics.RecordCoalescedWait(route.Upstream); - EndpointLog.WaitingForLeader(logger, token.Oid, route.Upstream); - - bool published = await ticket - .WaitForLeaderAsync(options.Value.Fetch.FollowerTimeout, cancellationToken) - .ConfigureAwait(false); - - long nowLength = 0; - Stream? nowCached = published - ? store.OpenRead(route.Upstream, token.Oid, out nowLength) - : null; - - if (nowCached is not null) - { - 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; - } - - EndpointLog.LeaderDidNotFinish(logger, token.Oid); - } - - bool stored = await StreamFromUpstreamAsync( - context, - route, - token, - range: null, - storeLocally: true, - cancellationToken).ConfigureAwait(false); - - if (ticket.IsLeader) - { - ticket.Complete(stored); - } - } - - private async Task StreamFromUpstreamAsync( - HttpContext context, - LfsRoute route, - HrefToken token, - string? range, - bool storeLocally, - CancellationToken cancellationToken) - { - EndpointLog.FetchingUpstream(logger, token.Oid, route.Upstream); - - using HttpRequestMessage request = UpstreamRequests.BuildObjectRequest(token, range); - using HttpResponseMessage response = await upstreamClient - .SendAsync(request, cancellationToken) - .ConfigureAwait(false); - - if (!response.IsSuccessStatusCode) - { - EndpointLog.UpstreamRefusedTransfer(logger, (int)response.StatusCode, token.Oid); - await CopyResponseAsync(response, context, cancellationToken).ConfigureAwait(false); - return false; - } - - CopyTransferHeaders(response, context); - - Stream upstreamBody = await response.Content - .ReadAsStreamAsync(cancellationToken) - .ConfigureAwait(false); - - await using ConfiguredAsyncDisposable upstreamBodyDisposal = upstreamBody.ConfigureAwait(false); - - if (!storeLocally) - { - long streamed = await StreamTee - .CopyAsync(upstreamBody, context.Response.Body, null, null, cancellationToken) - .ConfigureAwait(false); - - metrics.RecordMiss(route.Upstream, streamed); - return false; - } - - StagingHandle staging = store.OpenStaging(route.Upstream); - - await using (staging.ConfigureAwait(false)) - { - long streamed = await StreamTee.CopyAsync( - upstreamBody, - context.Response.Body, - staging.Stream, - failure => EndpointLog.StoreSinkFailed(logger, failure, token.Oid), - cancellationToken).ConfigureAwait(false); - - metrics.RecordMiss(route.Upstream, streamed); - - bool published = await store - .PublishAsync(staging, route.Upstream, token.Oid, cancellationToken) - .ConfigureAwait(false); - - if (published) - { - metrics.RecordStored(route.Upstream); - } - else - { - metrics.RecordVerificationFailure(route.Upstream); - } - - return published; - } - } - - private async Task UploadAsync(HttpContext context, LfsRoute route, CancellationToken cancellationToken) - { - if (!TryGetToken(context, route, TokenAction.Upload, out HrefToken? token)) - { - return; - } - - StagingHandle staging = store.OpenStaging(route.Upstream); - - await using (staging.ConfigureAwait(false)) - { - ReadTeeStream teed = new( - context.Request.Body, - staging.Stream, - failure => EndpointLog.StoreSinkFailed(logger, failure, token.Oid)); - - await using ConfiguredAsyncDisposable teedDisposal = teed.ConfigureAwait(false); - - using HttpRequestMessage request = UpstreamRequests.BuildUploadRequest( - token, - teed, - context.Request.ContentLength); - - using HttpResponseMessage response = await upstreamClient - .SendAsync(request, cancellationToken) - .ConfigureAwait(false); - - metrics.RecordUpload(route.Upstream, teed.BytesRead); - - // 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 (await store.PublishAsync(staging, route.Upstream, token.Oid, cancellationToken) - .ConfigureAwait(false)) - { - metrics.RecordStored(route.Upstream); - } - else - { - metrics.RecordVerificationFailure(route.Upstream); - } - } - - await CopyResponseAsync(response, context, cancellationToken).ConfigureAwait(false); - } - } - - private async Task VerifyAsync(HttpContext context, LfsRoute route, CancellationToken cancellationToken) - { - if (!TryGetToken(context, route, TokenAction.Verify, out HrefToken? token)) - { - return; - } - - using HttpRequestMessage request = UpstreamRequests.BuildVerifyRequest(token, context.Request.Body); - using HttpResponseMessage response = await upstreamClient - .SendAsync(request, cancellationToken) - .ConfigureAwait(false); - - await CopyResponseAsync(response, context, cancellationToken).ConfigureAwait(false); - } - - /// - /// Answers a lock listing, from the snapshot when that is both possible and permitted. - /// - private async Task LockListAsync( - HttpContext context, - LfsRoute route, - Uri upstreamBase, - CancellationToken cancellationToken) - { - LockSnapshotKey key = new( - route.Upstream, - route.RepositoryPath, - context.Request.Query["refspec"].FirstOrDefault()); - - LockListOutcome outcome = await lockLists - .ResolveAsync(key, upstreamBase, context.Request.Headers.Authorization.ToString(), cancellationToken) - .ConfigureAwait(false); - - switch (outcome.Kind) - { - case LockListOutcomeKind.Refuse: - // Upstream's own refusal, not a proxy interpretation of it. - context.Response.StatusCode = (int)outcome.Status!.Value; - return; - - case LockListOutcomeKind.Serve: - await WriteLockPageAsync(context, outcome.Snapshot!, cancellationToken).ConfigureAwait(false); + await relay.RelayAsync(context, route, upstreamBase, cancellationToken).ConfigureAwait(false); return; - - default: - await RelayAsync(context, route, upstreamBase, cancellationToken).ConfigureAwait(false); - return; - } - } - - /// - /// Writes one page of a snapshot, applying the filters and cursor the client asked for. - /// - private static async Task WriteLockPageAsync( - HttpContext context, - LockSnapshot snapshot, - CancellationToken cancellationToken) - { - IReadOnlyList matches = snapshot.Filter( - context.Request.Query["path"].FirstOrDefault(), - context.Request.Query["id"].FirstOrDefault()); - - int offset = 0; - - // A cursor from a snapshot that has since been replaced restarts the walk rather than being - // applied to a different ordering, which would silently skip or repeat locks. - if (LockCursor.TryDecode(context.Request.Query["cursor"].FirstOrDefault(), out LockCursor? cursor) - && cursor.SnapshotId == snapshot.Id) - { - offset = cursor.Offset; - } - - int? limit = int.TryParse( - context.Request.Query["limit"].FirstOrDefault(), - System.Globalization.NumberStyles.None, - System.Globalization.CultureInfo.InvariantCulture, - out int requested) - ? requested - : null; - - (IReadOnlyList page, int? nextOffset) = LockSnapshot.Paginate(matches, offset, limit); - - JsonObject body = LockListParser.BuildResponse( - page, - nextOffset is int next ? new LockCursor(snapshot.Id, next).Encode() : null); - - context.Response.StatusCode = StatusCodes.Status200OK; - context.Response.ContentType = UpstreamRequests.LfsMediaType; - - await context.Response - .WriteAsync(body.ToJsonString(), cancellationToken) - .ConfigureAwait(false); - } - - /// - /// Runs a batched lock or unlock, issuing the individual calls in parallel. - /// - /// - /// A proxy extension, so it is refused rather than relayed when the subsystem is switched off: an - /// upstream has no such endpoint, and relaying would turn a disabled feature into a confusing 404 - /// from the forge instead of a clear one from here. - /// - private async Task LockFanOutAsync( - HttpContext context, - LfsRoute route, - Uri upstreamBase, - CancellationToken cancellationToken) - { - if (!options.Value.Locks.Enabled) - { - context.Response.StatusCode = StatusCodes.Status404NotFound; - return; - } - - JsonNode? body; - - try - { - body = await JsonNode.ParseAsync(context.Request.Body, cancellationToken: cancellationToken) - .ConfigureAwait(false); - } - catch (System.Text.Json.JsonException) - { - context.Response.StatusCode = StatusCodes.Status400BadRequest; - return; - } - - if (!LockFanOutRequest.TryParse(body, out LockFanOutRequest? request)) - { - context.Response.StatusCode = StatusCodes.Status400BadRequest; - return; - } - - // Refused outright rather than accepted and throttled part way through, which would leave the - // caller reconciling a partial result they never asked for. - if (request.Targets.Count > options.Value.Locks.MaxFanOutPaths) - { - context.Response.StatusCode = StatusCodes.Status413PayloadTooLarge; - return; - } - - LockSnapshotKey key = new(route.Upstream, route.RepositoryPath, request.Ref); - - JsonObject results = await lockFanOut - .ExecuteAsync( - request, - key, - upstreamBase, - context.Request.Headers.Authorization.ToString(), - cancellationToken) - .ConfigureAwait(false); - - // Always 200 when the request itself was well formed. Partial success is the normal outcome, - // and a transport-level failure would discard the half that worked. - context.Response.StatusCode = StatusCodes.Status200OK; - context.Response.ContentType = UpstreamRequests.LfsMediaType; - - await context.Response - .WriteAsync(results.ToJsonString(), cancellationToken) - .ConfigureAwait(false); - } - - /// - /// Drops the snapshot when a relayed lock change actually took effect. - /// - /// - /// Gated on the response status, because invalidating after a refused creation would throw away a - /// perfectly good snapshot every time two people raced for the same file, which is exactly when - /// the cache is under the most load. - /// - private void InvalidateLocksIfChanged(HttpContext context, LfsRoute route) - { - if (context.Response.StatusCode is >= 200 and < 300) - { - lockSnapshots.Invalidate(new LockSnapshotKey( - route.Upstream, - route.RepositoryPath, - context.Request.Query["refspec"].FirstOrDefault())); - } - } - - private async Task RelayAsync( - HttpContext context, - LfsRoute route, - Uri upstreamBase, - CancellationToken cancellationToken) - { - IEnumerable>> headers = context.Request.Headers - .Select(header => new KeyValuePair>( - header.Key, - header.Value.Where(value => value is not null).Select(value => value!))); - - using HttpRequestMessage request = UpstreamRequests.BuildRelayRequest( - upstreamBase, - context.Request.Method, - route.RelayPath, - context.Request.QueryString.Value ?? string.Empty, - context.Request.Body, - headers); - - using HttpResponseMessage response = await upstreamClient - .SendAsync(request, cancellationToken) - .ConfigureAwait(false); - - EndpointLog.Relayed(logger, context.Request.Method, route.RelayPath, route.Upstream); - await CopyResponseAsync(response, context, cancellationToken).ConfigureAwait(false); - } - - private bool TryGetToken( - HttpContext context, - LfsRoute route, - string expectedAction, - [NotNullWhen(true)] out HrefToken? token) - { - token = null; - string? encoded = context.Request.Query[TokenQueryParameter]; - - if (!codec.TryDecode(encoded, out HrefToken? decoded, out string? failureReason)) - { - metrics.RecordRejectedToken(); - EndpointLog.RejectedToken(logger, route.Oid ?? "(none)", failureReason ?? "unspecified"); - - // No detail in the response: telling a caller which part of a token it got wrong only - // helps a caller who is guessing. - context.Response.StatusCode = StatusCodes.Status403Forbidden; - return false; - } - - if (decoded.Action != expectedAction) - { - metrics.RecordRejectedToken(); - EndpointLog.TokenActionMismatch(logger, decoded.Action, expectedAction); - context.Response.StatusCode = StatusCodes.Status403Forbidden; - return false; - } - - // A token is bound to one object, so it cannot be replayed against a different path. - if (decoded.Oid != route.Oid || decoded.Upstream != route.Upstream) - { - metrics.RecordRejectedToken(); - EndpointLog.RejectedToken(logger, route.Oid ?? "(none)", "token does not match the requested object"); - context.Response.StatusCode = StatusCodes.Status403Forbidden; - return false; - } - - token = decoded; - return true; - } - - private static async Task ServeFromStoreAsync( - HttpContext context, - Stream cached, - long length, - CancellationToken cancellationToken) - { - context.Response.ContentType = OctetStream; - context.Response.ContentLength = length; - - await StreamTee - .CopyAsync(cached, context.Response.Body, null, null, cancellationToken) - .ConfigureAwait(false); - } - - private static void CopyTransferHeaders(HttpResponseMessage response, HttpContext context) - { - context.Response.StatusCode = (int)response.StatusCode; - context.Response.ContentType = response.Content.Headers.ContentType?.ToString() ?? OctetStream; - - if (response.Content.Headers.ContentLength is long length) - { - context.Response.ContentLength = length; - } - - if (response.Content.Headers.ContentRange is not null) - { - context.Response.Headers.ContentRange = response.Content.Headers.ContentRange.ToString(); - } - - if (response.Headers.AcceptRanges.Count > 0) - { - context.Response.Headers.AcceptRanges = string.Join(", ", response.Headers.AcceptRanges); - } - } - - private static async Task CopyResponseAsync( - HttpResponseMessage response, - HttpContext context, - CancellationToken cancellationToken) - { - context.Response.StatusCode = (int)response.StatusCode; - - foreach ((string name, IEnumerable values) in response.Headers) - { - if (!UpstreamRequests.IsHopHeader(name)) - { - context.Response.Headers[name] = values.ToArray(); - } - } - - foreach ((string name, IEnumerable values) in response.Content.Headers) - { - if (!UpstreamRequests.IsHopHeader(name)) - { - context.Response.Headers[name] = values.ToArray(); - } - } - - if (response.StatusCode == HttpStatusCode.NoContent) - { - return; - } - - Stream body = await response.Content - .ReadAsStreamAsync(cancellationToken) - .ConfigureAwait(false); - - await using (body.ConfigureAwait(false)) - { - await body.CopyToAsync(context.Response.Body, cancellationToken).ConfigureAwait(false); } } } diff --git a/GitLfsCache/Endpoints/LockRouteHandler.cs b/GitLfsCache/Endpoints/LockRouteHandler.cs new file mode 100644 index 0000000..daefcbd --- /dev/null +++ b/GitLfsCache/Endpoints/LockRouteHandler.cs @@ -0,0 +1,240 @@ +// Copyright (c) 2023-2026 ktsu-dev contributors + +namespace ktsu.GitLfsCache.Endpoints; + +using System.Text.Json.Nodes; +using ktsu.GitLfsCache.Configuration; +using ktsu.GitLfsCache.Locks; +using ktsu.GitLfsCache.Upstreams; +using Microsoft.AspNetCore.Http; +using Microsoft.Extensions.Options; + +/// +/// Handles the lock routes: listing, batched locking, and the relayed changes that invalidate a +/// snapshot. +/// +/// +/// Separate from the object routes because it shares almost nothing with them. The lock routes need +/// the upstream and the allow-list check, both of which the dispatcher has already done by the time +/// one is called, and nothing about transfer tokens, the object store, or public URL resolution. +/// +/// Answers lock listings from a snapshot. +/// Holds lock snapshots, so a relayed change can invalidate one. +/// Runs the individual calls of a batched lock request. +/// Passes a request upstream when it cannot be terminated here. +/// The configured options. +internal sealed class LockRouteHandler( + LockListService lockLists, + ILockSnapshotStore lockSnapshots, + LockFanOut lockFanOut, + UpstreamRelay relay, + IOptions options) +{ + /// + /// Answers a lock listing, from the snapshot when that is both possible and permitted. + /// + /// The request context. + /// The parsed route. + /// The resolved upstream base URL. + /// Cancels the listing. + /// A task that completes when the response has been written. + public async Task ListAsync( + HttpContext context, + LfsRoute route, + Uri upstreamBase, + CancellationToken cancellationToken) + { + Ensure.NotNull(context); + Ensure.NotNull(route); + + LockSnapshotKey key = new( + route.Upstream, + route.RepositoryPath, + context.Request.Query["refspec"].FirstOrDefault()); + + LockListOutcome outcome = await lockLists + .ResolveAsync(key, upstreamBase, context.Request.Headers.Authorization.ToString(), cancellationToken) + .ConfigureAwait(false); + + switch (outcome.Kind) + { + case LockListOutcomeKind.Refuse: + // Upstream's own refusal, not a proxy interpretation of it. + context.Response.StatusCode = (int)outcome.Status!.Value; + return; + + case LockListOutcomeKind.Serve: + await WriteLockPageAsync(context, outcome.Snapshot!, cancellationToken).ConfigureAwait(false); + return; + + default: + await relay.RelayAsync(context, route, upstreamBase, cancellationToken).ConfigureAwait(false); + return; + } + } + + /// + /// Runs a batched lock or unlock, issuing the individual calls in parallel. + /// + /// + /// A proxy extension, so it is refused rather than relayed when the subsystem is switched off: an + /// upstream has no such endpoint, and relaying would turn a disabled feature into a confusing 404 + /// from the forge instead of a clear one from here. + /// + /// The request context. + /// The parsed route. + /// The resolved upstream base URL. + /// Cancels the fan-out. + /// A task that completes when the response has been written. + public async Task FanOutAsync( + HttpContext context, + LfsRoute route, + Uri upstreamBase, + CancellationToken cancellationToken) + { + Ensure.NotNull(context); + Ensure.NotNull(route); + + if (!options.Value.Locks.Enabled) + { + context.Response.StatusCode = StatusCodes.Status404NotFound; + return; + } + + JsonNode? body; + + try + { + body = await JsonNode.ParseAsync(context.Request.Body, cancellationToken: cancellationToken) + .ConfigureAwait(false); + } + catch (System.Text.Json.JsonException) + { + context.Response.StatusCode = StatusCodes.Status400BadRequest; + return; + } + + if (!LockFanOutRequest.TryParse(body, out LockFanOutRequest? request)) + { + context.Response.StatusCode = StatusCodes.Status400BadRequest; + return; + } + + // Refused outright rather than accepted and throttled part way through, which would leave the + // caller reconciling a partial result they never asked for. + if (request.Targets.Count > options.Value.Locks.MaxFanOutPaths) + { + context.Response.StatusCode = StatusCodes.Status413PayloadTooLarge; + return; + } + + LockSnapshotKey key = new(route.Upstream, route.RepositoryPath, request.Ref); + + JsonObject results = await lockFanOut + .ExecuteAsync( + request, + key, + upstreamBase, + context.Request.Headers.Authorization.ToString(), + cancellationToken) + .ConfigureAwait(false); + + // Always 200 when the request itself was well formed. Partial success is the normal outcome, + // and a transport-level failure would discard the half that worked. + context.Response.StatusCode = StatusCodes.Status200OK; + context.Response.ContentType = UpstreamRequests.LfsMediaType; + + await context.Response + .WriteAsync(results.ToJsonString(), cancellationToken) + .ConfigureAwait(false); + } + + /// + /// Relays a lock creation or release, then drops the snapshot if it took effect. + /// + /// + /// Creation and release are relayed, never terminated, because upstream is the only thing that may + /// grant or release a lock. The snapshot is dropped afterwards so the change this client just made + /// is visible to the next listing rather than waiting out the lifetime. + /// + /// The request context. + /// The parsed route. + /// The resolved upstream base URL. + /// Cancels the relay. + /// A task that completes when the response has been written. + public async Task RelayChangeAsync( + HttpContext context, + LfsRoute route, + Uri upstreamBase, + CancellationToken cancellationToken) + { + Ensure.NotNull(context); + Ensure.NotNull(route); + + await relay.RelayAsync(context, route, upstreamBase, cancellationToken).ConfigureAwait(false); + InvalidateIfChanged(context, route); + } + + /// + /// Writes one page of a snapshot, applying the filters and cursor the client asked for. + /// + private static async Task WriteLockPageAsync( + HttpContext context, + LockSnapshot snapshot, + CancellationToken cancellationToken) + { + IReadOnlyList matches = snapshot.Filter( + context.Request.Query["path"].FirstOrDefault(), + context.Request.Query["id"].FirstOrDefault()); + + int offset = 0; + + // A cursor from a snapshot that has since been replaced restarts the walk rather than being + // applied to a different ordering, which would silently skip or repeat locks. + if (LockCursor.TryDecode(context.Request.Query["cursor"].FirstOrDefault(), out LockCursor? cursor) + && cursor.SnapshotId == snapshot.Id) + { + offset = cursor.Offset; + } + + int? limit = int.TryParse( + context.Request.Query["limit"].FirstOrDefault(), + System.Globalization.NumberStyles.None, + System.Globalization.CultureInfo.InvariantCulture, + out int requested) + ? requested + : null; + + (IReadOnlyList page, int? nextOffset) = LockSnapshot.Paginate(matches, offset, limit); + + JsonObject body = LockListParser.BuildResponse( + page, + nextOffset is int next ? new LockCursor(snapshot.Id, next).Encode() : null); + + context.Response.StatusCode = StatusCodes.Status200OK; + context.Response.ContentType = UpstreamRequests.LfsMediaType; + + await context.Response + .WriteAsync(body.ToJsonString(), cancellationToken) + .ConfigureAwait(false); + } + + /// + /// Drops the snapshot when a relayed lock change actually took effect. + /// + /// + /// Gated on the response status, because invalidating after a refused creation would throw away a + /// perfectly good snapshot every time two people raced for the same file, which is exactly when + /// the cache is under the most load. + /// + private void InvalidateIfChanged(HttpContext context, LfsRoute route) + { + if (context.Response.StatusCode is >= 200 and < 300) + { + lockSnapshots.Invalidate(new LockSnapshotKey( + route.Upstream, + route.RepositoryPath, + context.Request.Query["refspec"].FirstOrDefault())); + } + } +} diff --git a/GitLfsCache/Endpoints/ObjectRouteHandler.cs b/GitLfsCache/Endpoints/ObjectRouteHandler.cs new file mode 100644 index 0000000..106b6d0 --- /dev/null +++ b/GitLfsCache/Endpoints/ObjectRouteHandler.cs @@ -0,0 +1,437 @@ +// Copyright (c) 2023-2026 ktsu-dev contributors + +namespace ktsu.GitLfsCache.Endpoints; + +using System.Diagnostics.CodeAnalysis; +using System.Runtime.CompilerServices; +using System.Text.Json.Nodes; +using ktsu.GitLfsCache.Batch; +using ktsu.GitLfsCache.Configuration; +using ktsu.GitLfsCache.Fetching; +using ktsu.GitLfsCache.Observability; +using ktsu.GitLfsCache.Storage; +using ktsu.GitLfsCache.Tokens; +using ktsu.GitLfsCache.Upstreams; +using Microsoft.AspNetCore.Http; +using Microsoft.Extensions.Logging; +using Microsoft.Extensions.Options; + +/// +/// Handles the object routes: batch, transfer and verify. +/// +/// +/// These three stay together because they genuinely share their front half — validate the transfer +/// token and decide whether the bytes come from the store or from upstream — which is the argument +/// the original single handler was built on. What they never shared was the lock subsystem, which now +/// lives in . +/// +/// Sends requests upstream. +/// Decodes transfer tokens. +/// Rewrites batch responses. +/// The local object store. +/// Keeps concurrent misses to one upstream fetch. +/// Resolves the base URL rewritten hrefs point at. +/// Cache counters. +/// The configured options. +/// Logger. +internal sealed class ObjectRouteHandler( + IUpstreamClient upstreamClient, + IHrefTokenCodec codec, + BatchRewriter rewriter, + IObjectStore store, + IFetchCoalescer coalescer, + PublicUrlResolver publicUrls, + CacheMetrics metrics, + IOptions options, + ILogger logger) +{ + private const string OctetStream = "application/octet-stream"; + private const string TokenQueryParameter = "t"; + + /// + /// Answers a Batch API call, rewriting the hrefs upstream returns to point back here. + /// + /// The request context. + /// The parsed route. + /// The resolved upstream base URL. + /// Cancels the call. + /// A task that completes when the response has been written. + public async Task BatchAsync( + HttpContext context, + LfsRoute route, + Uri upstreamBase, + CancellationToken cancellationToken) + { + Ensure.NotNull(context); + Ensure.NotNull(route); + + using HttpRequestMessage request = UpstreamRequests.BuildBatchRequest( + upstreamBase, + route.RepositoryPath, + context.Request.Body, + context.Request.Headers.Authorization.ToString()); + + using HttpResponseMessage response = await upstreamClient + .SendAsync(request, cancellationToken) + .ConfigureAwait(false); + + // Upstream is the authority on access. A refusal is relayed exactly as it arrived, so the + // client sees upstream's real answer rather than a proxy interpretation of it. + if (!response.IsSuccessStatusCode) + { + await UpstreamRelay.CopyResponseAsync(response, context, cancellationToken).ConfigureAwait(false); + return; + } + + JsonNode? upstreamBody; + + Stream batchBody = await response.Content + .ReadAsStreamAsync(cancellationToken) + .ConfigureAwait(false); + + await using (batchBody.ConfigureAwait(false)) + { + upstreamBody = await JsonNode.ParseAsync(batchBody, cancellationToken: cancellationToken) + .ConfigureAwait(false); + } + + if (upstreamBody is null) + { + context.Response.StatusCode = StatusCodes.Status502BadGateway; + return; + } + + JsonNode rewritten = rewriter.Rewrite(upstreamBody, new BatchRewriteContext + { + Upstream = route.Upstream, + RepositoryPath = route.RepositoryPath, + PublicBaseUrl = publicUrls.Resolve(context.Request), + }); + + context.Response.StatusCode = StatusCodes.Status200OK; + context.Response.ContentType = UpstreamRequests.LfsMediaType; + await context.Response + .WriteAsync(rewritten.ToJsonString(), cancellationToken) + .ConfigureAwait(false); + } + + /// + /// Serves an object, from the store when it is there and from upstream when it is not. + /// + /// The request context. + /// The parsed route. + /// Cancels the transfer. + /// A task that completes when the response has been written. + public async Task DownloadAsync(HttpContext context, LfsRoute route, CancellationToken cancellationToken) + { + Ensure.NotNull(context); + Ensure.NotNull(route); + + if (!TryGetToken(context, route, TokenAction.Download, out HrefToken? token)) + { + return; + } + + string range = context.Request.Headers.Range.ToString(); + + Stream? cached = store.OpenRead(route.Upstream, token.Oid, out long length); + + if (cached is not null) + { + await using (cached.ConfigureAwait(false)) + { + store.Touch(route.Upstream, token.Oid); + metrics.RecordHit(route.Upstream, length); + EndpointLog.ServedFromCache(logger, token.Oid, route.Upstream); + + await ServeFromStoreAsync(context, cached, length, cancellationToken).ConfigureAwait(false); + } + + return; + } + + if (!string.IsNullOrEmpty(range)) + { + // A partial response cannot be stored as a whole object, so the range is forwarded and the + // result streamed straight through. Rare enough not to be worth partial-object bookkeeping. + EndpointLog.RangeRequestNotStored(logger, token.Oid); + await StreamFromUpstreamAsync(context, route, token, range, storeLocally: false, cancellationToken) + .ConfigureAwait(false); + return; + } + + using IFetchTicket ticket = coalescer.Acquire(route.Upstream, token.Oid); + + if (!ticket.IsLeader) + { + metrics.RecordCoalescedWait(route.Upstream); + EndpointLog.WaitingForLeader(logger, token.Oid, route.Upstream); + + bool published = await ticket + .WaitForLeaderAsync(options.Value.Fetch.FollowerTimeout, cancellationToken) + .ConfigureAwait(false); + + long nowLength = 0; + Stream? nowCached = published + ? store.OpenRead(route.Upstream, token.Oid, out nowLength) + : null; + + if (nowCached is not null) + { + 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; + } + + EndpointLog.LeaderDidNotFinish(logger, token.Oid); + } + + bool stored = await StreamFromUpstreamAsync( + context, + route, + token, + range: null, + storeLocally: true, + cancellationToken).ConfigureAwait(false); + + if (ticket.IsLeader) + { + ticket.Complete(stored); + } + } + + /// + /// Sends an upload upstream, keeping a copy on the way through. + /// + /// The request context. + /// The parsed route. + /// Cancels the upload. + /// A task that completes when the response has been written. + public async Task UploadAsync(HttpContext context, LfsRoute route, CancellationToken cancellationToken) + { + Ensure.NotNull(context); + Ensure.NotNull(route); + + if (!TryGetToken(context, route, TokenAction.Upload, out HrefToken? token)) + { + return; + } + + StagingHandle staging = store.OpenStaging(route.Upstream); + + await using (staging.ConfigureAwait(false)) + { + ReadTeeStream teed = new( + context.Request.Body, + staging.Stream, + failure => EndpointLog.StoreSinkFailed(logger, failure, token.Oid)); + + await using ConfiguredAsyncDisposable teedDisposal = teed.ConfigureAwait(false); + + using HttpRequestMessage request = UpstreamRequests.BuildUploadRequest( + token, + teed, + context.Request.ContentLength); + + using HttpResponseMessage response = await upstreamClient + .SendAsync(request, cancellationToken) + .ConfigureAwait(false); + + metrics.RecordUpload(route.Upstream, teed.BytesRead); + + // 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 (await store.PublishAsync(staging, route.Upstream, token.Oid, cancellationToken) + .ConfigureAwait(false)) + { + metrics.RecordStored(route.Upstream); + } + else + { + metrics.RecordVerificationFailure(route.Upstream); + } + } + + await UpstreamRelay.CopyResponseAsync(response, context, cancellationToken).ConfigureAwait(false); + } + } + + /// + /// Passes an upload verification upstream, which is the authority on whether it succeeded. + /// + /// The request context. + /// The parsed route. + /// Cancels the verification. + /// A task that completes when the response has been written. + public async Task VerifyAsync(HttpContext context, LfsRoute route, CancellationToken cancellationToken) + { + Ensure.NotNull(context); + Ensure.NotNull(route); + + if (!TryGetToken(context, route, TokenAction.Verify, out HrefToken? token)) + { + return; + } + + using HttpRequestMessage request = UpstreamRequests.BuildVerifyRequest(token, context.Request.Body); + using HttpResponseMessage response = await upstreamClient + .SendAsync(request, cancellationToken) + .ConfigureAwait(false); + + await UpstreamRelay.CopyResponseAsync(response, context, cancellationToken).ConfigureAwait(false); + } + + private async Task StreamFromUpstreamAsync( + HttpContext context, + LfsRoute route, + HrefToken token, + string? range, + bool storeLocally, + CancellationToken cancellationToken) + { + EndpointLog.FetchingUpstream(logger, token.Oid, route.Upstream); + + using HttpRequestMessage request = UpstreamRequests.BuildObjectRequest(token, range); + using HttpResponseMessage response = await upstreamClient + .SendAsync(request, cancellationToken) + .ConfigureAwait(false); + + if (!response.IsSuccessStatusCode) + { + EndpointLog.UpstreamRefusedTransfer(logger, (int)response.StatusCode, token.Oid); + await UpstreamRelay.CopyResponseAsync(response, context, cancellationToken).ConfigureAwait(false); + return false; + } + + CopyTransferHeaders(response, context); + + Stream upstreamBody = await response.Content + .ReadAsStreamAsync(cancellationToken) + .ConfigureAwait(false); + + await using ConfiguredAsyncDisposable upstreamBodyDisposal = upstreamBody.ConfigureAwait(false); + + if (!storeLocally) + { + long streamed = await StreamTee + .CopyAsync(upstreamBody, context.Response.Body, null, null, cancellationToken) + .ConfigureAwait(false); + + metrics.RecordMiss(route.Upstream, streamed); + return false; + } + + StagingHandle staging = store.OpenStaging(route.Upstream); + + await using (staging.ConfigureAwait(false)) + { + long streamed = await StreamTee.CopyAsync( + upstreamBody, + context.Response.Body, + staging.Stream, + failure => EndpointLog.StoreSinkFailed(logger, failure, token.Oid), + cancellationToken).ConfigureAwait(false); + + metrics.RecordMiss(route.Upstream, streamed); + + bool published = await store + .PublishAsync(staging, route.Upstream, token.Oid, cancellationToken) + .ConfigureAwait(false); + + if (published) + { + metrics.RecordStored(route.Upstream); + } + else + { + metrics.RecordVerificationFailure(route.Upstream); + } + + return published; + } + } + + private bool TryGetToken( + HttpContext context, + LfsRoute route, + string expectedAction, + [NotNullWhen(true)] out HrefToken? token) + { + token = null; + string? encoded = context.Request.Query[TokenQueryParameter]; + + if (!codec.TryDecode(encoded, out HrefToken? decoded, out string? failureReason)) + { + metrics.RecordRejectedToken(); + EndpointLog.RejectedToken(logger, route.Oid ?? "(none)", failureReason ?? "unspecified"); + + // No detail in the response: telling a caller which part of a token it got wrong only + // helps a caller who is guessing. + context.Response.StatusCode = StatusCodes.Status403Forbidden; + return false; + } + + if (decoded.Action != expectedAction) + { + metrics.RecordRejectedToken(); + EndpointLog.TokenActionMismatch(logger, decoded.Action, expectedAction); + context.Response.StatusCode = StatusCodes.Status403Forbidden; + return false; + } + + // A token is bound to one object, so it cannot be replayed against a different path. + if (decoded.Oid != route.Oid || decoded.Upstream != route.Upstream) + { + metrics.RecordRejectedToken(); + EndpointLog.RejectedToken(logger, route.Oid ?? "(none)", "token does not match the requested object"); + context.Response.StatusCode = StatusCodes.Status403Forbidden; + return false; + } + + token = decoded; + return true; + } + + private static async Task ServeFromStoreAsync( + HttpContext context, + Stream cached, + long length, + CancellationToken cancellationToken) + { + context.Response.ContentType = OctetStream; + context.Response.ContentLength = length; + + await StreamTee + .CopyAsync(cached, context.Response.Body, null, null, cancellationToken) + .ConfigureAwait(false); + } + + private static void CopyTransferHeaders(HttpResponseMessage response, HttpContext context) + { + context.Response.StatusCode = (int)response.StatusCode; + context.Response.ContentType = response.Content.Headers.ContentType?.ToString() ?? OctetStream; + + if (response.Content.Headers.ContentLength is long length) + { + context.Response.ContentLength = length; + } + + if (response.Content.Headers.ContentRange is not null) + { + context.Response.Headers.ContentRange = response.Content.Headers.ContentRange.ToString(); + } + + if (response.Headers.AcceptRanges.Count > 0) + { + context.Response.Headers.AcceptRanges = string.Join(", ", response.Headers.AcceptRanges); + } + } +} diff --git a/GitLfsCache/Endpoints/UpstreamRelay.cs b/GitLfsCache/Endpoints/UpstreamRelay.cs new file mode 100644 index 0000000..3896551 --- /dev/null +++ b/GitLfsCache/Endpoints/UpstreamRelay.cs @@ -0,0 +1,111 @@ +// Copyright (c) 2023-2026 ktsu-dev contributors + +namespace ktsu.GitLfsCache.Endpoints; + +using System.Net; +using ktsu.GitLfsCache.Upstreams; +using Microsoft.AspNetCore.Http; +using Microsoft.Extensions.Logging; + +/// +/// Passes a request upstream verbatim and copies the answer back. +/// +/// +/// Shared rather than duplicated because both route groups end up here: the object routes when the +/// store is off or the method is one the proxy does not model, and the lock routes because creation +/// and release are upstream's alone to grant. +/// +/// Sends requests upstream. +/// Logger. +internal sealed class UpstreamRelay(IUpstreamClient upstreamClient, ILogger logger) +{ + /// + /// Relays one request and writes upstream's answer to the response unchanged. + /// + /// The request context. + /// The parsed route. + /// The resolved upstream base URL. + /// Cancels the relay. + /// A task that completes when the response has been written. + public async Task RelayAsync( + HttpContext context, + LfsRoute route, + Uri upstreamBase, + CancellationToken cancellationToken) + { + Ensure.NotNull(context); + Ensure.NotNull(route); + + IEnumerable>> headers = context.Request.Headers + .Select(header => new KeyValuePair>( + header.Key, + header.Value.Where(value => value is not null).Select(value => value!))); + + using HttpRequestMessage request = UpstreamRequests.BuildRelayRequest( + upstreamBase, + context.Request.Method, + route.RelayPath, + context.Request.QueryString.Value ?? string.Empty, + context.Request.Body, + headers); + + using HttpResponseMessage response = await upstreamClient + .SendAsync(request, cancellationToken) + .ConfigureAwait(false); + + EndpointLog.Relayed(logger, context.Request.Method, route.RelayPath, route.Upstream); + await CopyResponseAsync(response, context, cancellationToken).ConfigureAwait(false); + } + + /// + /// Copies an upstream response onto the outgoing response, minus the hop-by-hop headers. + /// + /// + /// Static and shared because every route that does not terminate a request itself ends by handing + /// upstream's own answer back, and they must all drop the same headers to do it. + /// + /// The upstream response. + /// The request context to write to. + /// Cancels the copy. + /// A task that completes when the response has been written. + public static async Task CopyResponseAsync( + HttpResponseMessage response, + HttpContext context, + CancellationToken cancellationToken) + { + Ensure.NotNull(response); + Ensure.NotNull(context); + + context.Response.StatusCode = (int)response.StatusCode; + + foreach ((string name, IEnumerable values) in response.Headers) + { + if (!UpstreamRequests.IsHopHeader(name)) + { + context.Response.Headers[name] = values.ToArray(); + } + } + + foreach ((string name, IEnumerable values) in response.Content.Headers) + { + if (!UpstreamRequests.IsHopHeader(name)) + { + context.Response.Headers[name] = values.ToArray(); + } + } + + if (response.StatusCode == HttpStatusCode.NoContent) + { + return; + } + + Stream body = await response.Content + .ReadAsStreamAsync(cancellationToken) + .ConfigureAwait(false); + + await using (body.ConfigureAwait(false)) + { + await body.CopyToAsync(context.Response.Body, cancellationToken).ConfigureAwait(false); + } + } +} diff --git a/GitLfsCache/GitLfsCacheServiceCollectionExtensions.cs b/GitLfsCache/GitLfsCacheServiceCollectionExtensions.cs index 24b15cd..b9cf8f3 100644 --- a/GitLfsCache/GitLfsCacheServiceCollectionExtensions.cs +++ b/GitLfsCache/GitLfsCacheServiceCollectionExtensions.cs @@ -98,6 +98,9 @@ public static IServiceCollection AddGitLfsCache( return readiness; }); services.AddSingleton(); + services.AddSingleton(); + services.AddSingleton(); + services.AddSingleton(); services.AddSingleton(); services.AddMetrics();