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,

Reply via email to