Skip to content

Commit ee40129

Browse files
committed
feat(cluster): preserve isolated replication over gRPC
Signed-off-by: Yordis Prieto <yordis.prieto@gmail.com>
1 parent c0537ba commit ee40129

83 files changed

Lines changed: 1060 additions & 2577 deletions

File tree

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

‎src/EventStore.ClusterNode/Components/Pages/Cluster.razor‎

Lines changed: 4 additions & 23 deletions
Original file line numberDiff line numberDiff line change
@@ -43,16 +43,15 @@
4343
<th class="px-5 py-4">Status</th>
4444
<th class="px-5 py-4">Timestamp (UTC)</th>
4545
<th class="px-5 py-4">Checkpoints</th>
46-
<th class="px-5 py-4">TCP</th>
47-
<th class="px-5 py-4">HTTP</th>
46+
<th class="px-5 py-4">HTTP / gRPC</th>
4847
<th class="px-5 py-4 text-right">Actions</th>
4948
</tr>
5049
</thead>
5150
<tbody class="divide-y divide-es-ink/10">
5251
@if (!ClusterMembers.Any())
5352
{
5453
<tr>
55-
<td class="px-5 py-4 text-es-muted" colspan="7">@ClusterEmptyMessage</td>
54+
<td class="px-5 py-4 text-es-muted" colspan="6">@ClusterEmptyMessage</td>
5655
</tr>
5756
}
5857
else
@@ -77,10 +76,6 @@
7776
<p class="mt-1">@EpochLabel(member)</p>
7877
}
7978
</td>
80-
<td class="px-5 py-4 font-mono text-xs text-es-muted">
81-
<p>Internal: @InternalTcpEndpoint(member)</p>
82-
<p class="mt-1">External: @ExternalTcpEndpoint(member)</p>
83-
</td>
8479
<td class="px-5 py-4 font-mono text-xs text-es-muted">@HttpEndpoint(member)</td>
8580
<td class="px-5 py-4 text-right">
8681
<div class="flex flex-wrap justify-end gap-2">
@@ -307,7 +302,7 @@
307302
<section class="mt-5 grid gap-4 md:grid-cols-2 xl:grid-cols-4">
308303
<SurfaceCard Eyebrow="Explore" Title="Navigator" Description="Jump into streams, subscriptions, projections, user management, and browser tools from one place." Href="/ui/navigator" LinkText="Open" />
309304
<SurfaceCard Eyebrow="Run" Title="Operations" Description="Reach privileged actions for scavenging, shutdown, reload, and node-level workflows." Href="/ui/operations" LinkText="Open" />
310-
<SurfaceCard Eyebrow="Watch" Title="Observability" Description="Inspect queues, replication, TCP, metrics, and health through a curated operator map." Href="/ui/observability" LinkText="Open" />
305+
<SurfaceCard Eyebrow="Watch" Title="Observability" Description="Inspect queues, replication, metrics, and health through a curated operator map." Href="/ui/observability" LinkText="Open" />
311306
<SurfaceCard Eyebrow="Review" Title="Configuration" Description="Find runtime information, loaded options, and subsystem metadata." Href="/ui/configuration" LinkText="Open" />
312307
</section>
313308
</div>
@@ -370,18 +365,14 @@
370365
.Append("Snapshot taken at ")
371366
.Append(TimestampLabel(ClusterReadAt ?? DateTime.UtcNow))
372367
.AppendLine()
373-
.Append(PadRight("Internal Tcp", 31)).Append(' ')
374-
.Append(PadRight("External Tcp", 31)).Append(' ')
375-
.Append(PadRight("Http", 23)).Append(' ')
368+
.Append(PadRight("HTTP / gRPC", 23)).Append(' ')
376369
.Append(PadRight("Status", 11)).Append(' ')
377370
.Append(PadRight("State", 18)).Append(' ')
378371
.Append(PadRight("Timestamp (UTC)", 19)).Append(" Checkpoints");
379372

380373
foreach (var member in ClusterMembers)
381374
{
382375
builder.AppendLine()
383-
.Append(PadRight(InternalTcpEndpoint(member), 31)).Append(' ')
384-
.Append(PadRight(ExternalTcpEndpoint(member), 31)).Append(' ')
385376
.Append(PadRight(HttpEndpoint(member), 23)).Append(' ')
386377
.Append(PadRight(MemberStatus(member), 11)).Append(' ')
387378
.Append(PadRight(member.State.ToString(), 18)).Append(' ')
@@ -500,16 +491,6 @@
500491
private static string MemberStatus(ClientClusterInfo.ClientMemberInfo member) =>
501492
member.IsAlive ? "Alive" : "Unreachable";
502493

503-
private static string InternalTcpEndpoint(ClientClusterInfo.ClientMemberInfo member) =>
504-
Endpoint(
505-
member.InternalTcpIp,
506-
member.InternalSecureTcpPort != 0 ? member.InternalSecureTcpPort : member.InternalTcpPort);
507-
508-
private static string ExternalTcpEndpoint(ClientClusterInfo.ClientMemberInfo member) =>
509-
Endpoint(
510-
member.ExternalTcpIp,
511-
member.ExternalSecureTcpPort != 0 ? member.ExternalSecureTcpPort : member.ExternalTcpPort);
512-
513494
private static string HttpEndpoint(ClientClusterInfo.ClientMemberInfo member) =>
514495
Endpoint(member.HttpEndPointIp, member.HttpEndPointPort);
515496

‎src/EventStore.ClusterNode/Components/Services/ClusterStatusService.cs‎

Lines changed: 3 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -171,7 +171,7 @@ private ClusterReplicaRow ParseReplicaRow(
171171
: Guid.Empty;
172172
var totalBytesSent = row.TotalBytesSent;
173173
var previousRow = _previousReplicas.GetValueOrDefault(connectionId);
174-
var replicaNode = FindMemberByInternalEndpoint(members, row.SubscriptionEndpoint);
174+
var replicaNode = FindMemberByEndpoint(members, row.SubscriptionEndpoint);
175175
var isCatchingUp = replicaNode?.State == VNodeState.CatchingUp;
176176
var catchupStartTime = now;
177177
var catchupStartBytesSent = totalBytesSent;
@@ -206,24 +206,19 @@ private ClusterReplicaRow ParseReplicaRow(
206206
private ClaimsPrincipal CurrentUser =>
207207
httpContextAccessor.HttpContext?.User ?? new ClaimsPrincipal(new ClaimsIdentity());
208208

209-
private static ClientClusterInfo.ClientMemberInfo FindMemberByInternalEndpoint(
209+
private static ClientClusterInfo.ClientMemberInfo FindMemberByEndpoint(
210210
IReadOnlyList<ClientClusterInfo.ClientMemberInfo> members,
211211
string endpoint)
212212
{
213213
var cleaned = endpoint.Replace("Unspecified/", "", StringComparison.OrdinalIgnoreCase);
214-
return members.FirstOrDefault(x => string.Equals(InternalTcpEndpoint(x), cleaned, StringComparison.OrdinalIgnoreCase));
214+
return members.FirstOrDefault(x => string.Equals(HttpEndpoint(x), cleaned, StringComparison.OrdinalIgnoreCase));
215215
}
216216

217217
private static Uri BuildLeaderAddress(
218218
HttpRequest request,
219219
ClientClusterInfo.ClientMemberInfo leader) =>
220220
new UriBuilder(request.Scheme, leader.HttpEndPointIp, leader.HttpEndPointPort).Uri;
221221

222-
private static string InternalTcpEndpoint(ClientClusterInfo.ClientMemberInfo member) =>
223-
Endpoint(
224-
member.InternalTcpIp,
225-
member.InternalSecureTcpPort != 0 ? member.InternalSecureTcpPort : member.InternalTcpPort);
226-
227222
private static string HttpEndpoint(ClientClusterInfo.ClientMemberInfo member) =>
228223
Endpoint(member.HttpEndPointIp, member.HttpEndPointPort);
229224

Lines changed: 278 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,278 @@
1+
using System;
2+
using System.Buffers;
3+
using System.Collections.Concurrent;
4+
using System.Collections.Generic;
5+
using System.IO.Pipelines;
6+
using System.Linq;
7+
using System.Threading;
8+
using System.Threading.Tasks;
9+
using Microsoft.AspNetCore.Connections;
10+
11+
namespace EventStore.ClusterNode.Components.Services;
12+
13+
public sealed class NodeConnectionTracker
14+
{
15+
private readonly ConcurrentDictionary<string, NodeConnectionState> _connections = new();
16+
17+
public IReadOnlyList<NodeConnectionSnapshot> Snapshot() =>
18+
_connections.Values.Select(x => x.Snapshot())
19+
.OrderBy(x => x.RemoteEndPoint, StringComparer.OrdinalIgnoreCase)
20+
.ThenBy(x => x.ConnectionId, StringComparer.Ordinal)
21+
.ToArray();
22+
23+
public async Task Track(ConnectionContext context, ConnectionDelegate next, bool isTls)
24+
{
25+
var state = new NodeConnectionState(
26+
context.ConnectionId,
27+
context.RemoteEndPoint?.ToString() ?? "",
28+
context.LocalEndPoint?.ToString() ?? "",
29+
isTls,
30+
DateTimeOffset.UtcNow);
31+
_connections[context.ConnectionId] = state;
32+
context.Transport = new CountingDuplexPipe(context.Transport, state);
33+
34+
try
35+
{
36+
await next(context);
37+
}
38+
finally
39+
{
40+
_connections.TryRemove(context.ConnectionId, out _);
41+
}
42+
}
43+
44+
public void ObserveRequest(
45+
string connectionId,
46+
string protocol,
47+
bool isGrpc,
48+
string connectionName,
49+
string userAgent)
50+
{
51+
if (_connections.TryGetValue(connectionId, out var connection))
52+
connection.ObserveRequest(protocol, isGrpc, connectionName, userAgent);
53+
}
54+
}
55+
56+
public sealed record NodeConnectionSnapshot(
57+
string ConnectionId,
58+
string RemoteEndPoint,
59+
string LocalEndPoint,
60+
string ClientName,
61+
string Application,
62+
string Protocol,
63+
bool IsTls,
64+
DateTimeOffset ConnectedAt,
65+
long TotalBytesSent,
66+
long TotalBytesReceived,
67+
long PendingSendBytes,
68+
long PendingReceivedBytes);
69+
70+
internal sealed class NodeConnectionState
71+
{
72+
private readonly object _metadataLock = new();
73+
private string _clientName = "";
74+
private bool _hasExplicitConnectionName;
75+
private bool _hasGrpcRequests;
76+
private bool _hasHttpRequests;
77+
private long _pendingReceivedBytes;
78+
private long _pendingSendBytes;
79+
private string _protocol = "";
80+
private long _totalBytesReceived;
81+
private long _totalBytesSent;
82+
83+
public NodeConnectionState(
84+
string connectionId,
85+
string remoteEndPoint,
86+
string localEndPoint,
87+
bool isTls,
88+
DateTimeOffset connectedAt)
89+
{
90+
ConnectionId = connectionId;
91+
RemoteEndPoint = remoteEndPoint;
92+
LocalEndPoint = localEndPoint;
93+
IsTls = isTls;
94+
ConnectedAt = connectedAt;
95+
}
96+
97+
private string ConnectionId { get; }
98+
private string RemoteEndPoint { get; }
99+
private string LocalEndPoint { get; }
100+
private bool IsTls { get; }
101+
private DateTimeOffset ConnectedAt { get; }
102+
103+
public void Received(long bytes, long pendingBytes)
104+
{
105+
Interlocked.Add(ref _totalBytesReceived, bytes);
106+
Interlocked.Exchange(ref _pendingReceivedBytes, pendingBytes);
107+
}
108+
109+
public void Reading(long pendingBytes) =>
110+
Interlocked.Exchange(ref _pendingReceivedBytes, pendingBytes);
111+
112+
public void Sending(int bytes)
113+
{
114+
Interlocked.Add(ref _totalBytesSent, bytes);
115+
Interlocked.Add(ref _pendingSendBytes, bytes);
116+
}
117+
118+
public void Sent() => Interlocked.Exchange(ref _pendingSendBytes, 0);
119+
120+
public void ObserveRequest(
121+
string protocol,
122+
bool isGrpc,
123+
string connectionName,
124+
string userAgent)
125+
{
126+
lock (_metadataLock)
127+
{
128+
_protocol = Merge(_protocol, protocol);
129+
_hasGrpcRequests |= isGrpc;
130+
_hasHttpRequests |= !isGrpc;
131+
132+
if (!string.IsNullOrWhiteSpace(connectionName))
133+
{
134+
_clientName = connectionName;
135+
_hasExplicitConnectionName = true;
136+
}
137+
else if (!_hasExplicitConnectionName && !string.IsNullOrWhiteSpace(userAgent))
138+
{
139+
_clientName = userAgent;
140+
}
141+
}
142+
}
143+
144+
public NodeConnectionSnapshot Snapshot()
145+
{
146+
lock (_metadataLock)
147+
{
148+
return new(
149+
ConnectionId,
150+
RemoteEndPoint,
151+
LocalEndPoint,
152+
_clientName,
153+
ApplicationLabel(),
154+
_protocol,
155+
IsTls,
156+
ConnectedAt,
157+
Interlocked.Read(ref _totalBytesSent),
158+
Interlocked.Read(ref _totalBytesReceived),
159+
Interlocked.Read(ref _pendingSendBytes),
160+
Interlocked.Read(ref _pendingReceivedBytes));
161+
}
162+
}
163+
164+
private string ApplicationLabel() => (_hasHttpRequests, _hasGrpcRequests) switch
165+
{
166+
(true, true) => "HTTP and gRPC",
167+
(false, true) => "gRPC",
168+
(true, false) => "HTTP",
169+
_ => "Awaiting request"
170+
};
171+
172+
private static string Merge(string current, string observed)
173+
{
174+
if (string.IsNullOrWhiteSpace(observed) || current == observed)
175+
return current;
176+
return string.IsNullOrWhiteSpace(current) ? observed : "Mixed";
177+
}
178+
}
179+
180+
internal sealed class CountingDuplexPipe : IDuplexPipe
181+
{
182+
public CountingDuplexPipe(IDuplexPipe inner, NodeConnectionState state)
183+
{
184+
Input = new CountingPipeReader(inner.Input, state);
185+
Output = new CountingPipeWriter(inner.Output, state);
186+
}
187+
188+
public PipeReader Input { get; }
189+
public PipeWriter Output { get; }
190+
}
191+
192+
internal sealed class CountingPipeReader : PipeReader
193+
{
194+
private readonly PipeReader _inner;
195+
private readonly NodeConnectionState _state;
196+
private ReadOnlySequence<byte> _currentBuffer;
197+
198+
public CountingPipeReader(PipeReader inner, NodeConnectionState state)
199+
{
200+
_inner = inner;
201+
_state = state;
202+
}
203+
204+
public override void AdvanceTo(SequencePosition consumed) => AdvanceTo(consumed, consumed);
205+
206+
public override void AdvanceTo(SequencePosition consumed, SequencePosition examined)
207+
{
208+
var consumedBytes = _currentBuffer.IsEmpty ? 0 : _currentBuffer.Slice(0, consumed).Length;
209+
var pendingBytes = _currentBuffer.IsEmpty ? 0 : _currentBuffer.Slice(consumed).Length;
210+
_state.Received(consumedBytes, pendingBytes);
211+
_currentBuffer = default;
212+
_inner.AdvanceTo(consumed, examined);
213+
}
214+
215+
public override void CancelPendingRead() => _inner.CancelPendingRead();
216+
217+
public override void Complete(Exception exception = null) => _inner.Complete(exception);
218+
219+
public override ValueTask CompleteAsync(Exception exception = null) => _inner.CompleteAsync(exception);
220+
221+
public override async ValueTask<ReadResult> ReadAsync(CancellationToken cancellationToken = default)
222+
{
223+
var result = await _inner.ReadAsync(cancellationToken);
224+
Observe(result);
225+
return result;
226+
}
227+
228+
public override bool TryRead(out ReadResult result)
229+
{
230+
if (!_inner.TryRead(out result))
231+
return false;
232+
233+
Observe(result);
234+
return true;
235+
}
236+
237+
private void Observe(ReadResult result)
238+
{
239+
_currentBuffer = result.Buffer;
240+
_state.Reading(result.Buffer.Length);
241+
}
242+
}
243+
244+
internal sealed class CountingPipeWriter : PipeWriter
245+
{
246+
private readonly PipeWriter _inner;
247+
private readonly NodeConnectionState _state;
248+
249+
public CountingPipeWriter(PipeWriter inner, NodeConnectionState state)
250+
{
251+
_inner = inner;
252+
_state = state;
253+
}
254+
255+
public override void Advance(int bytes)
256+
{
257+
_state.Sending(bytes);
258+
_inner.Advance(bytes);
259+
}
260+
261+
public override void CancelPendingFlush() => _inner.CancelPendingFlush();
262+
263+
public override void Complete(Exception exception = null) => _inner.Complete(exception);
264+
265+
public override ValueTask CompleteAsync(Exception exception = null) => _inner.CompleteAsync(exception);
266+
267+
public override async ValueTask<FlushResult> FlushAsync(CancellationToken cancellationToken = default)
268+
{
269+
var result = await _inner.FlushAsync(cancellationToken);
270+
if (!result.IsCanceled)
271+
_state.Sent();
272+
return result;
273+
}
274+
275+
public override Memory<byte> GetMemory(int sizeHint = 0) => _inner.GetMemory(sizeHint);
276+
277+
public override Span<byte> GetSpan(int sizeHint = 0) => _inner.GetSpan(sizeHint);
278+
}

0 commit comments

Comments
 (0)