feat(mesh): config-serve auth interceptor + dedicated h2c listener (fail-closed bearer)

Phase 3 Task 3. Central serves the deployment artifact over a dedicated h2c
Kestrel listener (ConfigServe:GrpcListenPort), gated by the ConfigServeAuthInterceptor
(shared bearer key, FixedTimeEquals, fail-closed, path-scoped to
/deployment_artifact.v1.DeploymentArtifactService/). The listener block is merged
with the LocalDb-sync listener so a fused admin+driver node re-applies its existing
HTTP surface exactly ONCE and adds both dedicated h2c ports on top — double-binding
would throw 'address already in use'; zero re-binding would silently unbind the
AdminUI behind Traefik. AddGrpc now registers under hasDriver || hasAdmin with both
path-scoped interceptors sharing one pipeline.

Claude-Session: https://claude.ai/code/session_01GASWkNEi68FSCtvr6rLoEW
This commit is contained in:
Joseph Doherty
2026-07-22 19:22:55 -04:00
parent b28c071bec
commit 4d88f67c0d
4 changed files with 474 additions and 26 deletions
@@ -0,0 +1,152 @@
using Grpc.Core;
using Microsoft.Extensions.Logging.Abstractions;
using Microsoft.Extensions.Options;
using Shouldly;
using Xunit;
using ZB.MOM.WW.OtOpcUa.Cluster;
using ZB.MOM.WW.OtOpcUa.Host.Configuration;
namespace ZB.MOM.WW.OtOpcUa.Host.IntegrationTests;
/// <summary>
/// Per-cluster mesh Phase 3 (Task 3) — the config-serve artifact endpoint's inbound gate.
/// </summary>
/// <remarks>
/// The <c>DeploymentArtifactService</c> streams a driver node the exact configuration it will
/// boot on. Anything able to reach central's serve port could otherwise pull the whole deployed
/// config, so the endpoint is gated by a shared bearer key, fail-closed, constant-time — the
/// same contract as <see cref="LocalDbSyncAuthInterceptor"/>, scoped to a different service path.
/// </remarks>
public sealed class ConfigServeAuthInterceptorTests
{
private const string FetchMethod = "/deployment_artifact.v1.DeploymentArtifactService/Fetch";
private const string OtherMethod = "/localdb_sync.v1.LocalDbSync/Sync";
private static ConfigServeAuthInterceptor CreateInterceptor(string? apiKey)
=> new(
Options.Create(new ConfigServeOptions { ApiKey = apiKey ?? string.Empty }),
NullLogger<ConfigServeAuthInterceptor>.Instance);
private static ServerCallContext CreateContext(string method, string? authorizationHeader)
{
var headers = new Metadata();
if (authorizationHeader is not null)
headers.Add("authorization", authorizationHeader);
return new FakeServerCallContext(method, headers);
}
/// <summary>Invokes the interceptor's server-streaming path (how <c>Fetch</c> actually runs).</summary>
private static Task Invoke(ConfigServeAuthInterceptor interceptor, ServerCallContext context)
=> interceptor.ServerStreamingServerHandler<string, string>(
"request", responseStream: null!, context, (_, _, _) => Task.CompletedTask);
[Fact]
public async Task NonServeMethod_PassesThrough_EvenWithNoKeyConfigured()
{
// The interceptor shares the AddGrpc pipeline with the LocalDb sync interceptor, so it sees
// every call. It must be scoped strictly to the artifact service — a sync call is not its
// business.
var interceptor = CreateInterceptor(apiKey: null);
var context = CreateContext(OtherMethod, authorizationHeader: null);
await Invoke(interceptor, context); // no throw
}
[Fact]
public async Task FetchMethod_WithNoKeyConfigured_IsDenied_EvenWithABearerToken()
{
// Fail-closed. "No key configured" is the DEFAULT — treating it as "no auth required" would
// expose the deployed config on exactly the shape most nodes ship with. A token must not help.
var interceptor = CreateInterceptor(apiKey: null);
var context = CreateContext(FetchMethod, "Bearer anything-at-all");
var ex = await Should.ThrowAsync<RpcException>(() => Invoke(interceptor, context));
ex.StatusCode.ShouldBe(StatusCode.Unauthenticated);
}
[Fact]
public async Task FetchMethod_WithNoBearerToken_IsDenied()
{
var interceptor = CreateInterceptor("the-shared-key");
var context = CreateContext(FetchMethod, authorizationHeader: null);
var ex = await Should.ThrowAsync<RpcException>(() => Invoke(interceptor, context));
ex.StatusCode.ShouldBe(StatusCode.Unauthenticated);
}
[Fact]
public async Task FetchMethod_WithWrongBearerToken_IsDenied()
{
var interceptor = CreateInterceptor("the-shared-key");
var context = CreateContext(FetchMethod, "Bearer the-wrong-key");
var ex = await Should.ThrowAsync<RpcException>(() => Invoke(interceptor, context));
ex.StatusCode.ShouldBe(StatusCode.Unauthenticated);
}
[Fact]
public async Task FetchMethod_WithCorrectBearerToken_PassesThrough()
{
var interceptor = CreateInterceptor("the-shared-key");
var context = CreateContext(FetchMethod, "Bearer the-shared-key");
await Invoke(interceptor, context); // no throw
}
[Fact]
public async Task FetchMethod_TokenComparison_IsNotAPrefixMatch()
{
// A prefix/StartsWith comparison would accept a truncated key and make the secret
// recoverable one character at a time.
var interceptor = CreateInterceptor("the-shared-key");
var context = CreateContext(FetchMethod, "Bearer the-shared-ke");
var ex = await Should.ThrowAsync<RpcException>(() => Invoke(interceptor, context));
ex.StatusCode.ShouldBe(StatusCode.Unauthenticated);
}
[Fact]
public async Task UnaryPath_IsAlsoGated()
{
// Fetch is server-streaming today, but gating only that path would leave any future unary
// method on the same service open. Gate every handler shape.
var interceptor = CreateInterceptor("the-shared-key");
var context = CreateContext(FetchMethod, "Bearer the-wrong-key");
var ex = await Should.ThrowAsync<RpcException>(() =>
interceptor.UnaryServerHandler<string, string>(
"request", context, (_, _) => Task.FromResult("ok")));
ex.StatusCode.ShouldBe(StatusCode.Unauthenticated);
}
/// <summary>
/// Minimal <see cref="ServerCallContext"/> carrying just a method name and request headers —
/// the only two things the interceptor reads. Hand-rolled because
/// <c>Grpc.Core.Testing.TestServerCallContext</c> ships only in the retired native package.
/// </summary>
private sealed class FakeServerCallContext(string method, Metadata requestHeaders)
: ServerCallContext
{
protected override string MethodCore => method;
protected override string HostCore => "localhost";
protected override string PeerCore => "ipv4:127.0.0.1:12345";
protected override DateTime DeadlineCore => DateTime.UtcNow.AddMinutes(1);
protected override Metadata RequestHeadersCore => requestHeaders;
protected override CancellationToken CancellationTokenCore => CancellationToken.None;
protected override Metadata ResponseTrailersCore { get; } = [];
protected override Status StatusCore { get; set; }
protected override WriteOptions? WriteOptionsCore { get; set; }
protected override AuthContext AuthContextCore { get; } =
new(null, new Dictionary<string, List<AuthProperty>>());
protected override ContextPropagationToken CreatePropagationTokenCore(
ContextPropagationOptions? options)
=> throw new NotSupportedException();
protected override Task WriteResponseHeadersAsyncCore(Metadata responseHeaders)
=> Task.CompletedTask;
}
}
@@ -0,0 +1,99 @@
using System.Net;
using Microsoft.AspNetCore.Builder;
using Microsoft.AspNetCore.Server.Kestrel.Core;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Hosting;
using Shouldly;
using Xunit;
using ZB.MOM.WW.OtOpcUa.Host.Configuration;
namespace ZB.MOM.WW.OtOpcUa.Host.IntegrationTests;
/// <summary>
/// Per-cluster mesh Phase 3 (Task 3) — the coexistence of the ConfigServe artifact h2c listener
/// with the LocalDb-sync h2c listener on a fused admin+driver node.
/// </summary>
/// <remarks>
/// <para>
/// <b>Why this is the load-bearing test.</b> A fused central node can carry BOTH dedicated
/// h2c ports (LocalDb sync + ConfigServe). Any explicit Kestrel <c>Listen*</c> makes Kestrel
/// discard <c>ASPNETCORE_URLS</c> wholesale, so the existing HTTP surface must be re-bound
/// exactly ONCE alongside the two dedicated ports. Re-binding it twice (one block per h2c
/// port) throws "address already in use"; re-binding it zero times silently unbinds the
/// AdminUI behind Traefik. This test pins that all three answer simultaneously.
/// </para>
/// <para>
/// Like <see cref="LocalDbSyncListenerTests"/>, this drives a minimal <c>WebApplication</c>
/// shaped exactly like <c>Program.cs</c>'s merged listener block rather than booting the real
/// host (which needs SQL, an Akka mesh and LDAP). What is under test is the Kestrel binding
/// contract, fully reproduced here; the real end-to-end fetch is the Task 8 boundary test.
/// </para>
/// </remarks>
public sealed class ConfigServeListenerTests
{
[Fact]
public async Task BothDedicatedH2cPorts_AndTheHttpSurface_AnswerSimultaneously()
{
var httpPort = GetFreePort();
var syncPort = GetFreePort();
var configServePort = GetFreePort();
var builder = WebApplication.CreateBuilder();
builder.Environment.ApplicationName = typeof(ConfigServeListenerTests).Assembly.GetName().Name!;
builder.Services.AddGrpc();
// Exactly the shape Program.cs's merged block applies: compute the existing surface ONCE,
// apply it ONCE, then add each configured dedicated h2c port.
var existingBindings = KestrelHttpBinding.Parse($"http://localhost:{httpPort}");
builder.WebHost.ConfigureKestrel(kestrel =>
{
foreach (var binding in existingBindings)
binding.Apply(kestrel);
kestrel.ListenAnyIP(syncPort, o => o.Protocols = HttpProtocols.Http2);
kestrel.ListenAnyIP(configServePort, o => o.Protocols = HttpProtocols.Http2);
});
await using var app = builder.Build();
app.MapGet("/healthz", () => "ok");
await app.StartAsync(TestContext.Current.CancellationToken);
try
{
// The re-bound HTTP/1.1 surface still answers (the AdminUI/Traefik guarantee).
using var http1 = new HttpClient();
(await http1.GetStringAsync(
$"http://localhost:{httpPort}/healthz", TestContext.Current.CancellationToken))
.ShouldBe("ok");
// Both dedicated ports negotiate prior-knowledge h2c. A 404 is a fine result — it proves
// the HTTP/2 connection was established and routed, which is all that is in question.
using var http2 = new HttpClient
{
DefaultRequestVersion = HttpVersion.Version20,
DefaultVersionPolicy = HttpVersionPolicy.RequestVersionExact,
};
using var syncResponse = await http2.GetAsync(
$"http://localhost:{syncPort}/", TestContext.Current.CancellationToken);
syncResponse.Version.ShouldBe(HttpVersion.Version20);
using var configServeResponse = await http2.GetAsync(
$"http://localhost:{configServePort}/", TestContext.Current.CancellationToken);
configServeResponse.Version.ShouldBe(HttpVersion.Version20);
}
finally
{
await app.StopAsync(TestContext.Current.CancellationToken);
}
}
private static int GetFreePort()
{
using var listener = new System.Net.Sockets.TcpListener(IPAddress.Loopback, 0);
listener.Start();
var port = ((IPEndPoint)listener.LocalEndpoint).Port;
listener.Stop();
return port;
}
}