Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
27 changes: 4 additions & 23 deletions src/EventStore.ClusterNode/Components/Pages/Cluster.razor
Original file line number Diff line number Diff line change
Expand Up @@ -43,16 +43,15 @@
<th class="px-5 py-4">Status</th>
<th class="px-5 py-4">Timestamp (UTC)</th>
<th class="px-5 py-4">Checkpoints</th>
<th class="px-5 py-4">TCP</th>
<th class="px-5 py-4">HTTP</th>
<th class="px-5 py-4">HTTP / gRPC</th>
<th class="px-5 py-4 text-right">Actions</th>
</tr>
</thead>
<tbody class="divide-y divide-es-ink/10">
@if (!ClusterMembers.Any())
{
<tr>
<td class="px-5 py-4 text-es-muted" colspan="7">@ClusterEmptyMessage</td>
<td class="px-5 py-4 text-es-muted" colspan="6">@ClusterEmptyMessage</td>
</tr>
}
else
Expand All @@ -77,10 +76,6 @@
<p class="mt-1">@EpochLabel(member)</p>
}
</td>
<td class="px-5 py-4 font-mono text-xs text-es-muted">
<p>Internal: @InternalTcpEndpoint(member)</p>
<p class="mt-1">External: @ExternalTcpEndpoint(member)</p>
</td>
<td class="px-5 py-4 font-mono text-xs text-es-muted">@HttpEndpoint(member)</td>
<td class="px-5 py-4 text-right">
<div class="flex flex-wrap justify-end gap-2">
Expand Down Expand Up @@ -307,7 +302,7 @@
<section class="mt-5 grid gap-4 md:grid-cols-2 xl:grid-cols-4">
<SurfaceCard Eyebrow="Explore" Title="Navigator" Description="Jump into streams, subscriptions, projections, user management, and browser tools from one place." Href="/ui/navigator" LinkText="Open" />
<SurfaceCard Eyebrow="Run" Title="Operations" Description="Reach privileged actions for scavenging, shutdown, reload, and node-level workflows." Href="/ui/operations" LinkText="Open" />
<SurfaceCard Eyebrow="Watch" Title="Observability" Description="Inspect queues, replication, TCP, metrics, and health through a curated operator map." Href="/ui/observability" LinkText="Open" />
<SurfaceCard Eyebrow="Watch" Title="Observability" Description="Inspect queues, replication, metrics, and health through a curated operator map." Href="/ui/observability" LinkText="Open" />
<SurfaceCard Eyebrow="Review" Title="Configuration" Description="Find runtime information, loaded options, and subsystem metadata." Href="/ui/configuration" LinkText="Open" />
</section>
</div>
Expand Down Expand Up @@ -370,18 +365,14 @@
.Append("Snapshot taken at ")
.Append(TimestampLabel(ClusterReadAt ?? DateTime.UtcNow))
.AppendLine()
.Append(PadRight("Internal Tcp", 31)).Append(' ')
.Append(PadRight("External Tcp", 31)).Append(' ')
.Append(PadRight("Http", 23)).Append(' ')
.Append(PadRight("HTTP / gRPC", 23)).Append(' ')
.Append(PadRight("Status", 11)).Append(' ')
.Append(PadRight("State", 18)).Append(' ')
.Append(PadRight("Timestamp (UTC)", 19)).Append(" Checkpoints");

foreach (var member in ClusterMembers)
{
builder.AppendLine()
.Append(PadRight(InternalTcpEndpoint(member), 31)).Append(' ')
.Append(PadRight(ExternalTcpEndpoint(member), 31)).Append(' ')
.Append(PadRight(HttpEndpoint(member), 23)).Append(' ')
.Append(PadRight(MemberStatus(member), 11)).Append(' ')
.Append(PadRight(member.State.ToString(), 18)).Append(' ')
Expand Down Expand Up @@ -500,16 +491,6 @@
private static string MemberStatus(ClientClusterInfo.ClientMemberInfo member) =>
member.IsAlive ? "Alive" : "Unreachable";

private static string InternalTcpEndpoint(ClientClusterInfo.ClientMemberInfo member) =>
Endpoint(
member.InternalTcpIp,
member.InternalSecureTcpPort != 0 ? member.InternalSecureTcpPort : member.InternalTcpPort);

private static string ExternalTcpEndpoint(ClientClusterInfo.ClientMemberInfo member) =>
Endpoint(
member.ExternalTcpIp,
member.ExternalSecureTcpPort != 0 ? member.ExternalSecureTcpPort : member.ExternalTcpPort);

private static string HttpEndpoint(ClientClusterInfo.ClientMemberInfo member) =>
Endpoint(member.HttpEndPointIp, member.HttpEndPointPort);

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -171,7 +171,7 @@ private ClusterReplicaRow ParseReplicaRow(
: Guid.Empty;
var totalBytesSent = row.TotalBytesSent;
var previousRow = _previousReplicas.GetValueOrDefault(connectionId);
var replicaNode = FindMemberByInternalEndpoint(members, row.SubscriptionEndpoint);
var replicaNode = FindMemberByEndpoint(members, row.SubscriptionEndpoint);
var isCatchingUp = replicaNode?.State == VNodeState.CatchingUp;
var catchupStartTime = now;
var catchupStartBytesSent = totalBytesSent;
Expand Down Expand Up @@ -206,24 +206,19 @@ private ClusterReplicaRow ParseReplicaRow(
private ClaimsPrincipal CurrentUser =>
httpContextAccessor.HttpContext?.User ?? new ClaimsPrincipal(new ClaimsIdentity());

private static ClientClusterInfo.ClientMemberInfo FindMemberByInternalEndpoint(
private static ClientClusterInfo.ClientMemberInfo FindMemberByEndpoint(
IReadOnlyList<ClientClusterInfo.ClientMemberInfo> members,
string endpoint)
{
var cleaned = endpoint.Replace("Unspecified/", "", StringComparison.OrdinalIgnoreCase);
return members.FirstOrDefault(x => string.Equals(InternalTcpEndpoint(x), cleaned, StringComparison.OrdinalIgnoreCase));
return members.FirstOrDefault(x => string.Equals(HttpEndpoint(x), cleaned, StringComparison.OrdinalIgnoreCase));
Comment thread
cursor[bot] marked this conversation as resolved.
}

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

private static string InternalTcpEndpoint(ClientClusterInfo.ClientMemberInfo member) =>
Endpoint(
member.InternalTcpIp,
member.InternalSecureTcpPort != 0 ? member.InternalSecureTcpPort : member.InternalTcpPort);

private static string HttpEndpoint(ClientClusterInfo.ClientMemberInfo member) =>
Endpoint(member.HttpEndPointIp, member.HttpEndPointPort);

Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,278 @@
using System;
using System.Buffers;
using System.Collections.Concurrent;
using System.Collections.Generic;
using System.IO.Pipelines;
using System.Linq;
using System.Threading;
using System.Threading.Tasks;
using Microsoft.AspNetCore.Connections;

namespace EventStore.ClusterNode.Components.Services;

public sealed class NodeConnectionTracker
{
private readonly ConcurrentDictionary<string, NodeConnectionState> _connections = new();

public IReadOnlyList<NodeConnectionSnapshot> Snapshot() =>
_connections.Values.Select(x => x.Snapshot())
.OrderBy(x => x.RemoteEndPoint, StringComparer.OrdinalIgnoreCase)
.ThenBy(x => x.ConnectionId, StringComparer.Ordinal)
.ToArray();

public async Task Track(ConnectionContext context, ConnectionDelegate next, bool isTls)
{
var state = new NodeConnectionState(
context.ConnectionId,
context.RemoteEndPoint?.ToString() ?? "",
context.LocalEndPoint?.ToString() ?? "",
isTls,
DateTimeOffset.UtcNow);
_connections[context.ConnectionId] = state;
context.Transport = new CountingDuplexPipe(context.Transport, state);

try
{
await next(context);
}
finally
{
_connections.TryRemove(context.ConnectionId, out _);
}
}

public void ObserveRequest(
string connectionId,
string protocol,
bool isGrpc,
string connectionName,
string userAgent)
{
if (_connections.TryGetValue(connectionId, out var connection))
connection.ObserveRequest(protocol, isGrpc, connectionName, userAgent);
}
}

public sealed record NodeConnectionSnapshot(
string ConnectionId,
string RemoteEndPoint,
string LocalEndPoint,
string ClientName,
string Application,
string Protocol,
bool IsTls,
DateTimeOffset ConnectedAt,
long TotalBytesSent,
long TotalBytesReceived,
long PendingSendBytes,
long PendingReceivedBytes);

internal sealed class NodeConnectionState
{
private readonly object _metadataLock = new();
private string _clientName = "";
private bool _hasExplicitConnectionName;
private bool _hasGrpcRequests;
private bool _hasHttpRequests;
private long _pendingReceivedBytes;
private long _pendingSendBytes;
private string _protocol = "";
private long _totalBytesReceived;
private long _totalBytesSent;

public NodeConnectionState(
string connectionId,
string remoteEndPoint,
string localEndPoint,
bool isTls,
DateTimeOffset connectedAt)
{
ConnectionId = connectionId;
RemoteEndPoint = remoteEndPoint;
LocalEndPoint = localEndPoint;
IsTls = isTls;
ConnectedAt = connectedAt;
}

private string ConnectionId { get; }
private string RemoteEndPoint { get; }
private string LocalEndPoint { get; }
private bool IsTls { get; }
private DateTimeOffset ConnectedAt { get; }

public void Received(long bytes, long pendingBytes)
{
Interlocked.Add(ref _totalBytesReceived, bytes);
Interlocked.Exchange(ref _pendingReceivedBytes, pendingBytes);
}

public void Reading(long pendingBytes) =>
Interlocked.Exchange(ref _pendingReceivedBytes, pendingBytes);

public void Sending(int bytes)
{
Interlocked.Add(ref _totalBytesSent, bytes);
Interlocked.Add(ref _pendingSendBytes, bytes);
}

public void Sent() => Interlocked.Exchange(ref _pendingSendBytes, 0);

public void ObserveRequest(
string protocol,
bool isGrpc,
string connectionName,
string userAgent)
{
lock (_metadataLock)
{
_protocol = Merge(_protocol, protocol);
_hasGrpcRequests |= isGrpc;
_hasHttpRequests |= !isGrpc;

if (!string.IsNullOrWhiteSpace(connectionName))
{
_clientName = connectionName;
_hasExplicitConnectionName = true;
}
else if (!_hasExplicitConnectionName && !string.IsNullOrWhiteSpace(userAgent))
{
_clientName = userAgent;
}
}
}

public NodeConnectionSnapshot Snapshot()
{
lock (_metadataLock)
{
return new(
ConnectionId,
RemoteEndPoint,
LocalEndPoint,
_clientName,
ApplicationLabel(),
_protocol,
IsTls,
ConnectedAt,
Interlocked.Read(ref _totalBytesSent),
Interlocked.Read(ref _totalBytesReceived),
Interlocked.Read(ref _pendingSendBytes),
Interlocked.Read(ref _pendingReceivedBytes));
}
}

private string ApplicationLabel() => (_hasHttpRequests, _hasGrpcRequests) switch
{
(true, true) => "HTTP and gRPC",
(false, true) => "gRPC",
(true, false) => "HTTP",
_ => "Awaiting request"
};

private static string Merge(string current, string observed)
{
if (string.IsNullOrWhiteSpace(observed) || current == observed)
return current;
return string.IsNullOrWhiteSpace(current) ? observed : "Mixed";
}
}

internal sealed class CountingDuplexPipe : IDuplexPipe
{
public CountingDuplexPipe(IDuplexPipe inner, NodeConnectionState state)
{
Input = new CountingPipeReader(inner.Input, state);
Output = new CountingPipeWriter(inner.Output, state);
}

public PipeReader Input { get; }
public PipeWriter Output { get; }
}

internal sealed class CountingPipeReader : PipeReader
{
private readonly PipeReader _inner;
private readonly NodeConnectionState _state;
private ReadOnlySequence<byte> _currentBuffer;

public CountingPipeReader(PipeReader inner, NodeConnectionState state)
{
_inner = inner;
_state = state;
}

public override void AdvanceTo(SequencePosition consumed) => AdvanceTo(consumed, consumed);

public override void AdvanceTo(SequencePosition consumed, SequencePosition examined)
{
var consumedBytes = _currentBuffer.IsEmpty ? 0 : _currentBuffer.Slice(0, consumed).Length;
var pendingBytes = _currentBuffer.IsEmpty ? 0 : _currentBuffer.Slice(consumed).Length;
_state.Received(consumedBytes, pendingBytes);
_currentBuffer = default;
_inner.AdvanceTo(consumed, examined);
}

public override void CancelPendingRead() => _inner.CancelPendingRead();

public override void Complete(Exception exception = null) => _inner.Complete(exception);

public override ValueTask CompleteAsync(Exception exception = null) => _inner.CompleteAsync(exception);

public override async ValueTask<ReadResult> ReadAsync(CancellationToken cancellationToken = default)
{
var result = await _inner.ReadAsync(cancellationToken);
Observe(result);
return result;
}

public override bool TryRead(out ReadResult result)
{
if (!_inner.TryRead(out result))
return false;

Observe(result);
return true;
}

private void Observe(ReadResult result)
{
_currentBuffer = result.Buffer;
_state.Reading(result.Buffer.Length);
}
}

internal sealed class CountingPipeWriter : PipeWriter
{
private readonly PipeWriter _inner;
private readonly NodeConnectionState _state;

public CountingPipeWriter(PipeWriter inner, NodeConnectionState state)
{
_inner = inner;
_state = state;
}

public override void Advance(int bytes)
{
_state.Sending(bytes);
_inner.Advance(bytes);
}

public override void CancelPendingFlush() => _inner.CancelPendingFlush();

public override void Complete(Exception exception = null) => _inner.Complete(exception);

public override ValueTask CompleteAsync(Exception exception = null) => _inner.CompleteAsync(exception);

public override async ValueTask<FlushResult> FlushAsync(CancellationToken cancellationToken = default)
{
var result = await _inner.FlushAsync(cancellationToken);
if (!result.IsCanceled)
_state.Sent();
return result;
}

public override Memory<byte> GetMemory(int sizeHint = 0) => _inner.GetMemory(sizeHint);

public override Span<byte> GetSpan(int sizeHint = 0) => _inner.GetSpan(sizeHint);
}
Loading
Loading