From d7d3a5d15e70ef6d20d20b8d5e7d173f3fc41a99 Mon Sep 17 00:00:00 2001 From: Ben Date: Tue, 18 Aug 2026 16:49:01 -0600 Subject: [PATCH 1/2] Move Redis publishing off caller threads --- .../servercomm/redis/RedisHandler.java | 57 ++++++++++++++- .../tests/redis/RedisHandlerTest.java | 72 +++++++++++++++++++ 2 files changed, 127 insertions(+), 2 deletions(-) create mode 100644 SimpleAPI/src/test/java/com/bencodez/simpleapi/tests/redis/RedisHandlerTest.java diff --git a/SimpleAPI/src/main/java/com/bencodez/simpleapi/servercomm/redis/RedisHandler.java b/SimpleAPI/src/main/java/com/bencodez/simpleapi/servercomm/redis/RedisHandler.java index d0dc3e9..9a206cf 100644 --- a/SimpleAPI/src/main/java/com/bencodez/simpleapi/servercomm/redis/RedisHandler.java +++ b/SimpleAPI/src/main/java/com/bencodez/simpleapi/servercomm/redis/RedisHandler.java @@ -2,7 +2,11 @@ import java.util.Map; import java.util.Objects; +import java.util.concurrent.ArrayBlockingQueue; import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.RejectedExecutionException; +import java.util.concurrent.ThreadPoolExecutor; +import java.util.concurrent.TimeUnit; import java.util.function.BiConsumer; import com.bencodez.simpleapi.servercomm.codec.JsonEnvelope; @@ -12,11 +16,17 @@ import redis.clients.jedis.HostAndPort; import redis.clients.jedis.Jedis; import redis.clients.jedis.JedisClientConfig; +import redis.clients.jedis.JedisPool; public abstract class RedisHandler { + private static final int PUBLISH_QUEUE_CAPACITY = 1024; + private static final long PUBLISHER_SHUTDOWN_TIMEOUT_SECONDS = 3L; + private final HostAndPort endpoint; private final JedisClientConfig clientConfig; + private final JedisPool publisherPool; + private final ThreadPoolExecutor publisherExecutor; private final Map listenerThreads = new ConcurrentHashMap<>(); private volatile boolean shuttingDown = false; @@ -42,6 +52,13 @@ public RedisHandler(String host, int port, String username, String password, int } this.clientConfig = cfg.build(); + this.publisherPool = new JedisPool(endpoint, clientConfig); + this.publisherExecutor = new ThreadPoolExecutor(1, 1, 0L, TimeUnit.MILLISECONDS, + new ArrayBlockingQueue<>(PUBLISH_QUEUE_CAPACITY), runnable -> { + Thread thread = new Thread(runnable, "RedisPublishThread-" + endpoint); + thread.setDaemon(true); + return thread; + }, new ThreadPoolExecutor.AbortPolicy()); } public void close() { @@ -58,6 +75,22 @@ public void close() { } } listenerThreads.clear(); + + publisherExecutor.shutdown(); + boolean interrupted = false; + try { + if (!publisherExecutor.awaitTermination(PUBLISHER_SHUTDOWN_TIMEOUT_SECONDS, TimeUnit.SECONDS)) { + publisherExecutor.shutdownNow(); + } + } catch (InterruptedException e) { + publisherExecutor.shutdownNow(); + interrupted = true; + } finally { + publisherPool.close(); + if (interrupted) { + Thread.currentThread().interrupt(); + } + } } public void loadListener(RedisListener listener) { @@ -117,11 +150,31 @@ public void loadListener(RedisListener listener) { thread.start(); } - /** Publish an envelope as a single JSON string. */ + /** + * Queues an envelope for ordered asynchronous publishing. Network connection, + * authentication and publish I/O are never performed on the caller thread. + */ public void publishEnvelope(String channel, JsonEnvelope envelope) { + if (shuttingDown) { + return; + } + String payload = JsonEnvelopeCodec.encode(envelope); + try { + publisherExecutor.execute(() -> publishNow(channel, payload)); + } catch (RejectedExecutionException e) { + if (!shuttingDown) { + debug("Redis publish queue is full; dropping message for channel " + channel); + } + } + } - try (Jedis jedis = new Jedis(endpoint, clientConfig)) { + /** + * Performs one publish using the pooled publisher connection. Kept protected so + * transport scheduling can be regression-tested without a live Redis server. + */ + protected void publishNow(String channel, String payload) { + try (Jedis jedis = publisherPool.getResource()) { debug("Redis Send: " + channel + ", " + payload); jedis.publish(channel, payload); } catch (Exception e) { diff --git a/SimpleAPI/src/test/java/com/bencodez/simpleapi/tests/redis/RedisHandlerTest.java b/SimpleAPI/src/test/java/com/bencodez/simpleapi/tests/redis/RedisHandlerTest.java new file mode 100644 index 0000000..90b32a7 --- /dev/null +++ b/SimpleAPI/src/test/java/com/bencodez/simpleapi/tests/redis/RedisHandlerTest.java @@ -0,0 +1,72 @@ +package com.bencodez.simpleapi.tests.redis; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNotEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import java.util.List; +import java.util.concurrent.CopyOnWriteArrayList; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; + +import org.junit.jupiter.api.Test; + +import com.bencodez.simpleapi.servercomm.codec.JsonEnvelope; +import com.bencodez.simpleapi.servercomm.redis.RedisHandler; + +public class RedisHandlerTest { + + @Test + public void publishRunsOffCallerThreadAndPreservesOrder() throws Exception { + CountDownLatch firstStarted = new CountDownLatch(1); + CountDownLatch releaseFirst = new CountDownLatch(1); + CountDownLatch completed = new CountDownLatch(2); + List channels = new CopyOnWriteArrayList<>(); + List threadNames = new CopyOnWriteArrayList<>(); + + RedisHandler handler = new RedisHandler("127.0.0.1", 6379, "", "", 0) { + @Override + public void debug(String message) { + // no-op + } + + @Override + protected void publishNow(String channel, String payload) { + threadNames.add(Thread.currentThread().getName()); + channels.add(channel); + if (channels.size() == 1) { + firstStarted.countDown(); + try { + releaseFirst.await(2, TimeUnit.SECONDS); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + } + completed.countDown(); + } + }; + + try { + String callerThread = Thread.currentThread().getName(); + JsonEnvelope envelope = JsonEnvelope.builder("Presence").put("server", "survival").build(); + + handler.publishEnvelope("first", envelope); + assertTrue(firstStarted.await(1, TimeUnit.SECONDS)); + + long startNanos = System.nanoTime(); + handler.publishEnvelope("second", envelope); + long elapsedMillis = TimeUnit.NANOSECONDS.toMillis(System.nanoTime() - startNanos); + + assertTrue(elapsedMillis < 250, "publishEnvelope should only enqueue work"); + releaseFirst.countDown(); + assertTrue(completed.await(2, TimeUnit.SECONDS)); + assertEquals(List.of("first", "second"), channels); + assertEquals(2, threadNames.size()); + assertNotEquals(callerThread, threadNames.get(0)); + assertEquals(threadNames.get(0), threadNames.get(1)); + } finally { + releaseFirst.countDown(); + handler.close(); + } + } +} From 4d00e3734d66e37ef788e5428cc6438095112c18 Mon Sep 17 00:00:00 2001 From: Ben Date: Tue, 18 Aug 2026 16:55:27 -0600 Subject: [PATCH 2/2] Validate pooled Redis connections before publish --- .../bencodez/simpleapi/servercomm/redis/RedisHandler.java | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/SimpleAPI/src/main/java/com/bencodez/simpleapi/servercomm/redis/RedisHandler.java b/SimpleAPI/src/main/java/com/bencodez/simpleapi/servercomm/redis/RedisHandler.java index 9a206cf..447f0af 100644 --- a/SimpleAPI/src/main/java/com/bencodez/simpleapi/servercomm/redis/RedisHandler.java +++ b/SimpleAPI/src/main/java/com/bencodez/simpleapi/servercomm/redis/RedisHandler.java @@ -17,6 +17,7 @@ import redis.clients.jedis.Jedis; import redis.clients.jedis.JedisClientConfig; import redis.clients.jedis.JedisPool; +import redis.clients.jedis.JedisPoolConfig; public abstract class RedisHandler { @@ -52,7 +53,9 @@ public RedisHandler(String host, int port, String username, String password, int } this.clientConfig = cfg.build(); - this.publisherPool = new JedisPool(endpoint, clientConfig); + JedisPoolConfig publisherPoolConfig = new JedisPoolConfig(); + publisherPoolConfig.setTestOnBorrow(true); + this.publisherPool = new JedisPool(publisherPoolConfig, endpoint, clientConfig); this.publisherExecutor = new ThreadPoolExecutor(1, 1, 0L, TimeUnit.MILLISECONDS, new ArrayBlockingQueue<>(PUBLISH_QUEUE_CAPACITY), runnable -> { Thread thread = new Thread(runnable, "RedisPublishThread-" + endpoint);