From 24f1932bd2536c5fdf051512b6cc9c5aaf628b3d Mon Sep 17 00:00:00 2001 From: godotg Date: Sun, 10 Sep 2023 15:19:04 +0800 Subject: [PATCH] fix[net]: executor hash now equal in client and server --- .../main/java/com/zfoo/net/task/TaskBus.java | 11 ++++------ .../zfoo/net/core/provider/ProviderTest.java | 6 ++--- .../net/core/tcp/client/TcpClientTest.java | 22 ++++++++++++++++--- 3 files changed, 26 insertions(+), 13 deletions(-) diff --git a/net/src/main/java/com/zfoo/net/task/TaskBus.java b/net/src/main/java/com/zfoo/net/task/TaskBus.java index 7cad0317..6438e404 100644 --- a/net/src/main/java/com/zfoo/net/task/TaskBus.java +++ b/net/src/main/java/com/zfoo/net/task/TaskBus.java @@ -15,9 +15,6 @@ package com.zfoo.net.task; import com.zfoo.event.manager.EventBus; import com.zfoo.net.NetContext; -import com.zfoo.net.router.attachment.GatewayAttachment; -import com.zfoo.net.router.attachment.HttpAttachment; -import com.zfoo.net.router.attachment.SignalAttachment; import com.zfoo.protocol.collection.concurrent.CopyOnWriteHashMapLongObject; import com.zfoo.protocol.util.AssertionUtils; import com.zfoo.protocol.util.RandomUtils; @@ -128,7 +125,7 @@ public final class TaskBus { execute(taskExecutorHash, task); } - public static int calTaskExecutorHash(int taskExecutorHash) { + private static int calTaskExecutorIndex(int taskExecutorHash) { // Other hash algorithms can be customized to make the distribution more uniform return Math.abs(taskExecutorHash) % EXECUTOR_SIZE; } @@ -142,11 +139,11 @@ public final class TaskBus { } else { hash = argument.hashCode(); } - return calTaskExecutorHash(hash); + return hash; } public static void execute(int taskExecutorHash, Runnable runnable) { - executors[calTaskExecutorHash(taskExecutorHash)].execute(ThreadUtils.safeRunnable(runnable)); + executors[calTaskExecutorIndex(taskExecutorHash)].execute(ThreadUtils.safeRunnable(runnable)); } public static void execute(Object argument, Runnable runnable) { @@ -171,7 +168,7 @@ public final class TaskBus { return schedulerExecutor; } - return executors[calTaskExecutorHash(RandomUtils.randomInt())]; + return executors[calTaskExecutorIndex(RandomUtils.randomInt())]; } } diff --git a/net/src/test/java/com/zfoo/net/core/provider/ProviderTest.java b/net/src/test/java/com/zfoo/net/core/provider/ProviderTest.java index 5b27e8f4..61458f01 100644 --- a/net/src/test/java/com/zfoo/net/core/provider/ProviderTest.java +++ b/net/src/test/java/com/zfoo/net/core/provider/ProviderTest.java @@ -75,7 +75,7 @@ public class ProviderTest { var ask = new ProviderMessAsk(); ask.setMessage("Hello, this is the consumer!"); for (int i = 0; i < 1000; i++) { - ThreadUtils.sleep(3000); + ThreadUtils.sleep(1000); var response = NetContext.getConsumer().syncAsk(ask, ProviderMessAnswer.class, null).packet(); logger.info("消费者请求[{}]收到消息[{}]", i, JsonUtils.object2String(response)); } @@ -96,7 +96,7 @@ public class ProviderTest { var atomicInteger = new AtomicInteger(0); for (int i = 0; i < 1000; i++) { - ThreadUtils.sleep(3000); + ThreadUtils.sleep(1000); NetContext.getConsumer().asyncAsk(ask, ProviderMessAnswer.class, null).whenComplete(answer -> { logger.info("消费者请求[{}]收到消息[{}]", atomicInteger.incrementAndGet(), JsonUtils.object2String(answer)); }); @@ -118,7 +118,7 @@ public class ProviderTest { var atomicInteger = new AtomicInteger(0); for (int i = 0; i < 1000; i++) { - ThreadUtils.sleep(3000); + ThreadUtils.sleep(1000); NetContext.getConsumer().asyncAsk(ask, ProviderMessAnswer.class, 100).whenComplete(answer -> { logger.info("消费者请求[{}]收到消息[{}]", atomicInteger.incrementAndGet(), JsonUtils.object2String(answer)); }); diff --git a/net/src/test/java/com/zfoo/net/core/tcp/client/TcpClientTest.java b/net/src/test/java/com/zfoo/net/core/tcp/client/TcpClientTest.java index 6ca6ffdb..55f1f322 100644 --- a/net/src/test/java/com/zfoo/net/core/tcp/client/TcpClientTest.java +++ b/net/src/test/java/com/zfoo/net/core/tcp/client/TcpClientTest.java @@ -16,7 +16,10 @@ package com.zfoo.net.core.tcp.client; import com.zfoo.net.NetContext; import com.zfoo.net.core.HostAndPort; import com.zfoo.net.core.tcp.TcpClient; +import com.zfoo.net.packet.json.JsonHelloResponse; import com.zfoo.net.packet.tcp.TcpHelloRequest; +import com.zfoo.net.packet.tcp.TcpHelloResponse; +import com.zfoo.protocol.util.JsonUtils; import com.zfoo.protocol.util.ThreadUtils; import org.junit.Ignore; import org.junit.Test; @@ -24,6 +27,8 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.context.support.ClassPathXmlApplicationContext; +import java.util.function.Consumer; + /** * @author godotg * @version 3.0 @@ -34,15 +39,26 @@ public class TcpClientTest { private static final Logger logger = LoggerFactory.getLogger(TcpClientTest.class); @Test - public void startClient() { + public void startClient() throws Exception { var context = new ClassPathXmlApplicationContext("config.xml"); var client = new TcpClient(HostAndPort.valueOf("127.0.0.1:9000")); var session = client.start(); + var request = TcpHelloRequest.valueOf("Hello, this is the tcp client!"); for (int i = 0; i < 1000; i++) { - ThreadUtils.sleep(2000); - NetContext.getRouter().send(session, TcpHelloRequest.valueOf("Hello, this is the tcp client!")); + ThreadUtils.sleep(1000); + NetContext.getRouter().send(session, request); + ThreadUtils.sleep(1000); + var response = NetContext.getRouter().syncAsk(session, request, TcpHelloResponse.class, null).packet(); + logger.info("sync client receive [packet:{}] from server", JsonUtils.object2String(response)); + NetContext.getRouter().asyncAsk(session, request, TcpHelloResponse.class, null) + .whenComplete(new Consumer() { + @Override + public void accept(TcpHelloResponse jsonHelloResponse) { + logger.info("async client receive [packet:{}] from server", JsonUtils.object2String(jsonHelloResponse)); + } + }); } ThreadUtils.sleep(Long.MAX_VALUE);