This is an automated email from the ASF dual-hosted git repository.
mmodzelewski pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iggy.git
The following commit(s) were added to refs/heads/master by this push:
new 8dde7e7f6 feat(java): let TCP clients run one I/O thread or share a
group (#4096)
8dde7e7f6 is described below
commit 8dde7e7f67ec09242207d6137255d17ad812136f
Author: Maciej Modzelewski <[email protected]>
AuthorDate: Wed Sep 9 07:32:41 2026 +0200
feat(java): let TCP clients run one I/O thread or share a group (#4096)
Every AsyncIggyTcpClient built its own event loop group with
Netty's default of 2 x CPU threads. Its pool holds one channel,
and a channel lives on one loop, so the other threads sat idle.
Nothing let clients share a group. An application that opened one
client per producer and consumer paid a thread count that grew
with the client count.
The group that a client creates for itself now has one thread by
default, sized through ioThreads on the builder. eventLoopGroup
registers the channel on a caller-owned group instead, and the
client never shuts that group down. The builder rejects a group
that cannot drive NIO channels or that already started to shut
down, and ignores ioThreads when a group is supplied. The
blocking builder forwards both options.
The connection takes one loop from the group at construction and
pins both the channel and the heartbeat to it. A next() call per
heartbeat tick rotates the timer across loops the channel never
uses. A second next() call, such as a pool handed the whole
group, takes a slot of the group's round-robin counter and lands
later connections off the sequence.
Shutting the owned group down swept every channel. A shared
group stays up, so close() now closes the tracked channels
itself before it closes the pool. The pool closes only idle
channels, and a login holds its lease until the reply, so the
pool alone left that channel open. If a caller shuts the group
down under a live connection, the heartbeat stops with a
warning and close() completes instead of failing on a
terminated loop.
Refs #4021
---
foreign/java/README.md | 30 ++
.../iggy/client/async/tcp/AsyncIggyTcpClient.java | 32 +-
.../async/tcp/AsyncIggyTcpClientBuilder.java | 79 +++-
.../iggy/client/async/tcp/AsyncTcpConnection.java | 88 +++-
.../client/blocking/tcp/IggyTcpClientBuilder.java | 31 ++
.../async/tcp/AsyncIggyTcpClientBuilderTest.java | 94 +++++
.../tcp/AsyncIggyTcpClientEventLoopGroupTest.java | 464 +++++++++++++++++++++
.../tcp/AsyncTcpConnectionConcurrencyTest.java | 4 +
8 files changed, 799 insertions(+), 23 deletions(-)
diff --git a/foreign/java/README.md b/foreign/java/README.md
index ec18b00c9..949d8c7cd 100644
--- a/foreign/java/README.md
+++ b/foreign/java/README.md
@@ -182,6 +182,36 @@ var client = Iggy.tcpClientBuilder()
.buildAndLogin();
```
+### Event Loop Threads
+
+Each TCP client drives a single connection, so by default it creates an event
loop group
+with one thread. An application that opens many clients can instead register
them all on
+one caller-owned group. The clients never shut that group down. Close the
clients first,
+then shut the group down:
+
+```java
+var group = new MultiThreadIoEventLoopGroup(2, NioIoHandler.newFactory());
+
+var producer = Iggy.tcpClientBuilder()
+ .blocking()
+ .eventLoopGroup(group)
+ .credentials("iggy", "iggy")
+ .buildAndLogin();
+var consumer = Iggy.tcpClientBuilder()
+ .blocking()
+ .eventLoopGroup(group)
+ .credentials("iggy", "iggy")
+ .buildAndLogin();
+
+// ... later
+producer.close();
+consumer.close();
+group.shutdownGracefully();
+```
+
+Do not block in a completion callback. Callbacks run on the group's loops, so
a blocked
+callback stalls every client that shares the group.
+
### Version Information
```java
diff --git
a/foreign/java/java-sdk/src/main/java/org/apache/iggy/client/async/tcp/AsyncIggyTcpClient.java
b/foreign/java/java-sdk/src/main/java/org/apache/iggy/client/async/tcp/AsyncIggyTcpClient.java
index d65622f53..e0a88922b 100644
---
a/foreign/java/java-sdk/src/main/java/org/apache/iggy/client/async/tcp/AsyncIggyTcpClient.java
+++
b/foreign/java/java-sdk/src/main/java/org/apache/iggy/client/async/tcp/AsyncIggyTcpClient.java
@@ -22,6 +22,7 @@ package org.apache.iggy.client.async.tcp;
import io.netty.buffer.ByteBuf;
import io.netty.buffer.Unpooled;
import io.netty.channel.ConnectTimeoutException;
+import io.netty.channel.IoEventLoopGroup;
import org.apache.iggy.IggyVersion;
import org.apache.iggy.client.ConnectionInfo;
import org.apache.iggy.client.async.ConsumerGroupsClient;
@@ -108,8 +109,12 @@ import java.util.stream.Stream;
* response handling is performed asynchronously.
*
* <h2>Resource Management</h2>
- * <p>Always call {@link #close()} when the client is no longer needed. This
shuts down
- * the Netty event loop group and releases all associated resources.
+ * <p>When the client is no longer needed, call {@link #close()}. This closes
the
+ * connection and shuts down the event loop group that the client created for
itself. A
+ * group supplied through {@link
AsyncIggyTcpClientBuilder#eventLoopGroup(IoEventLoopGroup)}
+ * stays up. The caller owns that group. After every client on the group is
closed, the
+ * caller shuts the group down. Do not block in a completion callback.
Callbacks run on
+ * that group's loops, so a blocked callback stalls every client that shares
the group.
*
* @see AsyncIggyTcpClientBuilder
* @see org.apache.iggy.Iggy#tcpClientBuilder()
@@ -141,6 +146,8 @@ public class AsyncIggyTcpClient {
private final Optional<RetryPolicy> retryPolicy;
private final boolean enableTls;
private final Optional<File> tlsCertificate;
+ private final Optional<IoEventLoopGroup> sharedEventLoopGroup;
+ private final int ioThreads;
private final TcpConnectionPoolConfig poolConfig;
private final ClientRoutingState routingState = new ClientRoutingState();
private final LoginRoutingHook loginRoutingHook = new LoginRoutingHook() {
@@ -207,7 +214,9 @@ public class AsyncIggyTcpClient {
VsrFrameDecoder.DEFAULT_MAX_FRAME_SIZE,
null,
false,
- Optional.empty());
+ Optional.empty(),
+ Optional.empty(),
+ AsyncTcpConnection.DEFAULT_IO_THREADS);
}
@SuppressWarnings("checkstyle:ParameterNumber")
@@ -223,7 +232,9 @@ public class AsyncIggyTcpClient {
int maxVsrFrameSize,
RetryPolicy retryPolicy,
boolean enableTls,
- Optional<File> tlsCertificate) {
+ Optional<File> tlsCertificate,
+ Optional<IoEventLoopGroup> sharedEventLoopGroup,
+ int ioThreads) {
this.connectionInfo = new ConnectionInfo(host, port);
this.seedConnectionInfo = this.connectionInfo;
this.username = Optional.ofNullable(username);
@@ -236,6 +247,8 @@ public class AsyncIggyTcpClient {
this.retryPolicy = Optional.ofNullable(retryPolicy);
this.enableTls = enableTls;
this.tlsCertificate = tlsCertificate;
+ this.sharedEventLoopGroup = sharedEventLoopGroup;
+ this.ioThreads = ioThreads;
var poolConfigBuilder = TcpConnectionPoolConfig.builder();
this.acquireTimeout.ifPresent(timeout ->
poolConfigBuilder.setAcquireTimeoutMillis(timeout.toMillis()));
@@ -460,10 +473,13 @@ public class AsyncIggyTcpClient {
}
/**
- * Closes the TCP connection and releases all Netty resources.
+ * Closes the TCP connection and releases the Netty resources that this
client owns.
*
- * <p>This shuts down the event loop group gracefully. After calling this
method,
- * the client cannot be reused — create a new instance if needed.
+ * <p>If the client created its own event loop group, it shuts that group
down gracefully.
+ * The client does not shut down a group supplied through
+ * {@link AsyncIggyTcpClientBuilder#eventLoopGroup(IoEventLoopGroup)}. The
caller shuts
+ * that group down. After you call this method, you cannot reuse the
client. If you need
+ * a client again, create a new instance.
*
* @return a {@link CompletableFuture} that completes when all resources
are released
*/
@@ -498,6 +514,8 @@ public class AsyncIggyTcpClient {
enableTls,
tlsCertificate,
poolConfig,
+ sharedEventLoopGroup,
+ ioThreads,
dialTimeout(),
requestTimeout,
heartbeatInterval,
diff --git
a/foreign/java/java-sdk/src/main/java/org/apache/iggy/client/async/tcp/AsyncIggyTcpClientBuilder.java
b/foreign/java/java-sdk/src/main/java/org/apache/iggy/client/async/tcp/AsyncIggyTcpClientBuilder.java
index e6a6ed510..95d983309 100644
---
a/foreign/java/java-sdk/src/main/java/org/apache/iggy/client/async/tcp/AsyncIggyTcpClientBuilder.java
+++
b/foreign/java/java-sdk/src/main/java/org/apache/iggy/client/async/tcp/AsyncIggyTcpClientBuilder.java
@@ -19,6 +19,10 @@
package org.apache.iggy.client.async.tcp;
+import io.netty.channel.IoEventLoop;
+import io.netty.channel.IoEventLoopGroup;
+import io.netty.channel.nio.NioIoHandle;
+import io.netty.util.concurrent.EventExecutor;
import org.apache.commons.lang3.StringUtils;
import org.apache.iggy.client.async.tcp.vsr.VsrFrameDecoder;
import org.apache.iggy.client.async.tcp.vsr.VsrHeaders;
@@ -61,6 +65,13 @@ import java.util.function.Function;
* .credentials("admin", "secret")
* .buildAndLogin()
* .join();
+ *
+ * // Many clients on one caller-owned event loop group
+ * var group = new MultiThreadIoEventLoopGroup(2, NioIoHandler.newFactory());
+ * var producer = AsyncIggyTcpClient.builder().eventLoopGroup(group).build();
+ * var consumer = AsyncIggyTcpClient.builder().eventLoopGroup(group).build();
+ * // ... close both clients, then:
+ * group.shutdownGracefully();
* }</pre>
*
* @see AsyncIggyTcpClient#builder()
@@ -78,6 +89,8 @@ public final class AsyncIggyTcpClientBuilder {
private Duration acquireTimeout;
private Duration heartbeatInterval = Duration.ofSeconds(5);
private long maxVsrFrameSize = VsrFrameDecoder.DEFAULT_MAX_FRAME_SIZE;
+ private int ioThreads = AsyncTcpConnection.DEFAULT_IO_THREADS;
+ private IoEventLoopGroup eventLoopGroup;
public AsyncIggyTcpClientBuilder() {}
@@ -227,6 +240,42 @@ public final class AsyncIggyTcpClientBuilder {
return this;
}
+ /**
+ * Sets the number of event loop threads in the group that the client
creates for itself.
+ *
+ * <p>The client drives a single channel, so the default of 1 is enough. If
+ * {@link #eventLoopGroup(IoEventLoopGroup)} is set, the client ignores
this value.
+ *
+ * @param ioThreads the event loop thread count, at least 1
+ * @return this builder
+ */
+ public AsyncIggyTcpClientBuilder ioThreads(int ioThreads) {
+ this.ioThreads = ioThreads;
+ return this;
+ }
+
+ /**
+ * Sets a caller-owned event loop group shared across clients.
+ *
+ * <p>The client pins its channel to one loop of the group and never shuts
the group down.
+ * After every client on the group is closed, the caller shuts the group
down. The group
+ * must drive NIO channels, for example
+ * {@code new MultiThreadIoEventLoopGroup(threads,
NioIoHandler.newFactory())}.
+ * Do not block in a completion callback. Callbacks run on the group's
loops, so a
+ * blocked callback stalls every client that shares the group.
+ *
+ * @param eventLoopGroup the group to register the client's channel on
+ * @return this builder
+ * @throws IggyInvalidArgumentException if the group is null
+ */
+ public AsyncIggyTcpClientBuilder eventLoopGroup(IoEventLoopGroup
eventLoopGroup) {
+ if (eventLoopGroup == null) {
+ throw new IggyInvalidArgumentException("EventLoopGroup cannot be
null");
+ }
+ this.eventLoopGroup = eventLoopGroup;
+ return this;
+ }
+
/**
* Builds and returns a configured AsyncIggyTcpClient instance.
* Note: You still need to call {@link AsyncIggyTcpClient#connect()} on
the returned client.
@@ -242,6 +291,8 @@ public final class AsyncIggyTcpClientBuilder {
validateRequestTimeout();
validateHeartbeatInterval();
validateMaxVsrFrameSize();
+ validateIoThreads();
+ validateEventLoopGroup();
return new AsyncIggyTcpClient(
host,
@@ -255,7 +306,9 @@ public final class AsyncIggyTcpClientBuilder {
(int) maxVsrFrameSize,
retryPolicy,
enableTls,
- Optional.ofNullable(tlsCertificate));
+ Optional.ofNullable(tlsCertificate),
+ Optional.ofNullable(eventLoopGroup),
+ ioThreads);
}
private void validateHost() {
@@ -308,6 +361,30 @@ public final class AsyncIggyTcpClientBuilder {
}
}
+ private void validateIoThreads() {
+ if (eventLoopGroup != null) {
+ return;
+ }
+ if (ioThreads < 1) {
+ throw new IggyInvalidArgumentException("IoThreads must be at least
1");
+ }
+ }
+
+ private void validateEventLoopGroup() {
+ if (eventLoopGroup == null) {
+ return;
+ }
+ for (EventExecutor executor : eventLoopGroup) {
+ if (!(executor instanceof IoEventLoop loop) ||
!loop.isCompatible(NioIoHandle.class)) {
+ throw new IggyInvalidArgumentException(
+ "EventLoopGroup must drive NIO channels, for example
MultiThreadIoEventLoopGroup with NioIoHandler");
+ }
+ }
+ if (eventLoopGroup.isShuttingDown()) {
+ throw new IggyInvalidArgumentException("EventLoopGroup shutdown
already started");
+ }
+ }
+
/**
* Builds, connects, and logs in using the provided credentials.
* This is a convenience method equivalent to calling {@code build()},
{@code connect()},
diff --git
a/foreign/java/java-sdk/src/main/java/org/apache/iggy/client/async/tcp/AsyncTcpConnection.java
b/foreign/java/java-sdk/src/main/java/org/apache/iggy/client/async/tcp/AsyncTcpConnection.java
index bdf97f90a..656103416 100644
---
a/foreign/java/java-sdk/src/main/java/org/apache/iggy/client/async/tcp/AsyncTcpConnection.java
+++
b/foreign/java/java-sdk/src/main/java/org/apache/iggy/client/async/tcp/AsyncTcpConnection.java
@@ -27,8 +27,11 @@ import io.netty.channel.ChannelFutureListener;
import io.netty.channel.ChannelOption;
import io.netty.channel.ChannelPipeline;
import io.netty.channel.ConnectTimeoutException;
+import io.netty.channel.EventLoop;
import io.netty.channel.IoEventLoopGroup;
import io.netty.channel.MultiThreadIoEventLoopGroup;
+import io.netty.channel.group.ChannelGroup;
+import io.netty.channel.group.DefaultChannelGroup;
import io.netty.channel.nio.NioIoHandler;
import io.netty.channel.pool.AbstractChannelPoolHandler;
import io.netty.channel.pool.ChannelHealthChecker;
@@ -37,7 +40,10 @@ import io.netty.channel.socket.nio.NioSocketChannel;
import io.netty.handler.ssl.SslContext;
import io.netty.handler.ssl.SslContextBuilder;
import io.netty.handler.ssl.SslHandler;
+import io.netty.util.concurrent.DefaultThreadFactory;
+import io.netty.util.concurrent.Future;
import io.netty.util.concurrent.FutureListener;
+import io.netty.util.concurrent.GlobalEventExecutor;
import io.netty.util.concurrent.ScheduledFuture;
import org.apache.iggy.client.ConnectionInfo;
import org.apache.iggy.client.async.tcp.vsr.ConsensusSession;
@@ -91,6 +97,8 @@ public class AsyncTcpConnection {
// and a transient one is not a rejected credential.
static final int TRANSIENT_NOT_COMMITTED = 57;
static final int TRANSIENT_NOT_ACCEPTED = 58;
+ // The pool holds one channel, and one channel lives on one loop.
+ static final int DEFAULT_IO_THREADS = 1;
private static final Logger log =
LoggerFactory.getLogger(AsyncTcpConnection.class);
private static final Duration DEFAULT_CONNECTION_TIMEOUT =
Duration.ofMillis(3000);
// A missing reply must not hold the single VSR-pinned channel forever.
@@ -98,9 +106,13 @@ public class AsyncTcpConnection {
private static final long TRANSIENT_RETRY_INTERVAL_MS = 50;
private static final Duration TRANSIENT_RETRY_BUDGET =
Duration.ofSeconds(30);
private static final Duration NOT_ACCEPTED_RETRY_BUDGET =
Duration.ofSeconds(2);
+ private static final String EVENT_LOOP_THREAD_PREFIX = "iggy-tcp-io";
private final IoEventLoopGroup eventLoopGroup;
+ private final boolean ownsEventLoopGroup;
+ private final EventLoop eventLoop;
private final FixedChannelPool channelPool;
+ private final ChannelGroup channels = new
DefaultChannelGroup(GlobalEventExecutor.INSTANCE, true);
private final AtomicBoolean isClosed = new AtomicBoolean(false);
private final AtomicLong authGeneration = new AtomicLong(0);
private final VsrRequestEncoder vsrEncoder;
@@ -130,6 +142,8 @@ public class AsyncTcpConnection {
enableTls,
tlsCertificate,
poolConfig,
+ Optional.empty(),
+ DEFAULT_IO_THREADS,
connectionTimeout,
Optional.empty(),
Duration.ofSeconds(5),
@@ -146,6 +160,8 @@ public class AsyncTcpConnection {
boolean enableTls,
Optional<File> tlsCertificate,
TcpConnectionPoolConfig poolConfig,
+ Optional<IoEventLoopGroup> sharedEventLoopGroup,
+ int ioThreads,
Optional<Duration> connectionTimeout,
Optional<Duration> requestTimeout,
Duration heartbeatInterval,
@@ -171,12 +187,17 @@ public class AsyncTcpConnection {
ConsensusSession consensusSession = new ConsensusSession();
this.vsrEncoder = new VsrRequestEncoder(consensusSession);
- this.eventLoopGroup = new
MultiThreadIoEventLoopGroup(NioIoHandler.newFactory());
+ this.ownsEventLoopGroup = sharedEventLoopGroup.isEmpty();
+ this.eventLoopGroup = sharedEventLoopGroup.orElseGet(() -> new
MultiThreadIoEventLoopGroup(
+ ioThreads,
+ new DefaultThreadFactory(EVENT_LOOP_THREAD_PREFIX, false,
Thread.MAX_PRIORITY),
+ NioIoHandler.newFactory()));
+ this.eventLoop = eventLoopGroup.next();
long dialTimeoutMillis =
connectionTimeout.orElse(DEFAULT_CONNECTION_TIMEOUT).toMillis();
var bootstrap = new Bootstrap()
- .group(eventLoopGroup)
+ .group(eventLoop)
.channel(NioSocketChannel.class)
.option(ChannelOption.TCP_NODELAY, true)
.option(ChannelOption.CONNECT_TIMEOUT_MILLIS, (int)
dialTimeoutMillis)
@@ -196,7 +217,8 @@ public class AsyncTcpConnection {
dialTimeoutMillis,
consensusSession,
maxVsrFrameSize,
- this::onSessionEvicted),
+ this::onSessionEvicted,
+ channels::add),
ChannelHealthChecker.ACTIVE,
FixedChannelPool.AcquireTimeoutAction.FAIL,
poolConfig.getAcquireTimeoutMillis(),
@@ -248,8 +270,13 @@ public class AsyncTcpConnection {
if (!heartbeatRunning || isClosed.get()) {
return;
}
- heartbeatTask =
- eventLoopGroup.next().schedule(this::sendHeartbeat,
heartbeatIntervalNanos, TimeUnit.NANOSECONDS);
+ try {
+ heartbeatTask = eventLoop.schedule(this::sendHeartbeat,
heartbeatIntervalNanos, TimeUnit.NANOSECONDS);
+ } catch (RejectedExecutionException loopGone) {
+ // Only a caller-owned group shuts down under a live
connection.
+ heartbeatRunning = false;
+ log.warn("Event loop rejected the heartbeat, stopping it: {}",
loopGone.getMessage());
+ }
}
}
@@ -289,6 +316,16 @@ public class AsyncTcpConnection {
}
}
+ boolean heartbeatScheduled() {
+ synchronized (heartbeatLock) {
+ return heartbeatTask != null;
+ }
+ }
+
+ EventLoop eventLoop() {
+ return eventLoop;
+ }
+
public <T> CompletableFuture<T> exchangeForEntity(
CommandCode commandCode, ByteBuf payload, Function<ByteBuf, T>
func) {
return send(commandCode, payload).thenApply(response -> {
@@ -920,18 +957,35 @@ public class AsyncTcpConnection {
stopHeartbeat();
releaseLoginPayload();
CompletableFuture<Void> shutdownFuture = new CompletableFuture<>();
- channelPool
- .closeAsync()
- .addListener(f ->
eventLoopGroup.shutdownGracefully().addListener(sf -> {
- if (sf.isSuccess()) {
- shutdownFuture.complete(null);
- } else {
- shutdownFuture.completeExceptionally(sf.cause());
- }
- }));
+ channels.close().addListener(channelsClosed ->
closePool(shutdownFuture));
return shutdownFuture;
}
+ private void closePool(CompletableFuture<Void> shutdownFuture) {
+ try {
+ channelPool.closeAsync().addListener(poolClosed -> {
+ if (!ownsEventLoopGroup) {
+ completeShutdown(shutdownFuture, poolClosed);
+ return;
+ }
+ eventLoopGroup
+ .shutdownGracefully()
+ .addListener(groupClosed ->
completeShutdown(shutdownFuture, groupClosed));
+ });
+ } catch (RejectedExecutionException loopGone) {
+ log.warn("Event loop rejected the pool close, channel already
gone: {}", loopGone.getMessage());
+ shutdownFuture.complete(null);
+ }
+ }
+
+ private static void completeShutdown(CompletableFuture<Void>
shutdownFuture, Future<?> step) {
+ if (step.isSuccess()) {
+ shutdownFuture.complete(null);
+ } else {
+ shutdownFuture.completeExceptionally(step.cause());
+ }
+ }
+
private static final class PoolChannelHandler extends
AbstractChannelPoolHandler {
private final String host;
private final int port;
@@ -941,6 +995,7 @@ public class AsyncTcpConnection {
private final ConsensusSession consensusSession;
private final int maxVsrFrameSize;
private final IntConsumer onEviction;
+ private final Consumer<Channel> onChannelCreated;
@SuppressWarnings("checkstyle:ParameterNumber")
PoolChannelHandler(
@@ -951,7 +1006,8 @@ public class AsyncTcpConnection {
long dialTimeoutMillis,
ConsensusSession consensusSession,
int maxVsrFrameSize,
- IntConsumer onEviction) {
+ IntConsumer onEviction,
+ Consumer<Channel> onChannelCreated) {
this.host = host;
this.port = port;
this.enableTls = enableTls;
@@ -960,10 +1016,12 @@ public class AsyncTcpConnection {
this.consensusSession = consensusSession;
this.maxVsrFrameSize = maxVsrFrameSize;
this.onEviction = onEviction;
+ this.onChannelCreated = onChannelCreated;
}
@Override
public void channelCreated(Channel ch) {
+ onChannelCreated.accept(ch);
ChannelPipeline pipeline = ch.pipeline();
if (enableTls) {
SslHandler ssl = sslContext.newHandler(ch.alloc(), host, port);
diff --git
a/foreign/java/java-sdk/src/main/java/org/apache/iggy/client/blocking/tcp/IggyTcpClientBuilder.java
b/foreign/java/java-sdk/src/main/java/org/apache/iggy/client/blocking/tcp/IggyTcpClientBuilder.java
index 65a09d215..e79ecfa3a 100644
---
a/foreign/java/java-sdk/src/main/java/org/apache/iggy/client/blocking/tcp/IggyTcpClientBuilder.java
+++
b/foreign/java/java-sdk/src/main/java/org/apache/iggy/client/blocking/tcp/IggyTcpClientBuilder.java
@@ -19,6 +19,7 @@
package org.apache.iggy.client.blocking.tcp;
+import io.netty.channel.IoEventLoopGroup;
import org.apache.iggy.client.async.tcp.AsyncIggyTcpClient;
import org.apache.iggy.client.async.tcp.AsyncIggyTcpClientBuilder;
import org.apache.iggy.config.RetryPolicy;
@@ -136,6 +137,36 @@ public final class IggyTcpClientBuilder {
return this;
}
+ /**
+ * Sets the number of event loop threads in the group that the client
creates for itself.
+ *
+ * <p>The client drives a single channel, so the default of 1 is enough. If
+ * {@link #eventLoopGroup(IoEventLoopGroup)} is set, the client ignores
this value.
+ *
+ * @param ioThreads the event loop thread count, at least 1
+ * @return this builder
+ */
+ public IggyTcpClientBuilder ioThreads(int ioThreads) {
+ asyncBuilder.ioThreads(ioThreads);
+ return this;
+ }
+
+ /**
+ * Sets a caller-owned event loop group shared across clients.
+ *
+ * <p>The client registers its channel on the group and never shuts the
group down.
+ * After every client on the group is closed, the caller shuts the group
down. The group
+ * must drive NIO channels.
+ *
+ * @param eventLoopGroup the group to register the client's channel on
+ * @return this builder
+ * @throws org.apache.iggy.exception.IggyInvalidArgumentException if the
group is null
+ */
+ public IggyTcpClientBuilder eventLoopGroup(IoEventLoopGroup
eventLoopGroup) {
+ asyncBuilder.eventLoopGroup(eventLoopGroup);
+ return this;
+ }
+
/**
* Enables or disables TLS for the TCP connection.
*
diff --git
a/foreign/java/java-sdk/src/test/java/org/apache/iggy/client/async/tcp/AsyncIggyTcpClientBuilderTest.java
b/foreign/java/java-sdk/src/test/java/org/apache/iggy/client/async/tcp/AsyncIggyTcpClientBuilderTest.java
index ce7234a35..248486b2d 100644
---
a/foreign/java/java-sdk/src/test/java/org/apache/iggy/client/async/tcp/AsyncIggyTcpClientBuilderTest.java
+++
b/foreign/java/java-sdk/src/test/java/org/apache/iggy/client/async/tcp/AsyncIggyTcpClientBuilderTest.java
@@ -19,6 +19,10 @@
package org.apache.iggy.client.async.tcp;
+import io.netty.channel.IoEventLoopGroup;
+import io.netty.channel.MultiThreadIoEventLoopGroup;
+import io.netty.channel.local.LocalIoHandler;
+import io.netty.channel.nio.NioIoHandler;
import org.apache.iggy.client.BaseIntegrationTest;
import org.apache.iggy.config.RetryPolicy;
import org.apache.iggy.exception.IggyAuthenticationException;
@@ -250,6 +254,96 @@ class AsyncIggyTcpClientBuilderTest extends
BaseIntegrationTest {
.isInstanceOf(IggyInvalidArgumentException.class);
}
+ @Test
+ void shouldRejectNonPositiveIoThreads() {
+ assertThatThrownBy(() ->
AsyncIggyTcpClient.builder().ioThreads(0).build())
+ .isInstanceOf(IggyInvalidArgumentException.class);
+ assertThatThrownBy(() ->
AsyncIggyTcpClient.builder().ioThreads(-1).build())
+ .isInstanceOf(IggyInvalidArgumentException.class);
+ }
+
+ @Test
+ void shouldIgnoreIoThreadsWhenAnEventLoopGroupIsSupplied() throws
Exception {
+ IoEventLoopGroup group = new MultiThreadIoEventLoopGroup(1,
NioIoHandler.newFactory());
+ try {
+
AsyncIggyTcpClient.builder().ioThreads(0).eventLoopGroup(group).build();
+ } finally {
+ group.shutdownGracefully(0, 1, TimeUnit.SECONDS).get(5,
TimeUnit.SECONDS);
+ }
+ }
+
+ @Test
+ void shouldRejectNullEventLoopGroup() {
+ assertThatThrownBy(() ->
AsyncIggyTcpClient.builder().eventLoopGroup(null))
+ .isInstanceOf(IggyInvalidArgumentException.class);
+ }
+
+ @Test
+ void shouldRejectAnEventLoopGroupThatCannotDriveNioChannels() throws
Exception {
+ IoEventLoopGroup localGroup = new MultiThreadIoEventLoopGroup(1,
LocalIoHandler.newFactory());
+ try {
+ assertThatThrownBy(() -> AsyncIggyTcpClient.builder()
+ .eventLoopGroup(localGroup)
+ .build())
+ .isInstanceOf(IggyInvalidArgumentException.class)
+ .hasMessageContaining("NIO");
+ } finally {
+ localGroup.shutdownGracefully(0, 1, TimeUnit.SECONDS).get(5,
TimeUnit.SECONDS);
+ }
+ }
+
+ @Test
+ void shouldRejectAShutDownEventLoopGroup() throws Exception {
+ IoEventLoopGroup group = new MultiThreadIoEventLoopGroup(1,
NioIoHandler.newFactory());
+ group.shutdownGracefully(0, 1, TimeUnit.SECONDS).get(5,
TimeUnit.SECONDS);
+
+ assertThatThrownBy(
+ () ->
AsyncIggyTcpClient.builder().eventLoopGroup(group).build())
+ .isInstanceOf(IggyInvalidArgumentException.class);
+ }
+
+ @Test
+ void shouldShareOneEventLoopGroupAcrossClients() throws Exception {
+ IoEventLoopGroup group = new MultiThreadIoEventLoopGroup(1,
NioIoHandler.newFactory());
+ AsyncIggyTcpClient other = null;
+ try {
+ client = AsyncIggyTcpClient.builder()
+ .host(serverHost())
+ .port(serverTcpPort())
+ .eventLoopGroup(group)
+ .build();
+ other = AsyncIggyTcpClient.builder()
+ .host(serverHost())
+ .port(serverTcpPort())
+ .eventLoopGroup(group)
+ .build();
+ client.connect().get(TEST_TIMEOUT_SECONDS, TimeUnit.SECONDS);
+ other.connect().get(TEST_TIMEOUT_SECONDS, TimeUnit.SECONDS);
+ client.users().login(TEST_USERNAME,
TEST_PASSWORD).get(TEST_TIMEOUT_SECONDS, TimeUnit.SECONDS);
+ other.users().login(TEST_USERNAME,
TEST_PASSWORD).get(TEST_TIMEOUT_SECONDS, TimeUnit.SECONDS);
+ assertThat(client.streams().getStreams().get(TEST_TIMEOUT_SECONDS,
TimeUnit.SECONDS))
+ .isNotNull();
+ assertThat(other.streams().getStreams().get(TEST_TIMEOUT_SECONDS,
TimeUnit.SECONDS))
+ .isNotNull();
+
+ client.close().get(TEST_TIMEOUT_SECONDS, TimeUnit.SECONDS);
+ client = null;
+
+ assertThat(group.isShuttingDown())
+ .as("closing a client must not take the shared group down")
+ .isFalse();
+ assertThat(other.streams().getStreams().get(TEST_TIMEOUT_SECONDS,
TimeUnit.SECONDS))
+ .as("the other client keeps working on the shared group")
+ .isNotNull();
+ } finally {
+ if (other != null) {
+ other.close().get(TEST_TIMEOUT_SECONDS, TimeUnit.SECONDS);
+ }
+ assertThat(group.isShuttingDown()).isFalse();
+ group.shutdownGracefully(0, 1,
TimeUnit.SECONDS).get(TEST_TIMEOUT_SECONDS, TimeUnit.SECONDS);
+ }
+ }
+
@Test
void shouldMaintainBackwardCompatibilityWithOldConstructor() throws
Exception {
// Given: Old constructor approach
diff --git
a/foreign/java/java-sdk/src/test/java/org/apache/iggy/client/async/tcp/AsyncIggyTcpClientEventLoopGroupTest.java
b/foreign/java/java-sdk/src/test/java/org/apache/iggy/client/async/tcp/AsyncIggyTcpClientEventLoopGroupTest.java
new file mode 100644
index 000000000..24b0146dd
--- /dev/null
+++
b/foreign/java/java-sdk/src/test/java/org/apache/iggy/client/async/tcp/AsyncIggyTcpClientEventLoopGroupTest.java
@@ -0,0 +1,464 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.iggy.client.async.tcp;
+
+import io.netty.buffer.ByteBuf;
+import io.netty.buffer.Unpooled;
+import io.netty.channel.EventLoop;
+import io.netty.channel.IoEventLoopGroup;
+import io.netty.channel.MultiThreadIoEventLoopGroup;
+import io.netty.channel.nio.NioIoHandler;
+import org.apache.iggy.client.ConnectionInfo;
+import org.junit.jupiter.api.Test;
+
+import java.io.EOFException;
+import java.io.IOException;
+import java.io.InputStream;
+import java.io.OutputStream;
+import java.net.InetAddress;
+import java.net.ServerSocket;
+import java.net.Socket;
+import java.nio.ByteBuffer;
+import java.nio.ByteOrder;
+import java.nio.charset.StandardCharsets;
+import java.time.Duration;
+import java.util.ArrayList;
+import java.util.List;
+import java.util.Optional;
+import java.util.Set;
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.CopyOnWriteArrayList;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.TimeoutException;
+import java.util.concurrent.atomic.AtomicInteger;
+import java.util.function.BooleanSupplier;
+import java.util.stream.Collectors;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatCode;
+
+/**
+ * One client drives one channel, so the group that it creates for itself runs
one loop
+ * and shuts down with the client. If the caller supplies a group, the client
pins its
+ * channel to one loop of that group and never shuts the group down.
+ */
+class AsyncIggyTcpClientEventLoopGroupTest {
+
+ private static final String OWNED_THREAD_PREFIX = "iggy-tcp-io-";
+ private static final Duration HEARTBEAT_INTERVAL = Duration.ofMillis(25);
+ // Graceful shutdown waits a two-second quiet period before the loop
thread exits.
+ private static final Duration SHUTDOWN_TIMEOUT = Duration.ofSeconds(15);
+
+ private static final int HEADER_SIZE = 256;
+ private static final int SIZE_OFFSET = 48;
+ private static final int COMMAND_OFFSET = 60;
+ private static final int REQUEST_ID_OFFSET = 168;
+ private static final int REQUEST_OPERATION_OFFSET = 176;
+ private static final int REQUEST_CODE_OFFSET = 196;
+ private static final int REPLY_REQUEST_ID_OFFSET = 200;
+ private static final int REPLY_OPERATION_OFFSET = 208;
+ private static final int REPLY_STATUS_OFFSET = 216;
+ private static final int COMMAND_REPLY = 8;
+ private static final int OPERATION_REGISTER = 1;
+ private static final int OPERATION_NON_REPLICATED = 2;
+ private static final int PING_CODE = 1;
+ private static final int LOGIN_CODE = 38;
+
+ @Test
+ void shouldRunOneEventLoopThreadPerClientByDefault() throws Exception {
+ try (MockVsrServer server = MockVsrServer.start()) {
+ Set<Thread> before = ownedEventLoopThreads();
+ AsyncIggyTcpClient client = builder(server).build();
+ client.connect().get(5, TimeUnit.SECONDS);
+
+ // Enough ticks for a round-robin heartbeat to wake more than one
loop.
+ server.awaitPings(8, Duration.ofSeconds(5));
+ assertThat(ownedEventLoopThreadsSince(before))
+ .as("one channel needs one loop")
+ .isEqualTo(1);
+
+ client.close().get(SHUTDOWN_TIMEOUT.toSeconds(), TimeUnit.SECONDS);
+ awaitOwnedEventLoopThreadsSince(before, 0, SHUTDOWN_TIMEOUT);
+ }
+ }
+
+ @Test
+ void shouldLeaveASharedGroupRunningWhenClientsClose() throws Exception {
+ IoEventLoopGroup group = new MultiThreadIoEventLoopGroup(1,
NioIoHandler.newFactory());
+ try (MockVsrServer server = MockVsrServer.start()) {
+ Set<Thread> before = ownedEventLoopThreads();
+ AsyncIggyTcpClient first =
builder(server).eventLoopGroup(group).build();
+ AsyncIggyTcpClient second =
builder(server).eventLoopGroup(group).build();
+ first.connect().get(5, TimeUnit.SECONDS);
+ second.connect().get(5, TimeUnit.SECONDS);
+ assertThat(ownedEventLoopThreadsSince(before))
+ .as("clients on a shared group create no group of their
own")
+ .isZero();
+
+ first.close().get(5, TimeUnit.SECONDS);
+ assertThat(group.isShuttingDown()).isFalse();
+ second.sendBinaryRequest(PING_CODE, new byte[0]).get(5,
TimeUnit.SECONDS);
+
+ second.close().get(5, TimeUnit.SECONDS);
+ assertThat(group.isShuttingDown()).isFalse();
+ } finally {
+ group.shutdownGracefully(0, 1, TimeUnit.SECONDS).get(5,
TimeUnit.SECONDS);
+ }
+ }
+
+ @Test
+ void shouldPinEachConnectionOnASharedGroupToTheNextLoop() throws Exception
{
+ IoEventLoopGroup group = new MultiThreadIoEventLoopGroup(4,
NioIoHandler.newFactory());
+ List<EventLoop> loops = new ArrayList<>();
+ group.forEach(executor -> loops.add((EventLoop) executor));
+ List<AsyncTcpConnection> connections = new ArrayList<>();
+ try (MockVsrServer server = MockVsrServer.start()) {
+ for (int i = 0; i < loops.size(); i++) {
+ AsyncTcpConnection connection = connection(server, group);
+ connection.connect().get(5, TimeUnit.SECONDS);
+ connections.add(connection);
+ }
+ server.awaitPings(loops.size() * 2, Duration.ofSeconds(5));
+
+ // A fresh group hands its loops out round-robin from the first.
Every
+ // extra next() call per connection, such as a pool handed the
whole
+ // group, skips loops and lands later connections off this
sequence.
+
assertThat(connections.stream().map(AsyncTcpConnection::eventLoop).toList())
+ .as("one connection takes one slot of the group's
round-robin chooser")
+ .containsExactlyElementsOf(loops);
+ } finally {
+ for (AsyncTcpConnection connection : connections) {
+ connection.close().get(5, TimeUnit.SECONDS);
+ }
+ group.shutdownGracefully(0, 1, TimeUnit.SECONDS).get(5,
TimeUnit.SECONDS);
+ }
+ }
+
+ @Test
+ void shouldCloseQuietlyAfterTheCallerShutTheSharedGroupDown() throws
Exception {
+ IoEventLoopGroup group = new MultiThreadIoEventLoopGroup(1,
NioIoHandler.newFactory());
+ try (MockVsrServer server = MockVsrServer.start()) {
+ AsyncIggyTcpClient client =
builder(server).eventLoopGroup(group).build();
+ client.connect().get(5, TimeUnit.SECONDS);
+ group.shutdownGracefully(0, 1, TimeUnit.SECONDS).get(5,
TimeUnit.SECONDS);
+
+ // The pool closes on its loop, and a terminated loop rejects that
task.
+ assertThatCode(() -> client.close().get(5, TimeUnit.SECONDS))
+ .as("nothing is left to release once the group took the
channel down")
+ .doesNotThrowAnyException();
+ }
+ }
+
+ @Test
+ void shouldCloseTheChannelHeldByAPendingLoginOnASharedGroup() throws
Exception {
+ IoEventLoopGroup group = new MultiThreadIoEventLoopGroup(1,
NioIoHandler.newFactory());
+ try (MockVsrServer server = MockVsrServer.start()) {
+ server.withholdRegisterReplies();
+ // Far past the test bound, so only close() can settle the login.
+ AsyncTcpConnection connection = connection(server, group,
Duration.ofMinutes(5));
+ connection.connect().get(5, TimeUnit.SECONDS);
+ CompletableFuture<ByteBuf> login = connection.send(LOGIN_CODE,
loginPayload());
+ await(() -> server.registers() == 1, "the server holds the login",
Duration.ofSeconds(5));
+
+ connection.close().get(5, TimeUnit.SECONDS);
+
+ // The pool only closes idle channels, and a login holds its lease
+ // until the reply, so close() must reach that channel itself.
+ await(() -> server.closedSockets() == 1, "the login's socket is
closed", Duration.ofSeconds(5));
+ await(login::isDone, "the pending login settles",
Duration.ofSeconds(5));
+ assertThat(login).isCompletedExceptionally();
+ assertThat(group.isShuttingDown()).isFalse();
+ } finally {
+ group.shutdownGracefully(0, 1, TimeUnit.SECONDS).get(5,
TimeUnit.SECONDS);
+ }
+ }
+
+ @Test
+ void
shouldCancelThePendingHeartbeatOnASharedGroupWhenTheConnectionCloses() throws
Exception {
+ IoEventLoopGroup group = new MultiThreadIoEventLoopGroup(1,
NioIoHandler.newFactory());
+ try (MockVsrServer server = MockVsrServer.start()) {
+ AsyncTcpConnection connection = connection(server, group);
+ connection.connect().get(5, TimeUnit.SECONDS);
+ server.awaitPings(2, Duration.ofSeconds(5));
+ await(connection::heartbeatScheduled, "a live connection keeps a
tick armed", Duration.ofSeconds(5));
+
+ connection.close().get(5, TimeUnit.SECONDS);
+
+ // A stray timer on a shared loop cannot be seen through the
server: the
+ // channel is gone, so a tick that still fired would fail before
sending.
+ assertThat(connection.heartbeatScheduled())
+ .as("close cancels the armed tick instead of leaving it to
the shared loop")
+ .isFalse();
+ Thread.sleep(HEARTBEAT_INTERVAL.multipliedBy(4).toMillis());
+ assertThat(connection.heartbeatScheduled())
+ .as("nothing re-arms the heartbeat after close")
+ .isFalse();
+ assertThat(group.isShuttingDown()).isFalse();
+ } finally {
+ group.shutdownGracefully(0, 1, TimeUnit.SECONDS).get(5,
TimeUnit.SECONDS);
+ }
+ }
+
+ @Test
+ void shouldReleaseTheGroupOfAConnectionReplacedByRetarget() throws
Exception {
+ try (MockVsrServer primary = MockVsrServer.start();
+ MockVsrServer survivor = MockVsrServer.start()) {
+ Set<Thread> before = ownedEventLoopThreads();
+ AsyncIggyTcpClient client = builder(primary).build();
+ client.connect().get(5, TimeUnit.SECONDS);
+
+ client.retarget(new ConnectionInfo(loopback(),
survivor.port())).get(5, TimeUnit.SECONDS);
+
assertThat(client.getConnectionInfo().port()).isEqualTo(survivor.port());
+
+ awaitOwnedEventLoopThreadsSince(before, 1, SHUTDOWN_TIMEOUT);
+ client.close().get(SHUTDOWN_TIMEOUT.toSeconds(), TimeUnit.SECONDS);
+ awaitOwnedEventLoopThreadsSince(before, 0, SHUTDOWN_TIMEOUT);
+ }
+ }
+
+ private static AsyncIggyTcpClientBuilder builder(MockVsrServer server) {
+ return AsyncIggyTcpClient.builder()
+ .host(loopback())
+ .port(server.port())
+ .heartbeatInterval(HEARTBEAT_INTERVAL)
+ .requestTimeout(Duration.ofSeconds(1));
+ }
+
+ private static AsyncTcpConnection connection(MockVsrServer server,
IoEventLoopGroup group) {
+ return connection(server, group, Duration.ofSeconds(1));
+ }
+
+ private static AsyncTcpConnection connection(
+ MockVsrServer server, IoEventLoopGroup group, Duration
requestTimeout) {
+ return new AsyncTcpConnection(
+ loopback(),
+ server.port(),
+ false,
+ Optional.empty(),
+ new AsyncTcpConnection.TcpConnectionPoolConfig(1000, 3000),
+ Optional.of(group),
+ 1,
+ Optional.of(Duration.ofSeconds(1)),
+ Optional.of(requestTimeout),
+ HEARTBEAT_INTERVAL,
+ 1024 * 1024,
+ null,
+ errorCode -> {},
+ ignored -> {});
+ }
+
+ private static String loopback() {
+ return InetAddress.getLoopbackAddress().getHostAddress();
+ }
+
+ private static ByteBuf loginPayload() {
+ ByteBuf payload = Unpooled.buffer();
+ payload.writeByte(4);
+ payload.writeBytes("iggy".getBytes(StandardCharsets.UTF_8));
+ payload.writeByte(4);
+ payload.writeBytes("iggy".getBytes(StandardCharsets.UTF_8));
+ payload.writeIntLE(0);
+ payload.writeIntLE(0);
+ return payload;
+ }
+
+ private static Set<Thread> ownedEventLoopThreads() {
+ return Thread.getAllStackTraces().keySet().stream()
+ .filter(Thread::isAlive)
+ .filter(thread ->
thread.getName().startsWith(OWNED_THREAD_PREFIX))
+ .collect(Collectors.toSet());
+ }
+
+ // Loops that earlier tests left draining die on their own schedule, so
+ // only threads born after the snapshot count.
+ private static int ownedEventLoopThreadsSince(Set<Thread> before) {
+ return (int) ownedEventLoopThreads().stream()
+ .filter(thread -> !before.contains(thread))
+ .count();
+ }
+
+ private static void awaitOwnedEventLoopThreadsSince(Set<Thread> before,
int expected, Duration timeout)
+ throws Exception {
+ await(
+ () -> ownedEventLoopThreadsSince(before) == expected,
+ "Expected " + expected + " owned event loop threads started by
this test",
+ timeout);
+ }
+
+ private static void await(BooleanSupplier condition, String description,
Duration timeout) throws Exception {
+ long deadline = System.nanoTime() + timeout.toNanos();
+ while (!condition.getAsBoolean()) {
+ if (System.nanoTime() > deadline) {
+ throw new TimeoutException(description + " within " + timeout);
+ }
+ Thread.sleep(20);
+ }
+ }
+
+ /**
+ * Replies success to every request and counts the pings that it saw
across all sockets.
+ * Register replies can be withheld to keep a login pending on the client.
+ */
+ private static final class MockVsrServer implements AutoCloseable {
+ private final ServerSocket serverSocket;
+ private final List<Socket> accepted = new CopyOnWriteArrayList<>();
+ private final AtomicInteger pings = new AtomicInteger();
+ private final AtomicInteger registers = new AtomicInteger();
+ private final AtomicInteger closedSockets = new AtomicInteger();
+ private volatile boolean withholdRegisterReplies;
+ private volatile boolean closed;
+
+ private MockVsrServer(ServerSocket serverSocket) {
+ this.serverSocket = serverSocket;
+ }
+
+ static MockVsrServer start() throws IOException {
+ MockVsrServer server = new MockVsrServer(new ServerSocket(0, 4,
InetAddress.getLoopbackAddress()));
+ Thread acceptor = new Thread(server::acceptLoop,
"mock-vsr-acceptor-" + server.port());
+ acceptor.setDaemon(true);
+ acceptor.start();
+ return server;
+ }
+
+ int port() {
+ return serverSocket.getLocalPort();
+ }
+
+ int pings() {
+ return pings.get();
+ }
+
+ int registers() {
+ return registers.get();
+ }
+
+ int closedSockets() {
+ return closedSockets.get();
+ }
+
+ void withholdRegisterReplies() {
+ withholdRegisterReplies = true;
+ }
+
+ void awaitPings(int expected, Duration timeout) throws Exception {
+ long deadline = System.nanoTime() + timeout.toNanos();
+ while (pings.get() < expected) {
+ if (System.nanoTime() > deadline) {
+ throw new TimeoutException("Expected " + expected + "
pings, saw " + pings.get());
+ }
+ Thread.sleep(5);
+ }
+ }
+
+ private void acceptLoop() {
+ while (!closed) {
+ try {
+ Socket socket = serverSocket.accept();
+ accepted.add(socket);
+ Thread exchange = new Thread(() -> exchange(socket),
"mock-vsr-exchange-" + port());
+ exchange.setDaemon(true);
+ exchange.start();
+ } catch (IOException stopped) {
+ return;
+ }
+ }
+ }
+
+ private void exchange(Socket socket) {
+ try (socket) {
+ InputStream input = socket.getInputStream();
+ OutputStream output = socket.getOutputStream();
+ Request request;
+ while (!closed && (request = readRequest(input)) != null) {
+ if (request.operation() == OPERATION_NON_REPLICATED &&
request.commandCode() == PING_CODE) {
+ pings.incrementAndGet();
+ }
+ if (request.operation() == OPERATION_REGISTER) {
+ registers.incrementAndGet();
+ if (withholdRegisterReplies) {
+ continue;
+ }
+ }
+ byte[] body = request.operation() == OPERATION_REGISTER ?
registerBody() : new byte[0];
+ writeResponse(output, request, body);
+ }
+ } catch (IOException clientWentAway) {
+ // A closed client and a closed server look the same here.
+ } finally {
+ closedSockets.incrementAndGet();
+ }
+ }
+
+ @Override
+ public void close() throws IOException {
+ closed = true;
+ for (Socket socket : accepted) {
+ socket.close();
+ }
+ serverSocket.close();
+ }
+ }
+
+ private static Request readRequest(InputStream input) throws IOException {
+ byte[] header = input.readNBytes(HEADER_SIZE);
+ if (header.length == 0) {
+ return null;
+ }
+ if (header.length != HEADER_SIZE) {
+ throw new EOFException("Truncated VSR request header");
+ }
+ ByteBuffer fields =
ByteBuffer.wrap(header).order(ByteOrder.LITTLE_ENDIAN);
+ int size = fields.getInt(SIZE_OFFSET);
+ byte[] body = input.readNBytes(size - HEADER_SIZE);
+ if (body.length != size - HEADER_SIZE) {
+ throw new EOFException("Truncated VSR request body");
+ }
+ return new Request(
+ Byte.toUnsignedInt(header[REQUEST_OPERATION_OFFSET]),
+ fields.getInt(REQUEST_CODE_OFFSET),
+ fields.getLong(REQUEST_ID_OFFSET));
+ }
+
+ private static void writeResponse(OutputStream output, Request request,
byte[] body) throws IOException {
+ byte[] header = new byte[HEADER_SIZE];
+ ByteBuffer fields =
ByteBuffer.wrap(header).order(ByteOrder.LITTLE_ENDIAN);
+ fields.putInt(SIZE_OFFSET, HEADER_SIZE + body.length);
+ header[COMMAND_OFFSET] = COMMAND_REPLY;
+ fields.putLong(REPLY_REQUEST_ID_OFFSET, request.requestId());
+ header[REPLY_OPERATION_OFFSET] = (byte) request.operation();
+ fields.putInt(REPLY_STATUS_OFFSET, 0);
+ output.write(header);
+ output.write(body);
+ output.flush();
+ }
+
+ private static byte[] registerBody() {
+ return ByteBuffer.allocate(21)
+ .order(ByteOrder.LITTLE_ENDIAN)
+ .putInt(0)
+ .putInt(1)
+ .putLong(42)
+ .putInt(11 << 10)
+ .put((byte) 0)
+ .array();
+ }
+
+ private record Request(int operation, int commandCode, long requestId) {}
+}
diff --git
a/foreign/java/java-sdk/src/test/java/org/apache/iggy/client/async/tcp/AsyncTcpConnectionConcurrencyTest.java
b/foreign/java/java-sdk/src/test/java/org/apache/iggy/client/async/tcp/AsyncTcpConnectionConcurrencyTest.java
index 354f09d95..07812d853 100644
---
a/foreign/java/java-sdk/src/test/java/org/apache/iggy/client/async/tcp/AsyncTcpConnectionConcurrencyTest.java
+++
b/foreign/java/java-sdk/src/test/java/org/apache/iggy/client/async/tcp/AsyncTcpConnectionConcurrencyTest.java
@@ -106,6 +106,8 @@ class AsyncTcpConnectionConcurrencyTest {
false,
Optional.empty(),
new AsyncTcpConnection.TcpConnectionPoolConfig(1000, 50),
+ Optional.empty(),
+ 1,
Optional.of(Duration.ofSeconds(1)),
Optional.of(Duration.ofSeconds(2)),
Duration.ofHours(1),
@@ -303,6 +305,8 @@ class AsyncTcpConnectionConcurrencyTest {
false,
Optional.empty(),
new AsyncTcpConnection.TcpConnectionPoolConfig(1000,
acquireTimeoutMillis),
+ Optional.empty(),
+ 1,
Optional.of(Duration.ofSeconds(1)),
Optional.of(Duration.ofSeconds(2)),
heartbeatInterval,