Skip to content

Commit 0d1b2be

Browse files
authored
fix(grpc): preserve forwarded authentication failures (#510)
- Forwarded writes must report authentication failures consistently with direct gRPC requests. Signed-off-by: Yordis Prieto <yordis.prieto@gmail.com>
1 parent d30db04 commit 0d1b2be

3 files changed

Lines changed: 153 additions & 0 deletions

File tree

Lines changed: 141 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,141 @@
1+
using System;
2+
using System.Collections.Generic;
3+
using System.Reflection;
4+
using System.Security.Claims;
5+
using System.Threading;
6+
using System.Threading.Tasks;
7+
using EventStore.Client.Streams;
8+
using EventStore.Core.Authorization;
9+
using EventStore.Core.Bus;
10+
using EventStore.Core.Messages;
11+
using EventStore.Core.Messaging;
12+
using EventStore.Core.Services.Transport.Grpc;
13+
using Google.Protobuf;
14+
using Grpc.Core;
15+
using Microsoft.AspNetCore.Http;
16+
using NUnit.Framework;
17+
using GrpcMetadata = EventStore.Core.Services.Transport.Grpc.Constants.Metadata;
18+
using GrpcStreams = EventStore.Client.Streams.Streams;
19+
20+
namespace EventStore.Core.Tests.Services.Transport.Grpc.StreamsTests;
21+
22+
[TestFixture]
23+
public class ForwardedAuthenticationTests
24+
{
25+
[Test]
26+
public void append_reports_forwarded_authentication_failure()
27+
{
28+
var service = CreateService(new AuthenticationFailurePublisher());
29+
var requests = new EnumerableStreamReader<AppendReq>([
30+
new AppendReq {
31+
Options = new AppendReq.Types.Options {
32+
NoStream = new(),
33+
StreamIdentifier = "forwarded-auth-append"
34+
}
35+
},
36+
new AppendReq {
37+
ProposedMessage = new AppendReq.Types.ProposedMessage {
38+
Id = Uuid.NewUuid().ToDto(),
39+
Metadata = {
40+
[GrpcMetadata.Type] = "test",
41+
[GrpcMetadata.ContentType] = GrpcMetadata.ContentTypes.ApplicationJson
42+
},
43+
Data = ByteString.CopyFromUtf8("{}")
44+
}
45+
}
46+
]);
47+
48+
var exception = Assert.ThrowsAsync<RpcException>(() => service.Append(requests, new TestServerCallContext()));
49+
50+
Assert.That(exception!.StatusCode, Is.EqualTo(StatusCode.Unauthenticated));
51+
Assert.That(exception.Status.Detail, Does.Contain("forwarding denied"));
52+
}
53+
54+
[Test]
55+
public void delete_reports_forwarded_authentication_failure()
56+
{
57+
var service = CreateService(new AuthenticationFailurePublisher());
58+
var request = new DeleteReq
59+
{
60+
Options = new DeleteReq.Types.Options
61+
{
62+
NoStream = new(),
63+
StreamIdentifier = "forwarded-auth-delete"
64+
}
65+
};
66+
67+
var exception = Assert.ThrowsAsync<RpcException>(() => service.Delete(request, new TestServerCallContext()));
68+
69+
Assert.That(exception!.StatusCode, Is.EqualTo(StatusCode.Unauthenticated));
70+
Assert.That(exception.Status.Detail, Does.Contain("forwarding denied"));
71+
}
72+
73+
private static GrpcStreams.StreamsBase CreateService(IPublisher publisher)
74+
{
75+
var type = typeof(GrpcTrackers).Assembly.GetType(
76+
"EventStore.Core.Services.Transport.Grpc.Streams`1", throwOnError: true)!
77+
.MakeGenericType(typeof(string));
78+
return (GrpcStreams.StreamsBase)Activator.CreateInstance(type,
79+
BindingFlags.Instance | BindingFlags.Public | BindingFlags.NonPublic,
80+
binder: null,
81+
args: [publisher, 1024, TimeSpan.FromSeconds(1), null, new GrpcTrackers(), new PassthroughAuthorizationProvider()],
82+
culture: null)!;
83+
}
84+
85+
private sealed class AuthenticationFailurePublisher : IPublisher
86+
{
87+
public void Publish(Message message)
88+
{
89+
if (message is not ClientMessage.WriteRequestMessage request)
90+
{
91+
throw new InvalidOperationException($"Unexpected message {message.GetType().Name}");
92+
}
93+
94+
request.Envelope.ReplyWith(new ClientMessage.NotAuthenticated(request.CorrelationId, "forwarding denied"));
95+
}
96+
}
97+
98+
private sealed class EnumerableStreamReader<T>(IEnumerable<T> values) : IAsyncStreamReader<T>
99+
{
100+
private readonly IEnumerator<T> _values = values.GetEnumerator();
101+
public T Current { get; private set; } = default!;
102+
103+
public Task<bool> MoveNext(CancellationToken cancellationToken)
104+
{
105+
if (!_values.MoveNext())
106+
{
107+
return Task.FromResult(false);
108+
}
109+
110+
Current = _values.Current;
111+
return Task.FromResult(true);
112+
}
113+
}
114+
115+
private sealed class TestServerCallContext : ServerCallContext
116+
{
117+
public TestServerCallContext()
118+
{
119+
UserStateCore["__HttpContext"] = new DefaultHttpContext
120+
{
121+
User = new ClaimsPrincipal(new ClaimsIdentity())
122+
};
123+
}
124+
125+
protected override string MethodCore => "/event_store.client.streams.Streams/Append";
126+
protected override string HostCore => "localhost";
127+
protected override string PeerCore => "ipv4:127.0.0.1:2113";
128+
protected override DateTime DeadlineCore => DateTime.MaxValue;
129+
protected override Metadata RequestHeadersCore { get; } = new();
130+
protected override CancellationToken CancellationTokenCore => CancellationToken.None;
131+
protected override Metadata ResponseTrailersCore { get; } = new();
132+
protected override Status StatusCore { get; set; }
133+
protected override WriteOptions WriteOptionsCore { get; set; }
134+
protected override AuthContext AuthContextCore { get; } =
135+
new(string.Empty, new Dictionary<string, List<AuthProperty>>());
136+
protected override IDictionary<object, object> UserStateCore { get; } = new Dictionary<object, object>();
137+
protected override Task WriteResponseHeadersAsyncCore(Metadata responseHeaders) => Task.CompletedTask;
138+
protected override ContextPropagationToken CreatePropagationTokenCore(ContextPropagationOptions options) =>
139+
throw new NotSupportedException();
140+
}
141+
}

‎src/EventStore.Core/Services/Transport/Grpc/Streams.Append.cs‎

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -122,6 +122,12 @@ void HandleWriteEventsCompleted(Message message)
122122
appendResponseSource.TrySetException(ex);
123123
return;
124124
}
125+
if (message is ClientMessage.NotAuthenticated notAuthenticated)
126+
{
127+
appendResponseSource.TrySetException(new RpcException(
128+
new Status(StatusCode.Unauthenticated, notAuthenticated.Reason)));
129+
return;
130+
}
125131

126132
if (!(message is ClientMessage.WriteEventsCompleted completed))
127133
{

‎src/EventStore.Core/Services/Transport/Grpc/Streams.Delete.cs‎

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -148,6 +148,12 @@ void HandleStreamDeletedCompleted(Message message)
148148
deleteResponseSource.TrySetException(ex);
149149
return;
150150
}
151+
if (message is ClientMessage.NotAuthenticated notAuthenticated)
152+
{
153+
deleteResponseSource.TrySetException(new RpcException(
154+
new Status(StatusCode.Unauthenticated, notAuthenticated.Reason)));
155+
return;
156+
}
151157

152158
if (message is not ClientMessage.DeleteStreamCompleted completed)
153159
{

0 commit comments

Comments
 (0)