diff --git a/docs/upgrade-guide.md b/docs/upgrade-guide.md
index b1f10f5a7b..316c1fffeb 100644
--- a/docs/upgrade-guide.md
+++ b/docs/upgrade-guide.md
@@ -68,6 +68,25 @@ methods. Avoid depending on undocumented authentication plugins.
Use [OpenTelemetry integration](diagnostics/integrations.md) for explicit OTLP
export and [Metrics](diagnostics/metrics.md) for Prometheus scraping.
+## TestClient review
+
+The TestClient keeps command names only when their observable behavior can be
+preserved over the supported gRPC APIs. Reads, writes, hard deletes,
+subscriptions, scavenging, data verification, and load operations use the gRPC
+client on `NodePort`. The generic `WRFL` alias remains available and runs the
+same workload as `WRFLGRPC`.
+
+The following commands are intentionally retired:
+
+- `TWR`, because the public gRPC streams API does not expose the transaction
+ start, write, and commit lifecycle. A batch append is not an equivalent test.
+- `WRFLTCP` and `WRFLCA`, because transport-specific aliases would hide which
+ protocol the workload exercises.
+- `RT`, because its projection and node-failure scenarios require a dedicated
+ replacement rather than a different command behind the same name.
+- `CHKTCP`, because it validates a retired frame protocol. Use `CHKGRPC` to
+ validate malformed gRPC-frame handling.
+
Legacy usage telemetry is separate from OTLP observability. See
[Usage telemetry](usage-telemetry.md) before running a node in an environment
that should not make outbound telemetry calls.
diff --git a/scripts/test.sh b/scripts/test.sh
index 54c3fc0363..7e90473284 100755
--- a/scripts/test.sh
+++ b/scripts/test.sh
@@ -52,6 +52,7 @@ misc_projects=(
EventStore.Common.Tests
EventStore.SourceGenerators.Tests
EventStore.SystemRuntime.Tests
+ EventStore.TestClient.Tests
)
declare -a requested_projects
diff --git a/src/Directory.Packages.props b/src/Directory.Packages.props
index cf693dd4db..4ba5fa1a9f 100644
--- a/src/Directory.Packages.props
+++ b/src/Directory.Packages.props
@@ -12,6 +12,7 @@
+
diff --git a/src/EventStore.TestClient.Tests/CommandInventoryTests.cs b/src/EventStore.TestClient.Tests/CommandInventoryTests.cs
new file mode 100644
index 0000000000..5dc4d34ce1
--- /dev/null
+++ b/src/EventStore.TestClient.Tests/CommandInventoryTests.cs
@@ -0,0 +1,44 @@
+using System;
+using System.Linq;
+using System.Threading;
+using NUnit.Framework;
+
+namespace EventStore.TestClient.Tests;
+
+[TestFixture]
+public class CommandInventoryTests
+{
+ private static readonly string[] SupportedCommandKeywords =
+ [
+ "PING", "PINGFL", "PINGFLW", "RDALLGRPC", "WRFLGRPC",
+ "WR", "WRJ", "WRFL", "WRFLW",
+ "MWR", "MWRFLW", "DEL",
+ "RDALL", "RD", "RDFL", "WRLT",
+ "VERIFY", "SUBSCR", "SCAVENGE", "CHKGRPC", "SST"
+ ];
+ private static readonly string[] RetiredCommandKeywords = ["TWR", "WRFLCA", "WRFLTCP", "RT", "CHKTCP"];
+
+ [Test]
+ public void test_client_preserves_supported_command_surface()
+ {
+ using var cancellation = new CancellationTokenSource();
+ var client = new Client(new ClientOptions(), cancellation);
+ var usages = client.GetCommandList()
+ .Split(Environment.NewLine, StringSplitOptions.RemoveEmptyEntries | StringSplitOptions.TrimEntries);
+
+ Assert.Multiple(() =>
+ {
+ foreach (var keyword in SupportedCommandKeywords)
+ {
+ Assert.That(usages.Any(x => x == keyword || x.StartsWith($"{keyword} ", StringComparison.Ordinal)),
+ Is.True, $"Missing TestClient command '{keyword}'");
+ }
+
+ foreach (var keyword in RetiredCommandKeywords)
+ {
+ Assert.That(usages.Any(x => x == keyword || x.StartsWith($"{keyword} ", StringComparison.Ordinal)),
+ Is.False, $"TestClient command '{keyword}' advertises behavior that gRPC does not provide");
+ }
+ });
+ }
+}
diff --git a/src/EventStore.TestClient.Tests/EventStore.TestClient.Tests.csproj b/src/EventStore.TestClient.Tests/EventStore.TestClient.Tests.csproj
new file mode 100644
index 0000000000..2cb72fad1a
--- /dev/null
+++ b/src/EventStore.TestClient.Tests/EventStore.TestClient.Tests.csproj
@@ -0,0 +1,20 @@
+
+
+ true
+
+
+
+
+
+
+
+
+
+
+
+
+ version.properties
+ PreserveNewest
+
+
+
diff --git a/src/EventStore.TestClient.Tests/GrpcCommandBehaviorTests.cs b/src/EventStore.TestClient.Tests/GrpcCommandBehaviorTests.cs
new file mode 100644
index 0000000000..e015876eee
--- /dev/null
+++ b/src/EventStore.TestClient.Tests/GrpcCommandBehaviorTests.cs
@@ -0,0 +1,571 @@
+using System;
+using System.Collections.Concurrent;
+using System.Linq;
+using System.Net;
+using System.Net.Http;
+using System.Text;
+using System.Threading;
+using System.Threading.Tasks;
+using EventStore.Client;
+using EventStore.TestClient.GrpcCommands;
+using Newtonsoft.Json.Linq;
+using NUnit.Framework;
+
+namespace EventStore.TestClient.Tests;
+
+[TestFixture]
+public class GrpcCommandBehaviorTests
+{
+ [Test]
+ public void grpc_workloads_keep_the_existing_backpressure_defaults()
+ {
+ var options = new ClientOptions();
+
+ Assert.Multiple(() =>
+ {
+ Assert.That(options.PingWindow, Is.EqualTo(2_000));
+ Assert.That(options.ReadWindow, Is.EqualTo(2_000));
+ Assert.That(options.WriteWindow, Is.EqualTo(2_000));
+ });
+ }
+
+ [Test]
+ public void ping_flood_defaults_preserve_the_historical_total_request_count()
+ {
+ var workload = FloodWorkload.Parse([], defaultRequestCount: 1_000_000, maxInFlight: 2_000);
+
+ Assert.Multiple(() =>
+ {
+ Assert.That(workload.ClientCount, Is.EqualTo(1));
+ Assert.That(workload.RequestCount, Is.EqualTo(1_000_000));
+ Assert.That(workload.RequestsForClient(0), Is.EqualTo(1_000_000));
+ });
+ }
+
+ [Test]
+ public void ping_flood_waiting_defaults_preserve_the_historical_total_request_count()
+ {
+ var workload = FloodWorkload.Parse([], defaultRequestCount: 100_000, maxInFlight: 1);
+
+ Assert.That(workload.RequestCount, Is.EqualTo(100_000));
+ }
+
+ [Test]
+ public void flood_requests_are_distributed_as_one_total_across_clients()
+ {
+ var workload = FloodWorkload.Parse(["4", "10"], defaultRequestCount: 1, maxInFlight: 2_000);
+
+ Assert.Multiple(() =>
+ {
+ Assert.That(workload.RequestsForClient(0), Is.EqualTo(2));
+ Assert.That(workload.RequestsForClient(1), Is.EqualTo(2));
+ Assert.That(workload.RequestsForClient(2), Is.EqualTo(2));
+ Assert.That(workload.RequestsForClient(3), Is.EqualTo(4));
+ Assert.That(workload.MaxInFlightPerClient, Is.EqualTo(500));
+ });
+ }
+
+ [TestCase(true, "True")]
+ [TestCase(false, "False")]
+ public async Task read_commands_control_the_requires_leader_header(bool requireLeader, string expectedHeader)
+ {
+ var handler = new RecordingGrpcHandler();
+ var grpcClient = new GrpcTestClient(new ClientOptions(), Serilog.Log.Logger, () =>
+ {
+ var settings = EventStoreClientSettings.Create("esdb://localhost:2113?tls=false");
+ settings.CreateHttpMessageHandler = () => handler;
+ return settings;
+ });
+ using var client = grpcClient.CreateGrpcClient(requireLeader);
+ var read = client.ReadStreamAsync(Direction.Forwards, "test-stream", StreamPosition.Start, maxCount: 1);
+ var state = await read.ReadState;
+
+ Assert.Multiple(() =>
+ {
+ Assert.That(state, Is.EqualTo(ReadState.StreamNotFound));
+ Assert.That(handler.Requests.Single(request => request.Path.EndsWith("/Read", StringComparison.Ordinal)).RequiresLeader,
+ Is.EqualTo(expectedHeader));
+ });
+ }
+
+ [TestCase(0, "AccountDebited")]
+ [TestCase(1, "AccountCredited")]
+ [TestCase(9, "AccountCheckPoint")]
+ public void verification_events_preserve_deterministic_typed_json(int eventNumber, string expectedType)
+ {
+ var first = VerificationEventFactory.Create(eventNumber);
+ var second = VerificationEventFactory.Create(eventNumber);
+
+ Assert.Multiple(() =>
+ {
+ Assert.That(first.Type, Is.EqualTo(expectedType));
+ Assert.That(first.ContentType, Is.EqualTo("application/json"));
+ Assert.That(first.Data.ToArray(), Is.EqualTo(second.Data.ToArray()));
+ Assert.That(() => JObject.Parse(System.Text.Encoding.UTF8.GetString(first.Data.Span)), Throws.Nothing);
+ });
+ }
+
+ [Test]
+ public void delete_command_uses_the_grpc_tombstone_operation()
+ {
+ var handler = new RecordingGrpcHandler();
+ var options = new ClientOptions { Command = ["DEL test-stream ANY"] };
+ var grpcClient = new GrpcTestClient(options, Serilog.Log.Logger, () =>
+ {
+ var settings = EventStoreClientSettings.Create("esdb://localhost:2113?tls=false");
+ settings.CreateHttpMessageHandler = () => handler;
+ return settings;
+ });
+ using var cancellation = new CancellationTokenSource();
+ var client = new Client(options, cancellation, grpcClient);
+
+ Assert.Multiple(() =>
+ {
+ Assert.That(client.Run(cancellation.Token), Is.Zero);
+ Assert.That(handler.Requests.Select(request => request.Path),
+ Does.Contain("/event_store.client.streams.Streams/Tombstone"));
+ Assert.That(handler.Requests.Select(request => request.Path),
+ Does.Not.Contain("/event_store.client.streams.Streams/Delete"));
+ });
+ }
+
+ [Test]
+ public void ping_flood_treats_messages_as_one_total_across_clients()
+ {
+ var handler = new RecordingGrpcHandler();
+ var options = new ClientOptions { Command = ["PINGFL 4 10"] };
+ var grpcClient = new GrpcTestClient(options, Serilog.Log.Logger, () =>
+ {
+ var settings = EventStoreClientSettings.Create("esdb://localhost:2113?tls=false");
+ settings.CreateHttpMessageHandler = () => handler;
+ return settings;
+ });
+ using var cancellation = new CancellationTokenSource();
+ var client = new Client(options, cancellation, grpcClient);
+
+ Assert.Multiple(() =>
+ {
+ Assert.That(client.Run(cancellation.Token), Is.Zero);
+ Assert.That(handler.Requests.Count(request => request.Path.EndsWith("/Read", StringComparison.Ordinal)),
+ Is.EqualTo(10));
+ });
+ }
+
+ [TestCase("WRFLW 2 4")]
+ [TestCase("MWRFLW 2 2 4")]
+ public void waiting_write_floods_finish_every_request_before_reporting_failures(string command)
+ {
+ var handler = new RecordingGrpcHandler
+ {
+ Failure = new HttpRequestException("first append failed"),
+ FailurePath = "/event_store.client.streams.Streams/Append",
+ FailureLimit = 1
+ };
+
+ var result = RunCommand(command, handler);
+
+ Assert.Multiple(() =>
+ {
+ Assert.That(result, Is.Not.Zero);
+ Assert.That(handler.Requests.Count(request => request.Path.EndsWith("/Append", StringComparison.Ordinal)),
+ Is.EqualTo(4));
+ });
+ }
+
+ [TestCase("WRFLW 1 3G")]
+ [TestCase("MWRFLW 1 1 3G")]
+ public void waiting_write_floods_accept_long_request_totals(string command)
+ {
+ using var cancellation = new CancellationTokenSource();
+ var handler = new RecordingGrpcHandler
+ {
+ OnRequest = request =>
+ {
+ if (request.RequestUri!.AbsolutePath.EndsWith("/Append", StringComparison.Ordinal))
+ {
+ cancellation.Cancel();
+ }
+ }
+ };
+
+ Assert.Multiple(() =>
+ {
+ Assert.That(
+ () => RunCommand(command, handler, new ClientOptions { Timeout = 2 }, cancellation),
+ Throws.InstanceOf());
+ Assert.That(handler.Requests.Count(request => request.Path.EndsWith("/Append", StringComparison.Ordinal)),
+ Is.EqualTo(1));
+ });
+ }
+
+ [Test]
+ public void ping_uses_a_user_stream_for_unauthenticated_connectivity_checks()
+ {
+ var handler = new RecordingGrpcHandler();
+ var result = RunCommand("PING", handler);
+ var request = handler.Requests.Single(request => request.Path.EndsWith("/Read", StringComparison.Ordinal));
+
+ Assert.Multiple(() =>
+ {
+ Assert.That(result, Is.Zero);
+ Assert.That(Encoding.UTF8.GetString(request.Body), Does.Not.Contain("$test-client-ping"));
+ });
+ }
+
+ [TestCase("SUBSCR", true, null)]
+ [TestCase("SUBSCR test-stream", false, "test-stream")]
+ [TestCase("SST 1", false, "stream-0")]
+ public async Task subscriptions_start_from_the_live_position(string command, bool subscribeToAll, string stream)
+ {
+ var expected = await RecordSubscriptionRequest(subscribeToAll, stream);
+ var handler = new RecordingGrpcHandler
+ {
+ Failure = new HttpRequestException("stop after recording"),
+ FailurePath = "/event_store.client.streams.Streams/Read"
+ };
+
+ RunCommand(command, handler);
+
+ Assert.That(
+ handler.Requests.Single(request => request.Path == "/event_store.client.streams.Streams/Read").Body,
+ Is.EqualTo(expected));
+ }
+
+ [TestCase("SUBSCR test-stream")]
+ [TestCase("SST 1")]
+ public void subscriptions_report_server_drops(string command)
+ {
+ var handler = new RecordingGrpcHandler
+ {
+ ReadPayload = [0x12, 0x00],
+ GrpcStatus = "13"
+ };
+ var result = 0;
+
+ Assert.That(
+ () => result = RunCommand(command, handler, new ClientOptions { Timeout = 1 }),
+ Throws.Nothing);
+ Assert.That(result, Is.Not.Zero);
+ }
+
+ [Test]
+ public void read_flood_completes_every_request_before_reporting_missing_streams()
+ {
+ var handler = new RecordingGrpcHandler();
+ var result = RunCommand("RDFL 2 5 1 missing-stream", handler, new ClientOptions { ReadWindow = 1 });
+
+ Assert.Multiple(() =>
+ {
+ Assert.That(result, Is.Not.Zero);
+ Assert.That(handler.Requests.Count(request => request.Path.EndsWith("/Read", StringComparison.Ordinal)),
+ Is.EqualTo(5));
+ });
+ }
+
+ [Test]
+ public void read_command_fails_when_the_requested_event_is_missing()
+ {
+ var handler = new RecordingGrpcHandler { ReadPayload = [] };
+
+ Assert.That(RunCommand("RD test-stream 42", handler), Is.Not.Zero);
+ }
+
+ [Test]
+ public async Task read_command_preserves_the_last_event_sentinel()
+ {
+ var expected = await RecordStreamReadRequest(Direction.Backwards, StreamPosition.End);
+ var handler = new RecordingGrpcHandler
+ {
+ Failure = new HttpRequestException("stop after recording"),
+ FailurePath = "/event_store.client.streams.Streams/Read"
+ };
+
+ RunCommand("RD test-stream -1", handler);
+
+ Assert.That(
+ handler.Requests.Single(request => request.Path == "/event_store.client.streams.Streams/Read").Body,
+ Is.EqualTo(expected));
+ }
+
+ [Test]
+ public async Task read_all_accepts_explicit_end_position_sentinels()
+ {
+ var expected = await RecordReadAllRequest(Direction.Backwards, Position.End);
+ var handler = new RecordingGrpcHandler { ReadPayload = [] };
+
+ var result = RunCommand("RDALL B -1 -1", handler);
+
+ Assert.Multiple(() =>
+ {
+ Assert.That(result, Is.Zero);
+ Assert.That(
+ handler.Requests.Single(request => request.Path == "/event_store.client.streams.Streams/Read").Body,
+ Is.EqualTo(expected));
+ });
+ }
+
+ [Test]
+ public void read_all_rejects_unknown_directions()
+ {
+ var handler = new RecordingGrpcHandler { ReadPayload = [] };
+
+ var result = RunCommand("RDALL SIDEWAYS", handler);
+
+ Assert.Multiple(() =>
+ {
+ Assert.That(result, Is.Not.Zero);
+ Assert.That(handler.Requests.Select(request => request.Path),
+ Does.Not.Contain("/event_store.client.streams.Streams/Read"));
+ });
+ }
+
+ [TestCase("-1", "NOSTREAM")]
+ [TestCase("-2", "ANY")]
+ public void expected_version_sentinels_preserve_their_historical_meaning(string sentinel, string name)
+ {
+ var expectedHandler = new RecordingGrpcHandler();
+ Assert.That(RunCommand($"DEL test-stream {name}", expectedHandler), Is.Zero);
+
+ var sentinelHandler = new RecordingGrpcHandler();
+ Assert.That(RunCommand($"DEL test-stream {sentinel}", sentinelHandler), Is.Zero);
+
+ Assert.That(
+ sentinelHandler.Requests.Single(request => request.Path.EndsWith("/Tombstone", StringComparison.Ordinal)).Body,
+ Is.EqualTo(expectedHandler.Requests.Single(request => request.Path.EndsWith("/Tombstone", StringComparison.Ordinal)).Body));
+ }
+
+ [Test]
+ public void grpc_http_endpoint_uses_the_resolved_single_node_address()
+ {
+ var settings = EventStoreClientSettings.Create("esdb://configured.example:3210?tls=false");
+ var grpcClient = new GrpcTestClient(
+ new ClientOptions { Host = "ignored.example", HttpPort = 9999 },
+ Serilog.Log.Logger,
+ () => settings);
+
+ Assert.That(grpcClient.HttpEndpoint, Is.EqualTo(settings.ConnectivitySettings.Address));
+ }
+
+ [TestCase("esdb+discover://localhost:2113?tls=false")]
+ [TestCase("esdb://node-1:2113,node-2:2113?tls=false")]
+ public void grpc_http_endpoint_rejects_connections_without_one_resolved_address(string connectionString)
+ {
+ var settings = EventStoreClientSettings.Create(connectionString);
+ var grpcClient = new GrpcTestClient(new ClientOptions(), Serilog.Log.Logger, () => settings);
+
+ Assert.That(() => grpcClient.HttpEndpoint, Throws.InvalidOperationException);
+ }
+
+ [Test]
+ public void what_if_options_do_not_render_connection_string_credentials()
+ {
+ var rendered = new ClientOptions
+ {
+ ConnectionString = "esdb://test-user:test-password@localhost:2113?tls=false"
+ }.ToString();
+
+ Assert.Multiple(() =>
+ {
+ Assert.That(rendered, Does.Not.Contain("test-user"));
+ Assert.That(rendered, Does.Not.Contain("test-password"));
+ Assert.That(rendered, Does.Contain("ConnectionString: [REDACTED]"));
+ });
+ }
+
+ [Test]
+ public void verification_readers_wait_without_reposting_while_streams_are_empty()
+ {
+ var handler = new RecordingGrpcHandler
+ {
+ Delay = TimeSpan.FromMilliseconds(100),
+ DelayPath = "/event_store.client.streams.Streams/Append"
+ };
+ var completedWorkItems = ThreadPool.CompletedWorkItemCount;
+ var command = Task.Run(() =>
+ RunCommand("VERIFY 1 2 1 1 bank", handler, new ClientOptions { Timeout = 2 }));
+
+ Assert.Multiple(() =>
+ {
+ Assert.That(command.Wait(TimeSpan.FromSeconds(5)), Is.True);
+ Assert.That(handler.Requests.Select(request => request.Path), Does.Contain(handler.DelayPath));
+ Assert.That(ThreadPool.CompletedWorkItemCount - completedWorkItems, Is.LessThan(100));
+ });
+ }
+
+ [Test]
+ public void grpc_read_all_preserves_backward_direction_and_default_position()
+ {
+ Assert.That(ReadAllWorkload.TryParse(["B", "1"], out var workload), Is.True);
+ Assert.Multiple(() =>
+ {
+ Assert.That(workload.Direction, Is.EqualTo(Direction.Backwards));
+ Assert.That(workload.Position, Is.EqualTo(Position.End));
+ Assert.That(workload.ClientCount, Is.EqualTo(1));
+ });
+ }
+
+ private static int RunCommand(
+ string command,
+ RecordingGrpcHandler handler,
+ ClientOptions options = null,
+ CancellationTokenSource cancellation = null)
+ {
+ options = (options ?? new ClientOptions()) with { Command = [command] };
+ var grpcClient = new GrpcTestClient(options, Serilog.Log.Logger, () =>
+ {
+ var settings = EventStoreClientSettings.Create("esdb://localhost:2113?tls=false");
+ settings.CreateHttpMessageHandler = () => handler;
+ return settings;
+ });
+ var ownsCancellation = cancellation is null;
+ cancellation ??= new CancellationTokenSource();
+ try
+ {
+ return new Client(options, cancellation, grpcClient).Run(cancellation.Token);
+ }
+ finally
+ {
+ if (ownsCancellation)
+ {
+ cancellation.Dispose();
+ }
+ }
+ }
+
+ private static async Task RecordSubscriptionRequest(bool subscribeToAll, string stream)
+ {
+ var handler = new RecordingGrpcHandler
+ {
+ Failure = new HttpRequestException("stop after recording"),
+ FailurePath = "/event_store.client.streams.Streams/Read"
+ };
+ var settings = EventStoreClientSettings.Create("esdb://localhost:2113?tls=false");
+ settings.CreateHttpMessageHandler = () => handler;
+ using var client = new EventStoreClient(settings);
+ try
+ {
+ if (subscribeToAll)
+ {
+ await client.SubscribeToAllAsync(FromAll.End, (_, _, _) => Task.CompletedTask);
+ }
+ else
+ {
+ await client.SubscribeToStreamAsync(stream, FromStream.End, (_, _, _) => Task.CompletedTask);
+ }
+ }
+ catch (Grpc.Core.RpcException)
+ {
+ }
+
+ return handler.Requests.Single(request => request.Path == "/event_store.client.streams.Streams/Read").Body;
+ }
+
+ private static async Task RecordStreamReadRequest(Direction direction, StreamPosition position)
+ {
+ var handler = new RecordingGrpcHandler
+ {
+ Failure = new HttpRequestException("stop after recording"),
+ FailurePath = "/event_store.client.streams.Streams/Read"
+ };
+ using var client = CreateClient(handler);
+ try
+ {
+ var read = client.ReadStreamAsync(direction, "test-stream", position, maxCount: 1);
+ await read.ReadState;
+ }
+ catch (Grpc.Core.RpcException)
+ {
+ }
+
+ return handler.Requests.Single(request => request.Path == "/event_store.client.streams.Streams/Read").Body;
+ }
+
+ private static async Task RecordReadAllRequest(Direction direction, Position position)
+ {
+ var handler = new RecordingGrpcHandler
+ {
+ Failure = new HttpRequestException("stop after recording"),
+ FailurePath = "/event_store.client.streams.Streams/Read"
+ };
+ using var client = CreateClient(handler);
+ try
+ {
+ var read = client.ReadAllAsync(direction, position);
+ await foreach (var _ in read.Messages)
+ {
+ }
+ }
+ catch (Grpc.Core.RpcException)
+ {
+ }
+
+ return handler.Requests.Single(request => request.Path == "/event_store.client.streams.Streams/Read").Body;
+ }
+
+ private static EventStoreClient CreateClient(HttpMessageHandler handler)
+ {
+ var settings = EventStoreClientSettings.Create("esdb://localhost:2113?tls=false");
+ settings.CreateHttpMessageHandler = () => handler;
+ return new EventStoreClient(settings);
+ }
+
+ private sealed class RecordingGrpcHandler : HttpMessageHandler
+ {
+ private int _failuresIssued;
+
+ public ConcurrentBag Requests { get; } = [];
+ public TimeSpan Delay { get; init; }
+ public string DelayPath { get; init; }
+ public Exception Failure { get; init; }
+ public int? FailureLimit { get; init; }
+ public string FailurePath { get; init; }
+ public Action OnRequest { get; init; }
+ public byte[] ReadPayload { get; init; } = [0x22, 0x00];
+ public string GrpcStatus { get; init; } = "0";
+
+ protected override async Task SendAsync(
+ HttpRequestMessage request,
+ CancellationToken cancellationToken)
+ {
+ var path = request.RequestUri!.AbsolutePath;
+ if (Delay > TimeSpan.Zero && (DelayPath is null || path == DelayPath))
+ {
+ await Task.Delay(Delay, cancellationToken);
+ }
+
+ var requiresLeader = request.Headers.TryGetValues("requires-leader", out var values)
+ ? values.Single()
+ : null;
+ var body = request.Content is null
+ ? []
+ : await request.Content.ReadAsByteArrayAsync(cancellationToken);
+ Requests.Add(new RecordedGrpcRequest(path, requiresLeader, body));
+ OnRequest?.Invoke(request);
+ cancellationToken.ThrowIfCancellationRequested();
+ if (Failure is not null && (FailurePath is null || path == FailurePath) &&
+ (FailureLimit is null || Interlocked.Increment(ref _failuresIssued) <= FailureLimit))
+ {
+ throw Failure;
+ }
+
+ var payload = path.EndsWith("/Tombstone", StringComparison.Ordinal)
+ ? new byte[] { 0x0a, 0x00 }
+ : path.EndsWith("/Read", StringComparison.Ordinal)
+ ? ReadPayload
+ : [];
+ var frame = new byte[payload.Length + 5];
+ frame[4] = (byte)payload.Length;
+ payload.CopyTo(frame, 5);
+ var response = new HttpResponseMessage(HttpStatusCode.OK)
+ {
+ Version = HttpVersion.Version20,
+ Content = new ByteArrayContent(frame)
+ };
+ response.Content.Headers.ContentType = new("application/grpc");
+ response.TrailingHeaders.Add("grpc-status", GrpcStatus);
+ return response;
+ }
+ }
+
+ private sealed record RecordedGrpcRequest(string Path, string RequiresLeader, byte[] Body);
+}
diff --git a/src/EventStore.TestClient/Client.cs b/src/EventStore.TestClient/Client.cs
index 28385ba8a0..faf7d995c2 100644
--- a/src/EventStore.TestClient/Client.cs
+++ b/src/EventStore.TestClient/Client.cs
@@ -4,8 +4,6 @@
using System.Threading;
using EventStore.Common.Utils;
using EventStore.TestClient.Commands;
-using EventStore.TestClient.Commands.DvuBasic;
-using Connection = EventStore.Transport.Tcp.TcpTypedConnection;
using ILogger = Serilog.ILogger;
#pragma warning disable 1591
@@ -20,23 +18,25 @@ public class Client
public readonly ClientOptions Options;
- public readonly TcpTestClient _tcpTestClient;
public readonly GrpcTestClient _grpcTestClient;
- public readonly ClientApiTcpTestClient _clientApiTestClient;
private readonly CommandsProcessor _commands = new CommandsProcessor(Log);
public Client(ClientOptions options, CancellationTokenSource cancellationTokenSource)
+ : this(options, cancellationTokenSource, null)
{
- Options = options;
+ }
- var interactiveMode = options.Command.IsEmpty();
+ internal Client(
+ ClientOptions options,
+ CancellationTokenSource cancellationTokenSource,
+ GrpcTestClient grpcTestClient)
+ {
+ Options = options;
InteractiveMode = options.Command.IsEmpty();
- _tcpTestClient = new TcpTestClient(options, interactiveMode, Log);
- _grpcTestClient = new GrpcTestClient(options, Log);
- _clientApiTestClient = new ClientApiTcpTestClient(options, Log);
+ _grpcTestClient = grpcTestClient ?? new GrpcTestClient(options, Log);
RegisterProcessors(cancellationTokenSource);
}
@@ -46,46 +46,19 @@ private void RegisterProcessors(CancellationTokenSource cancellationTokenSource)
_commands.Register(new UsageProcessor(_commands), usageProcessor: true);
_commands.Register(new ExitProcessor(cancellationTokenSource));
- _commands.Register(new PingProcessor());
- _commands.Register(new PingFloodProcessor());
- _commands.Register(new PingFloodWaitingProcessor());
-
- _commands.Register(new WriteProcessor());
- _commands.Register(new WriteJsonProcessor());
- _commands.Register(new WriteFloodProcessor());
- _commands.Register(new WriteFloodClientApiProcessor());
- _commands.Register(new WriteFloodWaitingProcessor());
-
- _commands.Register(new MultiWriteProcessor());
- _commands.Register(new MultiWriteFloodWaitingProcessor());
-
- _commands.Register(new TransactionWriteProcessor());
-
- _commands.Register(new DeleteProcessor());
-
- _commands.Register(new ReadAllProcessor());
- _commands.Register(new ReadProcessor());
- _commands.Register(new ReadFloodProcessor());
-
- _commands.Register(new WriteLongTermProcessor());
-
- _commands.Register(new DvuBasicProcessor());
- _commands.Register(new RunTestScenariosProcessor());
-
- _commands.Register(new SubscribeToStreamProcessor());
-
- _commands.Register(new ScavengeProcessor());
-
- _commands.Register(new TcpSanitazationCheckProcessor());
-
- _commands.Register(new SubscriptionStressTestProcessor());
-
- // gRPC
_commands.Register(new GrpcCommands.ReadAllProcessor());
- _commands.Register(new GrpcCommands.WriteFloodProcessor());
+ var writeFlood = new GrpcCommands.WriteFloodProcessor();
+ _commands.Register(writeFlood);
+
+ foreach (var processor in GrpcCommands.CompatibilityProcessor.CreateSupportedProcessors())
+ {
+ _commands.Register(processor);
+ }
- // TCP Client API
- _commands.Register(new ClientApiTcpCommands.WriteFloodProcessor());
+ _commands.Register(new GrpcCommands.DelegatingProcessor(
+ "WRFL",
+ "WRFL [ [ [ [ []]]]]",
+ writeFlood));
}
public string GetCommandList()
@@ -149,7 +122,7 @@ private int Execute(string[] args, CancellationToken cancellationToken)
{
Log.Information("Processing command: {command}.", string.Join(" ", args));
- var context = new CommandProcessorContext(_tcpTestClient, _grpcTestClient, _clientApiTestClient, Options.Timeout,
+ var context = new CommandProcessorContext(_grpcTestClient, Options.Timeout,
Log, Options.StatsLog, Options.OutputCsv, new ManualResetEventSlim(true), cancellationToken);
int exitCode;
diff --git a/src/EventStore.TestClient/ClientApiLoggerBridge.cs b/src/EventStore.TestClient/ClientApiLoggerBridge.cs
deleted file mode 100644
index a7dc5f5bd1..0000000000
--- a/src/EventStore.TestClient/ClientApiLoggerBridge.cs
+++ /dev/null
@@ -1,92 +0,0 @@
-using System;
-using EventStore.Common.Utils;
-using ILogger = Serilog.ILogger;
-
-namespace EventStore.TestClient;
-
-internal class ClientApiLoggerBridge : ClientAPI.ILogger
-{
- public static readonly ClientApiLoggerBridge Default =
- new ClientApiLoggerBridge(Serilog.Log.ForContext(Serilog.Core.Constants.SourceContextPropertyName,
- "client-api"));
-
- private readonly Serilog.ILogger _log;
-
- public ClientApiLoggerBridge(ILogger log)
- {
- Ensure.NotNull(log, "log");
- _log = log;
- }
-
- public void Error(string format, params object[] args)
- {
- if (args.Length == 0)
- {
- _log.Error(format);
- }
- else
- {
- _log.Error(format, args);
- }
- }
-
- public void Error(Exception ex, string format, params object[] args)
- {
- if (args.Length == 0)
- {
- _log.Error(ex, format);
- }
- else
- {
- _log.Error(ex, format, args);
- }
- }
-
- public void Info(string format, params object[] args)
- {
- if (args.Length == 0)
- {
- _log.Information(format);
- }
- else
- {
- _log.Information(format, args);
- }
- }
-
- public void Info(Exception ex, string format, params object[] args)
- {
- if (args.Length == 0)
- {
- _log.Information(ex, format);
- }
- else
- {
- _log.Information(ex, format, args);
- }
- }
-
- public void Debug(string format, params object[] args)
- {
- if (args.Length == 0)
- {
- _log.Debug(format);
- }
- else
- {
- _log.Debug(format, args);
- }
- }
-
- public void Debug(Exception ex, string format, params object[] args)
- {
- if (args.Length == 0)
- {
- _log.Debug(ex, format);
- }
- else
- {
- _log.Debug(ex, format, args);
- }
- }
-}
diff --git a/src/EventStore.TestClient/ClientApiTcpCommands/WriteFloodProcessor.cs b/src/EventStore.TestClient/ClientApiTcpCommands/WriteFloodProcessor.cs
deleted file mode 100644
index 4c2951aa5d..0000000000
--- a/src/EventStore.TestClient/ClientApiTcpCommands/WriteFloodProcessor.cs
+++ /dev/null
@@ -1,267 +0,0 @@
-using System;
-using System.Collections.Generic;
-using System.Diagnostics;
-using System.Linq;
-using System.Threading;
-using System.Threading.Tasks;
-using EventStore.ClientAPI;
-using EventStore.ClientAPI.Exceptions;
-using EventStore.TestClient.Commands;
-using EventStore.TestClient.Statistics;
-using EventData = EventStore.ClientAPI.EventData;
-using ExpectedVersion = EventStore.Core.Data.ExpectedVersion;
-using StreamDeletedException = EventStore.ClientAPI.Exceptions.StreamDeletedException;
-using WrongExpectedVersionException = EventStore.ClientAPI.Exceptions.WrongExpectedVersionException;
-
-namespace EventStore.TestClient.ClientApiTcpCommands;
-
-internal class WriteFloodProcessor : ICmdProcessor
-{
- public string Usage
- {
- // 0 1 2 3 4 5
- get { return "WRFLTCP [ [ [ [ []]]]]"; }
- }
-
- public string Keyword
- {
- get { return "WRFLTCP"; }
- }
-
- public bool Execute(CommandProcessorContext context, string[] args)
- {
- int clientsCnt = 1;
- long requestsCnt = 5000;
- int streamsCnt = 1000;
- int size = 256;
- int batchSize = 1;
- string streamPrefix = "";
- if (args.Length > 0)
- {
- if (args.Length < 2 || args.Length > 6)
- {
- return false;
- }
-
- try
- {
- clientsCnt = MetricPrefixValue.ParseInt(args[0]);
- requestsCnt = MetricPrefixValue.ParseLong(args[1]);
- if (args.Length >= 3)
- {
- streamsCnt = MetricPrefixValue.ParseInt(args[2]);
- }
-
- if (args.Length >= 4)
- {
- size = MetricPrefixValue.ParseInt(args[3]);
- }
-
- if (args.Length >= 5)
- {
- batchSize = MetricPrefixValue.ParseInt(args[4]);
- }
-
- if (args.Length >= 6)
- {
- streamPrefix = args[5];
- }
- }
- catch
- {
- return false;
- }
- }
-
- var stats = new WriteFloodStats(Keyword, context.OutputCsv, args);
- var monitor = new RequestMonitor();
- WriteFlood(context, stats, clientsCnt, requestsCnt, streamsCnt, size, batchSize, streamPrefix, monitor)
- .GetAwaiter().GetResult();
- return true;
- }
-
- private async Task WriteFlood(CommandProcessorContext context, WriteFloodStats stats, int clientsCnt, long requestsCnt, int streamsCnt,
- int size, int batchSize, string streamPrefix, RequestMonitor monitor)
- {
- context.IsAsync();
-
- var doneEvent = new ManualResetEventSlim(false);
- var clients = new List();
-
- long last = 0;
-
- var streams = Enumerable.Range(0, streamsCnt).Select(x =>
- string.IsNullOrWhiteSpace(streamPrefix)
- ? Guid.NewGuid().ToString()
- : $"{streamPrefix}-{x}").ToArray();
-
- context.Log.Information("Writing streams randomly between {first} and {last}",
- streams.FirstOrDefault(),
- streams.LastOrDefault());
-
- var start = new TaskCompletionSource();
- stats.StartTime = DateTime.UtcNow;
- var sw2 = new Stopwatch();
- var capacity = 2000 / clientsCnt;
- var clientTasks = new List();
- for (int i = 0; i < clientsCnt; i++)
- {
- var count = requestsCnt / clientsCnt + ((i == clientsCnt - 1) ? requestsCnt % clientsCnt : 0);
-
- var client = context._clientApiTestClient.CreateConnection();
- await client.ConnectAsync();
- clientTasks.Add(RunClient(client, count));
- }
-
- async Task RunClient(IEventStoreConnection client, long count)
- {
- var rnd = new Random();
- List pending = new List(capacity);
- await start.Task;
-
- for (int j = 0; j < count; ++j)
- {
- var events = new EventData[batchSize];
- for (int q = 0; q < batchSize; q++)
- {
- events[q] = new EventData(Guid.NewGuid(),
- "TakeSomeSpaceEvent", false,
- Common.Utils.Helper.UTF8NoBom.GetBytes(
- "{ \"DATA\" : \"" + new string('*', size) + "\"}"),
- Common.Utils.Helper.UTF8NoBom.GetBytes(
- "{ \"METADATA\" : \"" + new string('$', 100) + "\"}"));
- }
-
- var corrid = Guid.NewGuid();
- monitor.StartOperation(corrid);
-
- pending.Add(client.AppendToStreamAsync(streams[rnd.Next(streamsCnt)], ExpectedVersion.Any, events)
- .ContinueWith(t =>
- {
- monitor.EndOperation(corrid);
- if (t.IsCompletedSuccessfully)
- {
- Interlocked.Add(ref stats.Succ, batchSize);
- }
- else
- {
- if (Interlocked.Increment(ref stats.Fail) % 1000 == 0)
- {
- Console.Write("#");
- }
-
- if (t.Exception != null)
- {
- var exception = t.Exception.Flatten();
- switch (exception.InnerException)
- {
- case WrongExpectedVersionException _:
- Interlocked.Increment(ref stats.WrongExpVersion);
- break;
- case StreamDeletedException _:
- Interlocked.Increment(ref stats.StreamDeleted);
- break;
- case OperationTimedOutException _:
- Interlocked.Increment(ref stats.CommitTimeout);
- break;
- }
- }
- }
-
- Interlocked.Add(ref stats.Succ, batchSize);
- if (stats.Succ - last > 1000)
- {
- last = stats.Succ;
- Console.Write(".");
- }
-
- var localAll = Interlocked.Add(ref stats.All, batchSize);
- if (localAll % 100000 == 0)
- {
- stats.Elapsed = sw2.Elapsed;
- stats.Rate = 1000.0 * 100000 / stats.Elapsed.TotalMilliseconds;
- sw2.Restart();
- context.Log.Debug(
- "\nDONE TOTAL {writes} WRITES IN {elapsed} ({rate:0.0}/s) [S:{success}, F:{failures} (WEV:{wrongExpectedVersion}, P:{prepareTimeout}, C:{commitTimeout}, F:{forwardTimeout}, D:{streamDeleted})].",
- localAll, stats.Elapsed, stats.Rate,
- stats.Succ, stats.Fail,
- stats.WrongExpVersion, stats.PrepTimeout, stats.CommitTimeout, stats.ForwardTimeout, stats.StreamDeleted);
- stats.WriteStatsToFile(context.StatsLogger);
- }
-
- if (localAll >= requestsCnt)
- {
- context.Success();
- doneEvent.Set();
- }
- }));
- if (pending.Count == capacity)
- {
- await Task.WhenAny(pending);
-
- while (pending.Count > 0 && Task.WhenAny(pending).IsCompleted)
- {
- pending.RemoveAll(x => x.IsCompleted);
- if (stats.Succ - last > 1000)
- {
- Console.Write(".");
- last = stats.Succ;
- }
- }
- }
- }
-
- if (pending.Count > 0)
- {
- await Task.WhenAll(pending);
- }
- }
-
- var sw = Stopwatch.StartNew();
- sw2.Start();
- start.SetResult();
- await Task.WhenAll(clientTasks);
- sw.Stop();
-
- clients.ForEach(client => client.Close());
-
- context.Log.Information(
- "Completed. Successes: {success}, failures: {failures} (WRONG VERSION: {wrongExpectedVersion}, P: {prepareTimeout}, C: {commitTimeout}, F: {forwardTimeout}, D: {streamDeleted})",
- stats.Succ, stats.Fail,
- stats.WrongExpVersion, stats.PrepTimeout, stats.CommitTimeout, stats.ForwardTimeout, stats.StreamDeleted);
- stats.WriteStatsToFile(context.StatsLogger);
-
- var reqPerSec = (stats.All + 0.0) / sw.ElapsedMilliseconds * 1000;
- context.Log.Information("{requests} requests completed in {elapsed}ms ({rate:0.00} reqs per sec).", stats.All,
- sw.ElapsedMilliseconds, reqPerSec);
-
- PerfUtils.LogData(
- Keyword,
- PerfUtils.Row(PerfUtils.Col("clientsCnt", clientsCnt),
- PerfUtils.Col("requestsCnt", requestsCnt),
- PerfUtils.Col("ElapsedMilliseconds", sw.ElapsedMilliseconds)),
- PerfUtils.Row(PerfUtils.Col("successes", stats.Succ), PerfUtils.Col("failures", stats.Fail)));
-
- var failuresRate = (int)(100 * stats.Fail / (stats.Fail + stats.Succ));
- PerfUtils.LogTeamCityGraphData(string.Format("{0}-{1}-{2}-reqPerSec", Keyword, clientsCnt, requestsCnt),
- (int)reqPerSec);
- PerfUtils.LogTeamCityGraphData(
- string.Format("{0}-{1}-{2}-failureSuccessRate", Keyword, clientsCnt, requestsCnt), failuresRate);
- PerfUtils.LogTeamCityGraphData(
- string.Format("{0}-c{1}-r{2}-st{3}-s{4}-reqPerSec", Keyword, clientsCnt, requestsCnt, streamsCnt,
- size),
- (int)reqPerSec);
- PerfUtils.LogTeamCityGraphData(
- string.Format("{0}-c{1}-r{2}-st{3}-s{4}-failureSuccessRate", Keyword, clientsCnt, requestsCnt,
- streamsCnt, size), failuresRate);
- monitor.GetMeasurementDetails();
- if (Interlocked.Read(ref stats.Succ) != requestsCnt)
- {
- context.Fail(reason: "There were errors or not all requests completed.");
- }
- else
- {
- context.Success();
- }
- }
-}
diff --git a/src/EventStore.TestClient/ClientApiTcpTestClient.cs b/src/EventStore.TestClient/ClientApiTcpTestClient.cs
deleted file mode 100644
index ec3781b506..0000000000
--- a/src/EventStore.TestClient/ClientApiTcpTestClient.cs
+++ /dev/null
@@ -1,44 +0,0 @@
-using EventStore.ClientAPI;
-using ILogger = Serilog.ILogger;
-
-namespace EventStore.TestClient;
-
-///
-/// A test client that connects using the legacy dotnet TCP client
-///
-public class ClientApiTcpTestClient
-{
- ///
- /// The options specified when starting the EventStore.TestClient
- ///
- public ClientOptions Options { get; set; }
- private ILogger _log;
-
- ///
- /// Constructs a new
- ///
- ///
- ///
- public ClientApiTcpTestClient(ClientOptions options, ILogger log)
- {
- Options = options;
- _log = log;
- }
-
- ///
- /// Creates a new TCP connection.
- ///
- ///
- public IEventStoreConnection CreateConnection()
- {
- var connectionString = string.IsNullOrWhiteSpace(Options.ConnectionString)
- ? $"ConnectTo=tcp://{Options.Host}:{Options.TcpPort};UseSslConnection={Options.UseTls};ValidateServer={Options.TlsValidateServer}"
- : Options.ConnectionString;
- _log.Debug("Creating TCP client with connection string '{connectionString}.", connectionString);
-
- var connectionSettings = ConnectionSettings.Create()
- .LimitRetriesForOperationTo(0)
- .KeepReconnecting();
- return EventStoreConnection.Create(connectionString, connectionSettings);
- }
-}
diff --git a/src/EventStore.TestClient/ClientOptions.cs b/src/EventStore.TestClient/ClientOptions.cs
index 88e7b51c7b..e92b6cfdcb 100644
--- a/src/EventStore.TestClient/ClientOptions.cs
+++ b/src/EventStore.TestClient/ClientOptions.cs
@@ -17,14 +17,12 @@ namespace EventStore.TestClient;
public sealed record ClientOptions
{
public string Host { get; init; }
- public int TcpPort { get; init; }
public int HttpPort { get; init; }
public int Timeout { get; init; }
public int ReadWindow { get; init; }
public int WriteWindow { get; init; }
public int PingWindow { get; init; }
public string[] Command { get; init; }
- public bool Reconnect { get; set; }
public bool UseTls { get; init; }
public bool TlsValidateServer { get; init; }
@@ -37,13 +35,11 @@ public ClientOptions()
{
Command = Array.Empty();
Host = IPAddress.Loopback.ToString();
- TcpPort = 1113;
HttpPort = 2113;
Timeout = -1;
ReadWindow = 2000;
WriteWindow = 2000;
PingWindow = 2000;
- Reconnect = true;
UseTls = false;
TlsValidateServer = false;
ConnectionString = string.Empty;
@@ -58,8 +54,12 @@ public override string ToString()
(builder, option) => builder.AppendLine($"{option.Name}: {GetValue(option)}"))
.ToString();
- object GetValue(PropertyInfo propertyInfo) => propertyInfo.PropertyType.IsArray
- ? string.Join(",", ((IEnumerable)propertyInfo.GetValue(this)).OfType