From 55c25e495b2b1eef77f81cfc37f3645361b5927b Mon Sep 17 00:00:00 2001 From: Xiaoyumuxi <3075514079@qq.com> Date: Mon, 14 Sep 2026 14:30:08 +0800 Subject: [PATCH 01/18] fix: support primitive RPC parameter and return types --- .../com/xiaoyu/rpc/core/util/TypeUtils.java | 55 +++++++++++++++++++ 1 file changed, 55 insertions(+) create mode 100644 rpc-core/src/main/java/com/xiaoyu/rpc/core/util/TypeUtils.java diff --git a/rpc-core/src/main/java/com/xiaoyu/rpc/core/util/TypeUtils.java b/rpc-core/src/main/java/com/xiaoyu/rpc/core/util/TypeUtils.java new file mode 100644 index 0000000..76a3027 --- /dev/null +++ b/rpc-core/src/main/java/com/xiaoyu/rpc/core/util/TypeUtils.java @@ -0,0 +1,55 @@ +package com.xiaoyu.rpc.core.util; + +import java.util.Map; + +/** + * RPC 反射调用相关的类型工具。 + *

+ * Java 基本类型(如 int、boolean)不能直接通过 Class.forName("int") 解析, + * 同时多数序列化器也需要使用对应的包装类型进行反序列化。 + */ +public final class TypeUtils { + + private static final Map> PRIMITIVE_TYPES = Map.ofEntries( + Map.entry("boolean", boolean.class), + Map.entry("byte", byte.class), + Map.entry("short", short.class), + Map.entry("int", int.class), + Map.entry("long", long.class), + Map.entry("float", float.class), + Map.entry("double", double.class), + Map.entry("char", char.class), + Map.entry("void", void.class)); + + private static final Map, Class> WRAPPER_TYPES = Map.ofEntries( + Map.entry(boolean.class, Boolean.class), + Map.entry(byte.class, Byte.class), + Map.entry(short.class, Short.class), + Map.entry(int.class, Integer.class), + Map.entry(long.class, Long.class), + Map.entry(float.class, Float.class), + Map.entry(double.class, Double.class), + Map.entry(char.class, Character.class), + Map.entry(void.class, Void.class)); + + private TypeUtils() { + } + + /** + * 按 JVM/Java 类型名解析 Class,兼容基本类型名称。 + */ + public static Class resolveClass(String typeName) throws ClassNotFoundException { + Class primitiveType = PRIMITIVE_TYPES.get(typeName); + return primitiveType != null ? primitiveType : Class.forName(typeName); + } + + /** + * 将基本类型转换成包装类型,普通引用类型保持不变。 + */ + public static Class wrapPrimitive(Class type) { + if (type == null || !type.isPrimitive()) { + return type; + } + return WRAPPER_TYPES.get(type); + } +} From d20ce82fcab67c78836c56f73389d007fd153b8c Mon Sep 17 00:00:00 2001 From: Xiaoyumuxi <3075514079@qq.com> Date: Mon, 14 Sep 2026 14:30:31 +0800 Subject: [PATCH 02/18] fix: deserialize primitive RPC return types --- .../com/xiaoyu/rpc/core/client/RpcClient.java | 34 ++++++++++++------- 1 file changed, 21 insertions(+), 13 deletions(-) diff --git a/rpc-core/src/main/java/com/xiaoyu/rpc/core/client/RpcClient.java b/rpc-core/src/main/java/com/xiaoyu/rpc/core/client/RpcClient.java index 3dd96ab..acc93fb 100644 --- a/rpc-core/src/main/java/com/xiaoyu/rpc/core/client/RpcClient.java +++ b/rpc-core/src/main/java/com/xiaoyu/rpc/core/client/RpcClient.java @@ -1,14 +1,19 @@ package com.xiaoyu.rpc.core.client; import com.xiaoyu.rpc.common.extension.ExtensionLoader; +import com.xiaoyu.rpc.common.serialization.Serializer; +import com.xiaoyu.rpc.common.serialization.SerializerCode; import com.xiaoyu.rpc.common.vo.RpcRequest; +import com.xiaoyu.rpc.common.vo.RpcResponse; import com.xiaoyu.rpc.core.config.RpcConfig; import com.xiaoyu.rpc.core.registry.ServiceDiscovery; import com.xiaoyu.rpc.core.transport.Transport; import com.xiaoyu.rpc.core.transport.TransportClient; +import com.xiaoyu.rpc.core.util.TypeUtils; import java.net.InetSocketAddress; import java.util.Objects; +import java.util.concurrent.CompletableFuture; public class RpcClient { @@ -31,37 +36,40 @@ public RpcClient() { this.serviceDiscovery = Objects.requireNonNull(serviceDiscovery, "serviceDiscovery"); } - public java.util.concurrent.CompletableFuture sendRequest(RpcRequest request, Class returnType) { + public CompletableFuture sendRequest(RpcRequest request, Class returnType) { try { // 先做一次服务发现(同步查找,通常会命中本地缓存) InetSocketAddress address = serviceDiscovery.lookupService(request.getInterfaceName()); if (address == null) { - java.util.concurrent.CompletableFuture future = new java.util.concurrent.CompletableFuture<>(); + CompletableFuture future = new CompletableFuture<>(); future.completeExceptionally(new RuntimeException("未发现服务: " + request.getInterfaceName())); return future; } // 交给传输层发送,返回异步 Future - java.util.concurrent.CompletableFuture transportFuture = transportClient.sendRequest(request, - address); + CompletableFuture transportFuture = transportClient.sendRequest(request, address); // 在回调里把响应体反序列化成目标返回类型 return transportFuture.thenApply(result -> { - if (result instanceof com.xiaoyu.rpc.common.vo.RpcResponse) { - com.xiaoyu.rpc.common.vo.RpcResponse response = (com.xiaoyu.rpc.common.vo.RpcResponse) result; - - byte[] data = response.getData().toByteArray(); + if (!(result instanceof RpcResponse)) { + throw new RuntimeException("Unexpected response type: " + result.getClass()); + } - com.xiaoyu.rpc.common.serialization.Serializer serializer = com.xiaoyu.rpc.common.serialization.SerializerCode - .getSerializerByCode(RpcConfig.getInstance().getSerializerCode()); - return serializer.deserialize(data, returnType); + RpcResponse response = (RpcResponse) result; + if (returnType == void.class || returnType == Void.class) { + return null; } - throw new RuntimeException("Unexpected response type: " + result.getClass()); + + byte[] data = response.getData().toByteArray(); + Serializer serializer = SerializerCode + .getSerializerByCode(RpcConfig.getInstance().getSerializerCode()); + Class deserializeType = TypeUtils.wrapPrimitive(returnType); + return serializer.deserialize(data, deserializeType); }); } catch (Exception e) { - java.util.concurrent.CompletableFuture future = new java.util.concurrent.CompletableFuture<>(); + CompletableFuture future = new CompletableFuture<>(); future.completeExceptionally(e); return future; } From f2efa0b81277dcb116930ae65da40a4a76a5efa6 Mon Sep 17 00:00:00 2001 From: Xiaoyumuxi <3075514079@qq.com> Date: Mon, 14 Sep 2026 14:30:56 +0800 Subject: [PATCH 03/18] fix: reuse RpcClient in JDK proxies --- .../xiaoyu/rpc/core/client/JdkProxyFactory.java | 16 +++++++++++++--- 1 file changed, 13 insertions(+), 3 deletions(-) diff --git a/rpc-core/src/main/java/com/xiaoyu/rpc/core/client/JdkProxyFactory.java b/rpc-core/src/main/java/com/xiaoyu/rpc/core/client/JdkProxyFactory.java index 1f30ac4..4e6d083 100644 --- a/rpc-core/src/main/java/com/xiaoyu/rpc/core/client/JdkProxyFactory.java +++ b/rpc-core/src/main/java/com/xiaoyu/rpc/core/client/JdkProxyFactory.java @@ -1,18 +1,29 @@ package com.xiaoyu.rpc.core.client; +import com.google.protobuf.ByteString; import com.xiaoyu.rpc.common.serialization.Serializer; import com.xiaoyu.rpc.common.serialization.SerializerCode; import com.xiaoyu.rpc.common.vo.RpcRequest; -import com.google.protobuf.ByteString; import com.xiaoyu.rpc.core.config.RpcConfig; import java.lang.reflect.InvocationHandler; import java.lang.reflect.Method; import java.lang.reflect.Proxy; +import java.util.Objects; import java.util.concurrent.CompletableFuture; public class JdkProxyFactory implements ProxyFactory { + private final RpcClient rpcClient; + + public JdkProxyFactory() { + this(new RpcClient()); + } + + JdkProxyFactory(RpcClient rpcClient) { + this.rpcClient = Objects.requireNonNull(rpcClient, "rpcClient"); + } + @Override @SuppressWarnings("unchecked") public T getProxy(Class clazz) { @@ -34,7 +45,6 @@ public Object invoke(Object proxy, Method method, Object[] args) throws Throwabl } if (args != null) { - // 获取配置的序列化器 Serializer serializer = SerializerCode .getSerializerByCode(RpcConfig.getInstance().getSerializerCode()); for (Object arg : args) { @@ -44,7 +54,7 @@ public Object invoke(Object proxy, Method method, Object[] args) throws Throwabl } RpcRequest request = builder.build(); - CompletableFuture future = new RpcClient().sendRequest(request, method.getReturnType()); + CompletableFuture future = rpcClient.sendRequest(request, method.getReturnType()); // 如果业务接口声明的返回类型是异步的,直接返回 Future;否则阻塞等待结果 if (CompletableFuture.class.isAssignableFrom(method.getReturnType())) { return future; From 1036eecbb6a2342ed53f5b48cf2da5192707eee2 Mon Sep 17 00:00:00 2001 From: Xiaoyumuxi <3075514079@qq.com> Date: Mon, 14 Sep 2026 14:31:41 +0800 Subject: [PATCH 04/18] fix: unify RPC configuration overrides --- .../com/xiaoyu/rpc/core/config/RpcConfig.java | 162 +++++++++++------- 1 file changed, 100 insertions(+), 62 deletions(-) diff --git a/rpc-core/src/main/java/com/xiaoyu/rpc/core/config/RpcConfig.java b/rpc-core/src/main/java/com/xiaoyu/rpc/core/config/RpcConfig.java index cd708cc..4b093ba 100644 --- a/rpc-core/src/main/java/com/xiaoyu/rpc/core/config/RpcConfig.java +++ b/rpc-core/src/main/java/com/xiaoyu/rpc/core/config/RpcConfig.java @@ -22,10 +22,8 @@ public class RpcConfig { // 序列化类型: JAVA, KRYO, PROTOBUF private String serializerType; - // 服务端口 private Integer serverPort; - // 服务端地址 private String serverHost; // 协议名称 @@ -48,6 +46,12 @@ public class RpcConfig { private Integer bossThreads = 1; // 最大连接数 private Integer maxConnections = 100; + // 单次 RPC 请求超时时间 + private Integer requestTimeoutMillis = 5000; + // 服务端业务线程数 (0 = CPU cores) + private Integer businessThreads = 0; + // 服务端业务线程池队列容量 + private Integer businessQueueCapacity = 1000; private RpcConfig() { loadConfig(); @@ -64,7 +68,8 @@ public static synchronized RpcConfig getInstance() { } /** - * 从YAML配置文件加载配置 + * 从YAML配置文件加载配置。 + * 配置优先级:System Properties > Nacos > 本地 rpc-config.yaml > 默认值。 */ private void loadConfig() { Yaml yaml = new Yaml(); @@ -74,28 +79,11 @@ private void loadConfig() { if (inputStream != null) { Map config = yaml.load(inputStream); - // 读取 rpc 配置节点 if (config != null && config.containsKey("rpc")) { @SuppressWarnings("unchecked") Map rpcConfig = (Map) config.get("rpc"); - - this.serializerType = (String) rpcConfig.getOrDefault("serializer", "PROTOBUF"); - this.serverPort = (Integer) rpcConfig.getOrDefault("server-port", 8080); - this.serverHost = (String) rpcConfig.getOrDefault("server-host", "127.0.0.1"); - this.protocol = (String) rpcConfig.getOrDefault("protocol", "netty"); - this.registryAddress = (String) rpcConfig.getOrDefault("registry-address", "127.0.0.1:8848"); - this.registryType = (String) rpcConfig.getOrDefault("registry", "nacos"); - this.proxyType = (String) rpcConfig.getOrDefault("proxy", "jdk"); - this.loadBalancer = (String) rpcConfig.getOrDefault("load-balancer", "roundrobin"); - this.transport = (String) rpcConfig.getOrDefault("transport", "netty"); - this.maxMessageSize = (Integer) rpcConfig.getOrDefault("max-message-size", 8 * 1024 * 1024); - this.workerThreads = (Integer) rpcConfig.getOrDefault("worker-threads", 0); - this.bossThreads = (Integer) rpcConfig.getOrDefault("boss-threads", 1); - this.maxConnections = (Integer) rpcConfig.getOrDefault("max-connections", 100); - - log.info("配置加载成功: 序列化方式={}, 服务器={}:{},使用的协议={}, 注册中心={}, 代理方式={}, 负载均衡={}, 传输层={}, 最大报文={}", - serializerType, serverHost, serverPort, protocol, registryAddress, proxyType, loadBalancer, - transport, maxMessageSize); + updateConfigFields(rpcConfig); + log.info("本地 RPC 配置加载成功"); } else { log.warn("配置文件格式错误,使用默认配置"); setDefaultConfig(); @@ -110,40 +98,16 @@ private void loadConfig() { setDefaultConfig(); } - String portStr = System.getProperty("rpc.server-port"); - if (portStr != null) { - this.serverPort = Integer.parseInt(portStr); - log.info("检测到 System Property 覆盖端口: {}", this.serverPort); - } - - String registryTypeStr = System.getProperty("rpc.registry"); - if (registryTypeStr != null) { - this.registryType = registryTypeStr; - log.info("检测到 System Property 覆盖注册中心类型: {}", this.registryType); - } - - String serializerStr = System.getProperty("rpc.serializer"); - if (serializerStr != null) { - this.serializerType = serializerStr; - log.info("检测到 System Property 覆盖序列化方式: {}", this.serializerType); - } - - String transportStr = System.getProperty("rpc.transport"); - if (transportStr != null) { - this.transport = transportStr; - log.info("检测到 System Property 覆盖传输层: {}", this.transport); - } - - String protocolStr = System.getProperty("rpc.protocol"); - if (protocolStr != null) { - this.protocol = protocolStr; - log.info("检测到 System Property 覆盖协议: {}", this.protocol); - } + // 先应用一次系统属性,使 rpc.registry / rpc.registry-address 能决定是否以及从哪里连接 Nacos。 + applySystemPropertyOverrides(); - // 加载Nacos配置并注册监听 if ("nacos".equalsIgnoreCase(this.registryType)) { loadNacosConfig(); } + + // Nacos 配置加载后再次应用,保证 System Properties / Spring Boot 配置拥有最高优先级。 + applySystemPropertyOverrides(); + logCurrentConfig("配置加载完成"); } /** @@ -151,8 +115,8 @@ private void loadConfig() { */ private void loadNacosConfig() { try { - // Nacos 配置参数 - String serverAddr = this.registryAddress != null && !this.registryAddress.isEmpty() ? this.registryAddress + String serverAddr = this.registryAddress != null && !this.registryAddress.isEmpty() + ? this.registryAddress : "127.0.0.1:8848"; String dataId = "rpc-config.yaml"; String group = "DEFAULT_GROUP"; @@ -162,22 +126,23 @@ private void loadNacosConfig() { ConfigService configService = NacosFactory.createConfigService(properties); - // 首次获取配置 String configInfo = configService.getConfig(dataId, group, 5000); if (configInfo != null && !configInfo.isEmpty()) { - log.info("从Nacos加载配置文件成功!\n{}", configInfo); + log.info("从Nacos加载配置文件成功"); parseYamlConfigString(configInfo); } else { log.info("Nacos中不存在配置 dataId={}, 将使用本地配置", dataId); } - // 添加监听器,实现热切换 configService.addListener(dataId, group, new Listener() { @Override public void receiveConfigInfo(String configInfo) { - log.info("检测到Nacos配置更新!\n{}", configInfo); + log.info("检测到Nacos配置更新"); if (configInfo != null && !configInfo.isEmpty()) { parseYamlConfigString(configInfo); + // 动态配置也不能覆盖显式的 JVM / Spring Boot 配置。 + applySystemPropertyOverrides(); + logCurrentConfig("Nacos 配置更新完成"); } } @@ -240,10 +205,52 @@ private void updateConfigFields(Map rpcConfig) { this.bossThreads = (Integer) rpcConfig.get("boss-threads"); if (rpcConfig.containsKey("max-connections")) this.maxConnections = (Integer) rpcConfig.get("max-connections"); + if (rpcConfig.containsKey("request-timeout-ms")) + this.requestTimeoutMillis = (Integer) rpcConfig.get("request-timeout-ms"); + if (rpcConfig.containsKey("business-threads")) + this.businessThreads = (Integer) rpcConfig.get("business-threads"); + if (rpcConfig.containsKey("business-queue-capacity")) + this.businessQueueCapacity = (Integer) rpcConfig.get("business-queue-capacity"); + } - log.info("配置更新完毕: 序列化方式={}, 服务器={}:{},使用的协议={}, 注册中心={}, 代理方式={}, 负载均衡={}, 传输层={}, 最大报文={}", - serializerType, serverHost, serverPort, protocol, registryAddress, proxyType, loadBalancer, - transport, maxMessageSize); + /** + * Spring Boot Starter 通过 System Properties 将配置同步到核心模块。 + */ + private void applySystemPropertyOverrides() { + this.serverPort = getIntegerOverride("rpc.server-port", this.serverPort); + this.serverHost = getStringOverride("rpc.server-host", this.serverHost); + this.registryType = getStringOverride("rpc.registry", this.registryType); + this.registryAddress = getStringOverride("rpc.registry-address", this.registryAddress); + this.serializerType = getStringOverride("rpc.serializer", this.serializerType); + this.transport = getStringOverride("rpc.transport", this.transport); + this.protocol = getStringOverride("rpc.protocol", this.protocol); + this.proxyType = getStringOverride("rpc.proxy", this.proxyType); + this.loadBalancer = getStringOverride("rpc.load-balancer", this.loadBalancer); + this.maxMessageSize = getIntegerOverride("rpc.max-message-size", this.maxMessageSize); + this.workerThreads = getIntegerOverride("rpc.worker-threads", this.workerThreads); + this.bossThreads = getIntegerOverride("rpc.boss-threads", this.bossThreads); + this.maxConnections = getIntegerOverride("rpc.max-connections", this.maxConnections); + this.requestTimeoutMillis = getIntegerOverride("rpc.request-timeout-ms", this.requestTimeoutMillis); + this.businessThreads = getIntegerOverride("rpc.business-threads", this.businessThreads); + this.businessQueueCapacity = getIntegerOverride("rpc.business-queue-capacity", this.businessQueueCapacity); + } + + private String getStringOverride(String key, String currentValue) { + String value = System.getProperty(key); + return value == null || value.trim().isEmpty() ? currentValue : value.trim(); + } + + private Integer getIntegerOverride(String key, Integer currentValue) { + String value = System.getProperty(key); + if (value == null || value.trim().isEmpty()) { + return currentValue; + } + try { + return Integer.parseInt(value.trim()); + } catch (NumberFormatException e) { + log.warn("忽略非法整数系统属性 {}={}", key, value); + return currentValue; + } } /** @@ -254,8 +261,25 @@ private void setDefaultConfig() { this.serverPort = 8080; this.serverHost = "127.0.0.1"; this.protocol = "netty"; + this.registryAddress = "127.0.0.1:8848"; + this.registryType = "nacos"; + this.proxyType = "jdk"; + this.loadBalancer = "roundrobin"; this.transport = "netty"; this.maxMessageSize = 8 * 1024 * 1024; + this.workerThreads = 0; + this.bossThreads = 1; + this.maxConnections = 100; + this.requestTimeoutMillis = 5000; + this.businessThreads = 0; + this.businessQueueCapacity = 1000; + } + + private void logCurrentConfig(String prefix) { + log.info("{}: serializer={}, server={}:{}, protocol={}, registry={}@{}, proxy={}, loadBalancer={}, transport={}, " + + "maxMessageSize={}, requestTimeoutMs={}, businessThreads={}, businessQueueCapacity={}", + prefix, serializerType, serverHost, serverPort, protocol, registryType, registryAddress, proxyType, + loadBalancer, transport, maxMessageSize, requestTimeoutMillis, businessThreads, businessQueueCapacity); } /** @@ -266,7 +290,6 @@ public byte getSerializerCode() { .getCode(); } - // Getters public String getSerializerType() { return serializerType; } @@ -319,6 +342,18 @@ public Integer getMaxConnections() { return maxConnections; } + public Integer getRequestTimeoutMillis() { + return requestTimeoutMillis; + } + + public Integer getBusinessThreads() { + return businessThreads; + } + + public Integer getBusinessQueueCapacity() { + return businessQueueCapacity; + } + @Override public String toString() { return "RpcConfig{" + @@ -332,6 +367,9 @@ public String toString() { ", loadBalancer='" + loadBalancer + '\'' + ", transport='" + transport + '\'' + ", maxMessageSize=" + maxMessageSize + + ", requestTimeoutMillis=" + requestTimeoutMillis + + ", businessThreads=" + businessThreads + + ", businessQueueCapacity=" + businessQueueCapacity + '}'; } } From acc44a65ca0e261d865039daa31a9129155bf14d Mon Sep 17 00:00:00 2001 From: Xiaoyumuxi <3075514079@qq.com> Date: Mon, 14 Sep 2026 14:31:47 +0800 Subject: [PATCH 05/18] config: add timeout and business executor settings --- rpc-core/src/main/resources/rpc-config.yaml | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/rpc-core/src/main/resources/rpc-config.yaml b/rpc-core/src/main/resources/rpc-config.yaml index 5d977d7..efbdc4e 100644 --- a/rpc-core/src/main/resources/rpc-config.yaml +++ b/rpc-core/src/main/resources/rpc-config.yaml @@ -1,4 +1,5 @@ rpc: + transport: "netty" protocol: "netty" server-host: "127.0.0.1" server-port: 8080 @@ -8,6 +9,9 @@ rpc: proxy: "bytebuddy" load-balancer: "roundrobin" max-message-size: 8388608 # 8MB + request-timeout-ms: 5000 # 单次 RPC 请求超时 worker-threads: 0 # 0 = CPU cores * 2 boss-threads: 1 + business-threads: 0 # 0 = CPU cores + business-queue-capacity: 1000 max-connections: 100 From ed8a231b2e163bb21b5a4bc03d8e1ab0d3a9d9ad Mon Sep 17 00:00:00 2001 From: Xiaoyumuxi <3075514079@qq.com> Date: Mon, 14 Sep 2026 14:32:00 +0800 Subject: [PATCH 06/18] fix: expose core reliability settings in Spring Boot starter --- .../rpc/spring/config/RpcProperties.java | 61 ++++++++++--------- 1 file changed, 31 insertions(+), 30 deletions(-) diff --git a/rpc-spring-boot-starter/src/main/java/com/xiaoyu/rpc/spring/config/RpcProperties.java b/rpc-spring-boot-starter/src/main/java/com/xiaoyu/rpc/spring/config/RpcProperties.java index 1f4f134..49d457a 100644 --- a/rpc-spring-boot-starter/src/main/java/com/xiaoyu/rpc/spring/config/RpcProperties.java +++ b/rpc-spring-boot-starter/src/main/java/com/xiaoyu/rpc/spring/config/RpcProperties.java @@ -10,53 +10,54 @@ @ConfigurationProperties(prefix = "rpc") public class RpcProperties { - /** - * 传输层实现 (netty) - */ + /** 传输层实现 (netty) */ private String transport = "netty"; - /** - * 协议类型 (netty, http, http2, grpc) - */ + /** 协议类型 (netty, http, http2, grpc) */ private String protocol = "netty"; - /** - * 服务端绑定主机 - */ + /** 服务端绑定主机 */ private String serverHost = "127.0.0.1"; - /** - * 服务端绑定端口 - */ + /** 服务端绑定端口 */ private int serverPort = 8080; - /** - * 注册中心类型 (nacos, local) - */ + /** 注册中心类型 (nacos, local) */ private String registry = "nacos"; - /** - * 注册中心地址 - */ + /** 注册中心地址 */ private String registryAddress = "127.0.0.1:8848"; - /** - * 序列化方式 (kryo, protobuf, json, java) - */ + /** 序列化方式 (kryo, protobuf, json, java) */ private String serializer = "kryo"; - /** - * 代理方式 (jdk, bytebuddy) - */ + /** 代理方式 (jdk, bytebuddy) */ private String proxy = "bytebuddy"; - /** - * 负载均衡策略 (roundrobin, random) - */ + /** 负载均衡策略 (roundrobin, random) */ private String loadBalancer = "roundrobin"; - /** - * 是否启用服务端 (Provider 模式) - */ + /** 最大 RPC 报文大小 */ + private int maxMessageSize = 8 * 1024 * 1024; + + /** 单次 RPC 请求超时(毫秒) */ + private int requestTimeoutMs = 5000; + + /** Netty worker 线程数,0 表示使用 CPU cores * 2 */ + private int workerThreads = 0; + + /** Netty boss 线程数 */ + private int bossThreads = 1; + + /** 服务端业务线程数,0 表示使用 CPU cores */ + private int businessThreads = 0; + + /** 服务端业务线程池等待队列容量 */ + private int businessQueueCapacity = 1000; + + /** 客户端最大缓存连接数 */ + private int maxConnections = 100; + + /** 是否启用服务端 (Provider 模式) */ private boolean serverEnabled = true; } From d4fe29af42c211b8f143e567e0593274dc14a856 Mon Sep 17 00:00:00 2001 From: Xiaoyumuxi <3075514079@qq.com> Date: Mon, 14 Sep 2026 14:32:14 +0800 Subject: [PATCH 07/18] fix: propagate Spring Boot RPC settings to core --- .../xiaoyu/rpc/spring/RpcAutoConfiguration.java | 17 ++++++++++++----- 1 file changed, 12 insertions(+), 5 deletions(-) diff --git a/rpc-spring-boot-starter/src/main/java/com/xiaoyu/rpc/spring/RpcAutoConfiguration.java b/rpc-spring-boot-starter/src/main/java/com/xiaoyu/rpc/spring/RpcAutoConfiguration.java index cc7e6f8..d3dd01f 100644 --- a/rpc-spring-boot-starter/src/main/java/com/xiaoyu/rpc/spring/RpcAutoConfiguration.java +++ b/rpc-spring-boot-starter/src/main/java/com/xiaoyu/rpc/spring/RpcAutoConfiguration.java @@ -24,7 +24,7 @@ public class RpcAutoConfiguration { */ @Bean public RpcConfig rpcConfig(RpcProperties properties) { - // 通过 System Properties 传递配置,让 RpcConfig 能够读取 + // 核心模块不依赖 Spring,通过 System Properties 作为两层之间的配置桥接。 System.setProperty("rpc.transport", properties.getTransport()); System.setProperty("rpc.protocol", properties.getProtocol()); System.setProperty("rpc.server-host", properties.getServerHost()); @@ -34,9 +34,17 @@ public RpcConfig rpcConfig(RpcProperties properties) { System.setProperty("rpc.serializer", properties.getSerializer()); System.setProperty("rpc.proxy", properties.getProxy()); System.setProperty("rpc.load-balancer", properties.getLoadBalancer()); + System.setProperty("rpc.max-message-size", String.valueOf(properties.getMaxMessageSize())); + System.setProperty("rpc.request-timeout-ms", String.valueOf(properties.getRequestTimeoutMs())); + System.setProperty("rpc.worker-threads", String.valueOf(properties.getWorkerThreads())); + System.setProperty("rpc.boss-threads", String.valueOf(properties.getBossThreads())); + System.setProperty("rpc.business-threads", String.valueOf(properties.getBusinessThreads())); + System.setProperty("rpc.business-queue-capacity", String.valueOf(properties.getBusinessQueueCapacity())); + System.setProperty("rpc.max-connections", String.valueOf(properties.getMaxConnections())); - log.info("RPC 配置已从 Spring Boot 同步: registry={}, port={}", - properties.getRegistry(), properties.getServerPort()); + log.info("RPC 配置已从 Spring Boot 同步: registry={}, server={}:{}, protocol={}, requestTimeoutMs={}", + properties.getRegistry(), properties.getServerHost(), properties.getServerPort(), + properties.getProtocol(), properties.getRequestTimeoutMs()); return RpcConfig.getInstance(); } @@ -81,9 +89,8 @@ public RpcServerRunner(RpcServer rpcServer) { } @Override - public void run(String... args) throws Exception { + public void run(String... args) { log.info("启动 RPC Server..."); - // 在新线程中启动,避免阻塞 Spring Boot 主线程 Thread serverThread = new Thread(() -> { try { rpcServer.start(); From a2856ac196095dbb411be86297b9e2ef85dee15a Mon Sep 17 00:00:00 2001 From: Xiaoyumuxi <3075514079@qq.com> Date: Mon, 14 Sep 2026 14:32:37 +0800 Subject: [PATCH 08/18] fix: add RPC request timeout and pending cleanup --- .../transport/netty/NettyTransportClient.java | 63 ++++++++++++------- 1 file changed, 42 insertions(+), 21 deletions(-) diff --git a/rpc-transport-netty/src/main/java/com/xiaoyu/rpc/core/transport/netty/NettyTransportClient.java b/rpc-transport-netty/src/main/java/com/xiaoyu/rpc/core/transport/netty/NettyTransportClient.java index 8aedb45..272e008 100644 --- a/rpc-transport-netty/src/main/java/com/xiaoyu/rpc/core/transport/netty/NettyTransportClient.java +++ b/rpc-transport-netty/src/main/java/com/xiaoyu/rpc/core/transport/netty/NettyTransportClient.java @@ -1,7 +1,5 @@ package com.xiaoyu.rpc.core.transport.netty; -import com.xiaoyu.rpc.common.serialization.Serializer; -import com.xiaoyu.rpc.common.serialization.SerializerCode; import com.xiaoyu.rpc.common.vo.RpcRequest; import com.xiaoyu.rpc.common.vo.RpcResponse; import com.xiaoyu.rpc.core.client.ChannelProvider; @@ -20,8 +18,11 @@ import lombok.extern.slf4j.Slf4j; import java.net.InetSocketAddress; +import java.util.UUID; import java.util.concurrent.CompletableFuture; +import java.util.concurrent.ScheduledFuture; import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; @Slf4j public class NettyTransportClient implements TransportClient { @@ -55,38 +56,59 @@ protected void initChannel(SocketChannel ch) { @Override public CompletableFuture sendRequest(RpcRequest request, InetSocketAddress address) { - String protocolName = RpcConfig.getInstance().getProtocol(); + RpcConfig config = RpcConfig.getInstance(); + String protocolName = config.getProtocol(); try { - // 使用 ChannelProvider 获取连接 Channel channel = ChannelProvider.get(address, getBootstrap()); if (channel == null || !channel.isActive()) { throw new RuntimeException("无法连接到服务器: " + address); } - // Reuse handler from pipeline - NettyRpcClientHandler clientHandler = channel.pipeline().get(NettyRpcClientHandler.class); - if (clientHandler == null) { - // Should be added by initChannel, but for safety in some custom protocols: - clientHandler = new NettyRpcClientHandler(); - channel.pipeline().addLast(clientHandler); + NettyRpcClientHandler handler = channel.pipeline().get(NettyRpcClientHandler.class); + if (handler == null) { + handler = new NettyRpcClientHandler(); + channel.pipeline().addLast(handler); } + final NettyRpcClientHandler clientHandler = handler; - // Generate ID and set to request - // requestId 是客户端关联响应的关键键值,必须在发送前写入 - String requestId = java.util.UUID.randomUUID().toString(); - RpcRequest.Builder builder = request.toBuilder(); - builder.setRequestId(requestId); - RpcRequest newRequest = builder.build(); + String requestId = UUID.randomUUID().toString(); + RpcRequest newRequest = request.toBuilder() + .setRequestId(requestId) + .build(); CompletableFuture resultFuture = new CompletableFuture<>(); - // 先注册 future 再发送,避免极端情况下响应先到导致找不到回调 + // 必须先注册 Future 再发送,避免极端情况下响应先到。 clientHandler.addFuture(requestId, resultFuture); + int timeoutMillis = Math.max(1, config.getRequestTimeoutMillis()); + final ScheduledFuture timeoutTask; + try { + timeoutTask = channel.eventLoop().schedule( + () -> clientHandler.failRequest(requestId, + new TimeoutException("RPC请求超时: requestId=" + requestId + + ", timeoutMs=" + timeoutMillis)), + timeoutMillis, + TimeUnit.MILLISECONDS); + } catch (Exception e) { + clientHandler.failRequest(requestId, e); + throw e; + } + + // 无论正常完成、超时还是异常,都取消定时任务并确保 pendingRequests 被清理。 + resultFuture.whenComplete((result, throwable) -> { + timeoutTask.cancel(false); + clientHandler.removeFuture(requestId); + }); + Protocol protocol = ProtocolFactory.getProtocol(protocolName); - protocol.sendRequest(channel, newRequest, clientHandler); + try { + protocol.sendRequest(channel, newRequest, clientHandler); + } catch (Exception e) { + clientHandler.failRequest(requestId, e); + throw e; + } - // 彻底移除 resultFuture.get(),直接返回异步 Future return resultFuture.thenApply(result -> { if (result instanceof RpcResponse) { RpcResponse rpcResponse = (RpcResponse) result; @@ -94,9 +116,8 @@ public CompletableFuture sendRequest(RpcRequest request, InetSocketAddre throw new RuntimeException("服务端报错: " + rpcResponse.getMessage()); } return rpcResponse; - } else { - throw new RuntimeException("服务端返回的不是 RpcResponse 类型"); } + throw new RuntimeException("服务端返回的不是 RpcResponse 类型"); }); } catch (Exception e) { log.error("RPC请求发起失败", e); From 7f3b43c47e7f4124b77ee1ac9d4aaf0fcf3e9d86 Mon Sep 17 00:00:00 2001 From: Xiaoyumuxi <3075514079@qq.com> Date: Mon, 14 Sep 2026 14:32:59 +0800 Subject: [PATCH 09/18] fix: offload RPC business calls from Netty event loop --- .../rpc/core/server/NettyRpcHandler.java | 115 +++++++++++------- 1 file changed, 74 insertions(+), 41 deletions(-) diff --git a/rpc-transport-netty/src/main/java/com/xiaoyu/rpc/core/server/NettyRpcHandler.java b/rpc-transport-netty/src/main/java/com/xiaoyu/rpc/core/server/NettyRpcHandler.java index 6b4e7f2..6b00266 100644 --- a/rpc-transport-netty/src/main/java/com/xiaoyu/rpc/core/server/NettyRpcHandler.java +++ b/rpc-transport-netty/src/main/java/com/xiaoyu/rpc/core/server/NettyRpcHandler.java @@ -1,90 +1,123 @@ package com.xiaoyu.rpc.core.server; -// 务必导入生成的类 -import com.xiaoyu.rpc.common.vo.RpcRequest; -import com.xiaoyu.rpc.common.vo.RpcResponse; +import com.google.protobuf.ByteString; import com.xiaoyu.rpc.common.serialization.Serializer; import com.xiaoyu.rpc.common.serialization.SerializerCode; +import com.xiaoyu.rpc.common.vo.RpcRequest; +import com.xiaoyu.rpc.common.vo.RpcResponse; import com.xiaoyu.rpc.core.config.RpcConfig; - -import com.google.protobuf.ByteString; +import com.xiaoyu.rpc.core.util.TypeUtils; import io.netty.channel.ChannelHandler; import io.netty.channel.ChannelHandlerContext; import io.netty.channel.SimpleChannelInboundHandler; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import java.lang.reflect.InvocationTargetException; import java.lang.reflect.Method; import java.util.List; -import java.util.Map; -import java.util.concurrent.ConcurrentHashMap; +import java.util.Objects; +import java.util.concurrent.Executor; +import java.util.concurrent.ForkJoinPool; +import java.util.concurrent.RejectedExecutionException; @ChannelHandler.Sharable public class NettyRpcHandler extends SimpleChannelInboundHandler { private static final Logger log = LoggerFactory.getLogger(NettyRpcHandler.class); - // 移除内部 Map,改用 ServiceRepository + private final Executor businessExecutor; + + /** + * 兼容直接构造场景。NettyTransportServer 会注入独立的有界业务线程池。 + */ + public NettyRpcHandler() { + this(ForkJoinPool.commonPool()); + } + + public NettyRpcHandler(Executor businessExecutor) { + this.businessExecutor = Objects.requireNonNull(businessExecutor, "businessExecutor"); + } @Override - protected void channelRead0(ChannelHandlerContext ctx, RpcRequest request) throws Exception { - RpcResponse.Builder responseBuilder = RpcResponse.newBuilder(); - responseBuilder.setRequestId(request.getRequestId()); + protected void channelRead0(ChannelHandlerContext ctx, RpcRequest request) { + try { + // 反序列化、反射调用以及用户业务逻辑都可能阻塞,不能占用 Netty EventLoop。 + businessExecutor.execute(() -> processRequest(ctx, request)); + } catch (RejectedExecutionException e) { + log.warn("RPC业务线程池已满,拒绝请求: interface={}, method={}, requestId={}", + request.getInterfaceName(), request.getMethodName(), request.getRequestId()); + writeErrorResponse(ctx, request, "服务器繁忙,请稍后重试"); + } + } + + private void processRequest(ChannelHandlerContext ctx, RpcRequest request) { + RpcResponse.Builder responseBuilder = RpcResponse.newBuilder() + .setRequestId(request.getRequestId()); try { - // 从 ServiceRepository 取到目标服务实现 Object serviceBean = ServiceRepository.getService(request.getInterfaceName()); if (serviceBean == null) { throw new RuntimeException("未找到服务实现: " + request.getInterfaceName()); } - // 将参数类型名还原为 Class[] - // Proto 存的是类名字符串,我们需要反射还原成 Class 对象 List paramTypeNames = request.getParamTypesList(); - Class[] parameterTypes = new Class[paramTypeNames.size()]; - for (int i = 0; i < paramTypeNames.size(); i++) { - // Class.forName 可能抛出 ClassNotFoundException - parameterTypes[i] = Class.forName(paramTypeNames.get(i)); + List paramByteList = request.getParametersList(); + if (paramTypeNames.size() != paramByteList.size()) { + throw new IllegalArgumentException("参数类型数量与参数数量不一致"); } - // 将参数字节反序列化为方法入参 - // Proto 存的是二进制,我们需要反序列化回 Java 对象 - List paramByteList = request.getParametersList(); + Class[] parameterTypes = new Class[paramTypeNames.size()]; Object[] parameters = new Object[paramByteList.size()]; - - // 获取序列化器 Serializer serializer = SerializerCode.getSerializerByCode(RpcConfig.getInstance().getSerializerCode()); - for (int i = 0; i < paramByteList.size(); i++) { + for (int i = 0; i < paramTypeNames.size(); i++) { + Class parameterType = TypeUtils.resolveClass(paramTypeNames.get(i)); + parameterTypes[i] = parameterType; + byte[] bytes = paramByteList.get(i).toByteArray(); - parameters[i] = serializer.deserialize(bytes, parameterTypes[i]); + Class deserializeType = TypeUtils.wrapPrimitive(parameterType); + parameters[i] = serializer.deserialize(bytes, deserializeType); } - // 通过反射调用目标方法 - Class serviceClass = serviceBean.getClass(); - Method method = serviceClass.getMethod(request.getMethodName(), parameterTypes); + Method method = serviceBean.getClass().getMethod(request.getMethodName(), parameterTypes); Object result = method.invoke(serviceBean, parameters); - // 把返回值序列化后写入响应 - byte[] resultBytes; - if (result == null) { - resultBytes = new byte[0]; - } else { - resultBytes = serializer.serialize(result); - } - + byte[] resultBytes = result == null ? new byte[0] : serializer.serialize(result); responseBuilder.setData(ByteString.copyFrom(resultBytes)); responseBuilder.setMessage("Success"); - } catch (Exception e) { + Throwable cause = unwrapInvocationException(e); log.error("Failed to process RPC request: interface={}, method={}, requestId={}", - request.getInterfaceName(), request.getMethodName(), request.getRequestId(), e); - responseBuilder.setMessage("Error: " + e.getMessage()); - // 可以在这里把异常对象也序列化传回去,或者只传错误信息 + request.getInterfaceName(), request.getMethodName(), request.getRequestId(), cause); + responseBuilder.setMessage("Error: " + safeMessage(cause)); responseBuilder.setData(ByteString.EMPTY); } - // 返回响应 ctx.writeAndFlush(responseBuilder.build()); } + + private void writeErrorResponse(ChannelHandlerContext ctx, RpcRequest request, String message) { + RpcResponse response = RpcResponse.newBuilder() + .setRequestId(request.getRequestId()) + .setMessage("Error: " + message) + .setData(ByteString.EMPTY) + .build(); + ctx.writeAndFlush(response); + } + + private Throwable unwrapInvocationException(Exception e) { + if (e instanceof InvocationTargetException) { + Throwable target = ((InvocationTargetException) e).getTargetException(); + if (target != null) { + return target; + } + } + return e; + } + + private String safeMessage(Throwable throwable) { + String message = throwable.getMessage(); + return message == null || message.isEmpty() ? throwable.getClass().getSimpleName() : message; + } } From a2e510d2f99f90b4285477d87f59560dfb57e3d6 Mon Sep 17 00:00:00 2001 From: Xiaoyumuxi <3075514079@qq.com> Date: Mon, 14 Sep 2026 14:33:14 +0800 Subject: [PATCH 10/18] fix: reuse server handler in auto protocol mode --- .../core/protocol/ProtocolDetectHandler.java | 92 +++++++------------ 1 file changed, 33 insertions(+), 59 deletions(-) diff --git a/rpc-transport-netty/src/main/java/com/xiaoyu/rpc/core/protocol/ProtocolDetectHandler.java b/rpc-transport-netty/src/main/java/com/xiaoyu/rpc/core/protocol/ProtocolDetectHandler.java index 84c9683..68de175 100644 --- a/rpc-transport-netty/src/main/java/com/xiaoyu/rpc/core/protocol/ProtocolDetectHandler.java +++ b/rpc-transport-netty/src/main/java/com/xiaoyu/rpc/core/protocol/ProtocolDetectHandler.java @@ -2,28 +2,17 @@ import com.xiaoyu.rpc.core.server.NettyRpcHandler; import io.netty.buffer.ByteBuf; +import io.netty.channel.ChannelHandler; import io.netty.channel.ChannelHandlerContext; import io.netty.handler.codec.ByteToMessageDecoder; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import java.util.List; +import java.util.Objects; /** - * 协议嗅探器 —— 服务端自动识别多种协议 - *

- * 原理:连接建立后,偷看(peek)入站数据的前几个字节,根据特征判断协议类型: - *

    - *
  • 0xAABBCCDD → TCP 私有协议(NettyProtocol)
  • - *
  • 0x50524920 ("PRI ") → HTTP/2 Connection Preface(Http2Protocol 或 - * GrpcProtocol)
  • - *
  • HTTP 方法名(GET / POST / PUT / HEAD / DELETE / OPTIONS / PATCH)→ - * HTTP/1.1(HttpProtocol)
  • - *
- * 识别后动态配置 pipeline 并移除自身,后续按确定的协议处理。 - *

- * 注意:gRPC 底层也是 HTTP/2 传输,无法在字节层面与普通 HTTP/2 区分。 - * 通过构造函数的 {@code http2ProtocolName} 参数控制 HTTP/2 连接的处理方式。 + * 协议嗅探器 —— 服务端自动识别多种协议。 */ public class ProtocolDetectHandler extends ByteToMessageDecoder { @@ -31,51 +20,43 @@ public class ProtocolDetectHandler extends ByteToMessageDecoder { /** TCP 私有协议魔数,与 NettyRpcEncoder/NettyRpcDecoder 一致 */ private static final int NETTY_MAGIC = 0xAABBCCDD; - - /** HTTP/2 Connection Preface 前 4 字节: "PRI " = 0x50524920 */ + /** HTTP/2 Connection Preface 前 4 字节: "PRI " */ private static final int HTTP2_MAGIC = 0x50524920; - // 常见 HTTP/1.1 方法的首字母 ASCII 码 - private static final byte BYTE_G = 'G'; // GET - private static final byte BYTE_P = 'P'; // POST, PUT, PATCH - private static final byte BYTE_D = 'D'; // DELETE - private static final byte BYTE_H = 'H'; // HEAD - private static final byte BYTE_O = 'O'; // OPTIONS - private static final byte BYTE_T = 'T'; // TRACE - private static final byte BYTE_C = 'C'; // CONNECT - - /** - * 当检测到 HTTP/2 Connection Preface 时使用的协议名。 - * 因为 gRPC 底层也是 HTTP/2,无法在字节层面区分, - * 所以通过这个参数指定:可以是 "http2" 或 "grpc"。 - * 默认为 "http2"。 - */ + private static final byte BYTE_G = 'G'; + private static final byte BYTE_P = 'P'; + private static final byte BYTE_D = 'D'; + private static final byte BYTE_H = 'H'; + private static final byte BYTE_O = 'O'; + private static final byte BYTE_T = 'T'; + private static final byte BYTE_C = 'C'; + private final String http2ProtocolName; + private final ChannelHandler serverHandler; - /** - * 默认构造函数,HTTP/2 连接使用 Http2Protocol 处理 - */ public ProtocolDetectHandler() { - this("http2"); + this("http2", new NettyRpcHandler()); } - /** - * 指定 HTTP/2 连接的处理协议 - * - * @param http2ProtocolName HTTP/2 连接使用的协议名("http2" 或 "grpc") - */ public ProtocolDetectHandler(String http2ProtocolName) { - this.http2ProtocolName = http2ProtocolName; + this(http2ProtocolName, new NettyRpcHandler()); + } + + public ProtocolDetectHandler(ChannelHandler serverHandler) { + this("http2", serverHandler); + } + + public ProtocolDetectHandler(String http2ProtocolName, ChannelHandler serverHandler) { + this.http2ProtocolName = Objects.requireNonNull(http2ProtocolName, "http2ProtocolName"); + this.serverHandler = Objects.requireNonNull(serverHandler, "serverHandler"); } @Override - protected void decode(ChannelHandlerContext ctx, ByteBuf in, List out) throws Exception { - // 至少需要 4 字节才能判断协议类型 + protected void decode(ChannelHandlerContext ctx, ByteBuf in, List out) { if (in.readableBytes() < 4) { return; } - // 偷看前 4 字节,不消费(不移动 readerIndex) int magic = in.getInt(in.readerIndex()); byte firstByte = in.getByte(in.readerIndex()); @@ -97,26 +78,19 @@ protected void decode(ChannelHandlerContext ctx, ByteBuf in, List out) t } } - /** - * 判断首字节是否可能是 HTTP/1.1 方法名的开头 - */ private boolean isHttpMethod(byte firstByte) { - return firstByte == BYTE_G // GET - || firstByte == BYTE_P // POST, PUT, PATCH - || firstByte == BYTE_D // DELETE - || firstByte == BYTE_H // HEAD - || firstByte == BYTE_O // OPTIONS - || firstByte == BYTE_T // TRACE - || firstByte == BYTE_C; // CONNECT + return firstByte == BYTE_G + || firstByte == BYTE_P + || firstByte == BYTE_D + || firstByte == BYTE_H + || firstByte == BYTE_O + || firstByte == BYTE_T + || firstByte == BYTE_C; } - /** - * 根据协议名动态配置 pipeline,然后移除自身 - */ private void configProtocol(ChannelHandlerContext ctx, String protocolName) { Protocol protocol = ProtocolFactory.getProtocol(protocolName); - // 先移除自身,再配置协议的编解码器和 Handler ctx.pipeline().remove(this); - protocol.config(ctx.pipeline(), true, new NettyRpcHandler()); + protocol.config(ctx.pipeline(), true, serverHandler); } } From df8ca43298995852ed36a7ae74ec6ff6378049d0 Mon Sep 17 00:00:00 2001 From: Xiaoyumuxi <3075514079@qq.com> Date: Mon, 14 Sep 2026 14:33:44 +0800 Subject: [PATCH 11/18] fix: add bounded business executor to Netty server --- .../transport/netty/NettyTransportServer.java | 88 ++++++++++++++++--- 1 file changed, 74 insertions(+), 14 deletions(-) diff --git a/rpc-transport-netty/src/main/java/com/xiaoyu/rpc/core/transport/netty/NettyTransportServer.java b/rpc-transport-netty/src/main/java/com/xiaoyu/rpc/core/transport/netty/NettyTransportServer.java index 13dd957..d120200 100644 --- a/rpc-transport-netty/src/main/java/com/xiaoyu/rpc/core/transport/netty/NettyTransportServer.java +++ b/rpc-transport-netty/src/main/java/com/xiaoyu/rpc/core/transport/netty/NettyTransportServer.java @@ -14,12 +14,19 @@ import io.netty.channel.socket.nio.NioServerSocketChannel; import lombok.extern.slf4j.Slf4j; +import java.util.concurrent.ArrayBlockingQueue; +import java.util.concurrent.ThreadFactory; +import java.util.concurrent.ThreadPoolExecutor; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; + @Slf4j public class NettyTransportServer implements TransportServer { private final int port; private EventLoopGroup bossGroup; private EventLoopGroup workerGroup; + private ThreadPoolExecutor businessExecutor; public NettyTransportServer(int port) { this.port = port; @@ -27,30 +34,52 @@ public NettyTransportServer(int port) { @Override public void start() throws InterruptedException { - // boss 负责接收连接,worker 负责连接上的读写事件 - bossGroup = new NioEventLoopGroup(); - workerGroup = new NioEventLoopGroup(); + RpcConfig config = RpcConfig.getInstance(); + int cpuCores = Runtime.getRuntime().availableProcessors(); + int bossThreads = Math.max(1, config.getBossThreads()); + int workerThreads = config.getWorkerThreads() != null && config.getWorkerThreads() > 0 + ? config.getWorkerThreads() + : Math.max(1, cpuCores * 2); + int businessThreads = config.getBusinessThreads() != null && config.getBusinessThreads() > 0 + ? config.getBusinessThreads() + : Math.max(1, cpuCores); + int businessQueueCapacity = Math.max(1, config.getBusinessQueueCapacity()); + + bossGroup = new NioEventLoopGroup(bossThreads); + workerGroup = new NioEventLoopGroup(workerThreads); + businessExecutor = new ThreadPoolExecutor( + businessThreads, + businessThreads, + 0L, + TimeUnit.MILLISECONDS, + new ArrayBlockingQueue<>(businessQueueCapacity), + new NamedThreadFactory("rpc-business-"), + new ThreadPoolExecutor.AbortPolicy()); + + // 一个服务端实例共享同一个无状态 Handler 和业务线程池,避免按连接创建线程资源。 + NettyRpcHandler serverHandler = new NettyRpcHandler(businessExecutor); + try { - ServerBootstrap b = new ServerBootstrap(); - b.group(bossGroup, workerGroup) + ServerBootstrap bootstrap = new ServerBootstrap(); + bootstrap.group(bossGroup, workerGroup) .channel(NioServerSocketChannel.class) .childHandler(new ChannelInitializer() { @Override protected void initChannel(SocketChannel ch) { String protocolName = RpcConfig.getInstance().getProtocol(); if ("auto".equalsIgnoreCase(protocolName)) { - // 自动嗅探模式:先放嗅探器,连接建立后根据首字节判断协议 - ch.pipeline().addLast(new ProtocolDetectHandler()); + ch.pipeline().addLast(new ProtocolDetectHandler(serverHandler)); } else { - // 指定协议模式:保持原有行为 Protocol protocol = ProtocolFactory.getProtocol(protocolName); - protocol.config(ch.pipeline(), true, new NettyRpcHandler()); + protocol.config(ch.pipeline(), true, serverHandler); } } }); - log.info("RPC Server (Netty) started on port {}...", port); - b.bind(port).sync().channel().closeFuture().sync(); + log.info("RPC Server (Netty) started on port {}, bossThreads={}, workerThreads={}, businessThreads={}, " + + "businessQueueCapacity={}", + port, bossThreads, workerThreads, businessThreads, businessQueueCapacity); + bootstrap.bind(port).sync().channel().closeFuture().sync(); } finally { stop(); } @@ -58,10 +87,41 @@ protected void initChannel(SocketChannel ch) { @Override public void stop() { - // shutdownGracefully 会等待队列任务处理后再退出,避免直接中断 I/O - if (bossGroup != null) + // 先停止接收新连接,并开始关闭 I/O 线程。 + if (bossGroup != null) { bossGroup.shutdownGracefully(); - if (workerGroup != null) + } + if (workerGroup != null) { workerGroup.shutdownGracefully(); + } + + // 不再接收新任务后,尽量等待已提交的业务请求执行完成。 + if (businessExecutor != null) { + businessExecutor.shutdown(); + try { + if (!businessExecutor.awaitTermination(5, TimeUnit.SECONDS)) { + businessExecutor.shutdownNow(); + } + } catch (InterruptedException e) { + businessExecutor.shutdownNow(); + Thread.currentThread().interrupt(); + } + } + } + + private static final class NamedThreadFactory implements ThreadFactory { + private final String prefix; + private final AtomicInteger sequence = new AtomicInteger(1); + + private NamedThreadFactory(String prefix) { + this.prefix = prefix; + } + + @Override + public Thread newThread(Runnable runnable) { + Thread thread = new Thread(runnable, prefix + sequence.getAndIncrement()); + thread.setDaemon(false); + return thread; + } } } From 8c6f19a7d5b1722f72c7cd58868d61e605b3be44 Mon Sep 17 00:00:00 2001 From: Xiaoyumuxi <3075514079@qq.com> Date: Mon, 14 Sep 2026 14:34:17 +0800 Subject: [PATCH 12/18] test: cover RPC configuration overrides --- .../xiaoyu/rpc/core/config/RpcConfigTest.java | 120 +++++++++++++----- 1 file changed, 87 insertions(+), 33 deletions(-) diff --git a/rpc-core/src/test/java/com/xiaoyu/rpc/core/config/RpcConfigTest.java b/rpc-core/src/test/java/com/xiaoyu/rpc/core/config/RpcConfigTest.java index 7d53bef..a7e7932 100644 --- a/rpc-core/src/test/java/com/xiaoyu/rpc/core/config/RpcConfigTest.java +++ b/rpc-core/src/test/java/com/xiaoyu/rpc/core/config/RpcConfigTest.java @@ -2,8 +2,8 @@ import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.BeforeEach; -import org.junit.jupiter.api.Test; import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Test; import java.lang.reflect.Field; @@ -15,27 +15,36 @@ @DisplayName("RpcConfig 配置测试") public class RpcConfigTest { + private static final String[] RPC_SYSTEM_PROPERTIES = { + "rpc.registry", + "rpc.serializer", + "rpc.server-host", + "rpc.server-port", + "rpc.registry-address", + "rpc.transport", + "rpc.protocol", + "rpc.proxy", + "rpc.load-balancer", + "rpc.max-message-size", + "rpc.request-timeout-ms", + "rpc.worker-threads", + "rpc.boss-threads", + "rpc.business-threads", + "rpc.business-queue-capacity", + "rpc.max-connections" + }; + @BeforeEach void setUp() throws Exception { - // 重置单例以便每个测试独立 resetSingleton(); - // 清理测试用的系统属性 - System.clearProperty("rpc.registry"); - System.clearProperty("rpc.serializer"); - System.clearProperty("rpc.server-port"); - System.clearProperty("rpc.transport"); - System.clearProperty("rpc.protocol"); + clearRpcSystemProperties(); + // 单元测试不依赖外部 Nacos。 + System.setProperty("rpc.registry", "local"); } @AfterEach void tearDown() throws Exception { - // 清理系统属性 - System.clearProperty("rpc.registry"); - System.clearProperty("rpc.serializer"); - System.clearProperty("rpc.server-port"); - System.clearProperty("rpc.transport"); - System.clearProperty("rpc.protocol"); - // 重置单例 + clearRpcSystemProperties(); resetSingleton(); } @@ -45,6 +54,12 @@ private void resetSingleton() throws Exception { instanceField.set(null, null); } + private void clearRpcSystemProperties() { + for (String key : RPC_SYSTEM_PROPERTIES) { + System.clearProperty(key); + } + } + @Test @DisplayName("测试单例模式") void testSingletonPattern() { @@ -59,10 +74,12 @@ void testSingletonPattern() { void testDefaultConfigValues() { RpcConfig config = RpcConfig.getInstance(); - assertNotNull(config.getSerializerType(), "Serializer type should not be null"); - assertNotNull(config.getServerHost(), "Server host should not be null"); - assertNotNull(config.getServerPort(), "Server port should not be null"); - assertNotNull(config.getProtocol(), "Protocol should not be null"); + assertNotNull(config.getSerializerType()); + assertNotNull(config.getServerHost()); + assertNotNull(config.getServerPort()); + assertNotNull(config.getProtocol()); + assertTrue(config.getRequestTimeoutMillis() > 0); + assertTrue(config.getBusinessQueueCapacity() > 0); } @Test @@ -71,8 +88,7 @@ void testSystemPropertyOverrideRegistry() throws Exception { System.setProperty("rpc.registry", "local"); resetSingleton(); - RpcConfig config = RpcConfig.getInstance(); - assertEquals("local", config.getRegistryType(), "Registry type should be overridden by system property"); + assertEquals("local", RpcConfig.getInstance().getRegistryType()); } @Test @@ -81,8 +97,7 @@ void testSystemPropertyOverrideSerializer() throws Exception { System.setProperty("rpc.serializer", "kryo"); resetSingleton(); - RpcConfig config = RpcConfig.getInstance(); - assertEquals("kryo", config.getSerializerType(), "Serializer should be overridden by system property"); + assertEquals("kryo", RpcConfig.getInstance().getSerializerType()); } @Test @@ -91,8 +106,7 @@ void testSystemPropertyOverridePort() throws Exception { System.setProperty("rpc.server-port", "9999"); resetSingleton(); - RpcConfig config = RpcConfig.getInstance(); - assertEquals(9999, config.getServerPort(), "Server port should be overridden by system property"); + assertEquals(9999, RpcConfig.getInstance().getServerPort()); } @Test @@ -101,8 +115,7 @@ void testSystemPropertyOverrideTransport() throws Exception { System.setProperty("rpc.transport", "netty"); resetSingleton(); - RpcConfig config = RpcConfig.getInstance(); - assertEquals("netty", config.getTransport(), "Transport should be overridden by system property"); + assertEquals("netty", RpcConfig.getInstance().getTransport()); } @Test @@ -111,14 +124,54 @@ void testSystemPropertyOverrideProtocol() throws Exception { System.setProperty("rpc.protocol", "grpc"); resetSingleton(); + assertEquals("grpc", RpcConfig.getInstance().getProtocol()); + } + + @Test + @DisplayName("测试 Spring Boot 使用的扩展系统属性全部生效") + void testExtendedSystemPropertyOverrides() throws Exception { + System.setProperty("rpc.server-host", "0.0.0.0"); + System.setProperty("rpc.registry-address", "10.0.0.8:8848"); + System.setProperty("rpc.proxy", "jdk"); + System.setProperty("rpc.load-balancer", "random"); + System.setProperty("rpc.max-message-size", "1048576"); + System.setProperty("rpc.request-timeout-ms", "2500"); + System.setProperty("rpc.worker-threads", "6"); + System.setProperty("rpc.boss-threads", "2"); + System.setProperty("rpc.business-threads", "8"); + System.setProperty("rpc.business-queue-capacity", "256"); + System.setProperty("rpc.max-connections", "64"); + resetSingleton(); + RpcConfig config = RpcConfig.getInstance(); - assertEquals("grpc", config.getProtocol(), "Protocol should be overridden by system property"); + assertEquals("0.0.0.0", config.getServerHost()); + assertEquals("10.0.0.8:8848", config.getRegistryAddress()); + assertEquals("jdk", config.getProxyType()); + assertEquals("random", config.getLoadBalancer()); + assertEquals(1048576, config.getMaxMessageSize()); + assertEquals(2500, config.getRequestTimeoutMillis()); + assertEquals(6, config.getWorkerThreads()); + assertEquals(2, config.getBossThreads()); + assertEquals(8, config.getBusinessThreads()); + assertEquals(256, config.getBusinessQueueCapacity()); + assertEquals(64, config.getMaxConnections()); + } + + @Test + @DisplayName("非法整数系统属性不会破坏配置加载") + void testInvalidIntegerOverrideFallsBack() throws Exception { + System.setProperty("rpc.request-timeout-ms", "not-a-number"); + resetSingleton(); + + RpcConfig config = RpcConfig.getInstance(); + assertEquals(5000, config.getRequestTimeoutMillis()); } @Test @DisplayName("测试 getSerializerCode 方法") - void testGetSerializerCode() { + void testGetSerializerCode() throws Exception { System.setProperty("rpc.serializer", "kryo"); + resetSingleton(); RpcConfig config = RpcConfig.getInstance(); byte code = config.getSerializerCode(); @@ -131,9 +184,10 @@ void testToString() { RpcConfig config = RpcConfig.getInstance(); String str = config.toString(); - assertNotNull(str, "toString should not return null"); - assertTrue(str.contains("RpcConfig"), "toString should contain class name"); - assertTrue(str.contains("serializerType"), "toString should contain serializerType"); - assertTrue(str.contains("serverPort"), "toString should contain serverPort"); + assertNotNull(str); + assertTrue(str.contains("RpcConfig")); + assertTrue(str.contains("serializerType")); + assertTrue(str.contains("serverPort")); + assertTrue(str.contains("requestTimeoutMillis")); } } From 21e2e6354f6b26efd9e30cbdbda7aeeb301326a6 Mon Sep 17 00:00:00 2001 From: Xiaoyumuxi <3075514079@qq.com> Date: Mon, 14 Sep 2026 14:34:39 +0800 Subject: [PATCH 13/18] test: cover primitive RPC return values --- .../xiaoyu/rpc/core/client/RpcClientTest.java | 57 ++++++++++++++----- 1 file changed, 44 insertions(+), 13 deletions(-) diff --git a/rpc-core/src/test/java/com/xiaoyu/rpc/core/client/RpcClientTest.java b/rpc-core/src/test/java/com/xiaoyu/rpc/core/client/RpcClientTest.java index 4e4403b..a38fb56 100644 --- a/rpc-core/src/test/java/com/xiaoyu/rpc/core/client/RpcClientTest.java +++ b/rpc-core/src/test/java/com/xiaoyu/rpc/core/client/RpcClientTest.java @@ -26,12 +26,14 @@ public class RpcClientTest { @BeforeEach void setUp() throws Exception { + System.setProperty("rpc.registry", "local"); System.setProperty("rpc.serializer", "java"); resetRpcConfigSingleton(); } @AfterEach void tearDown() throws Exception { + System.clearProperty("rpc.registry"); System.clearProperty("rpc.serializer"); resetRpcConfigSingleton(); } @@ -45,9 +47,9 @@ void testServiceNotFound() { CompletableFuture future = rpcClient.sendRequest(minimalRequest(), String.class); - assertTrue(future.isCompletedExceptionally(), "Future should be completed exceptionally"); + assertTrue(future.isCompletedExceptionally()); ExecutionException ex = assertThrows(ExecutionException.class, () -> future.get(1, TimeUnit.SECONDS)); - assertTrue(ex.getCause().getMessage().contains("未发现服务"), "Error should mention service not found"); + assertTrue(ex.getCause().getMessage().contains("未发现服务")); } @Test @@ -61,9 +63,9 @@ void testTransportThrows() { CompletableFuture future = rpcClient.sendRequest(minimalRequest(), String.class); - assertTrue(future.isCompletedExceptionally(), "Future should be completed exceptionally"); + assertTrue(future.isCompletedExceptionally()); ExecutionException ex = assertThrows(ExecutionException.class, () -> future.get(1, TimeUnit.SECONDS)); - assertTrue(ex.getCause().getMessage().contains("transport down"), "Error should keep transport failure"); + assertTrue(ex.getCause().getMessage().contains("transport down")); } @Test @@ -75,28 +77,57 @@ void testUnexpectedResponseType() { CompletableFuture future = rpcClient.sendRequest(minimalRequest(), String.class); - assertTrue(future.isCompletedExceptionally(), "Future should be completed exceptionally"); + assertTrue(future.isCompletedExceptionally()); ExecutionException ex = assertThrows(ExecutionException.class, () -> future.get(1, TimeUnit.SECONDS)); - assertTrue(ex.getCause().getMessage().contains("Unexpected response type"), "Error should mention type mismatch"); + assertTrue(ex.getCause().getMessage().contains("Unexpected response type")); } @Test @DisplayName("RpcResponse 正常反序列化返回目标类型") void testSuccessfulDeserialize() throws Exception { Serializer serializer = ExtensionLoader.getExtensionLoader(Serializer.class).getExtension("java"); - byte[] body = serializer.serialize("hello"); - RpcResponse response = RpcResponse.newBuilder() - .setRequestId("req-1") - .setMessage("Success") - .setData(ByteString.copyFrom(body)) - .build(); + RpcResponse response = successfulResponse(serializer.serialize("hello")); TransportClient transportClient = (request, address) -> CompletableFuture.completedFuture(response); ServiceDiscovery serviceDiscovery = serviceName -> new InetSocketAddress("127.0.0.1", 8080); RpcClient rpcClient = new RpcClient(transportClient, serviceDiscovery); Object result = rpcClient.sendRequest(minimalRequest(), String.class).get(1, TimeUnit.SECONDS); - assertEquals("hello", result, "Response payload should be deserialized to String"); + assertEquals("hello", result); + } + + @Test + @DisplayName("基本类型返回值使用包装类型完成反序列化") + void testPrimitiveReturnTypeDeserialize() throws Exception { + Serializer serializer = ExtensionLoader.getExtensionLoader(Serializer.class).getExtension("java"); + RpcResponse response = successfulResponse(serializer.serialize(42)); + + TransportClient transportClient = (request, address) -> CompletableFuture.completedFuture(response); + ServiceDiscovery serviceDiscovery = serviceName -> new InetSocketAddress("127.0.0.1", 8080); + RpcClient rpcClient = new RpcClient(transportClient, serviceDiscovery); + + Object result = rpcClient.sendRequest(minimalRequest(), int.class).get(1, TimeUnit.SECONDS); + assertEquals(42, result); + } + + @Test + @DisplayName("void 返回类型不尝试反序列化空响应体") + void testVoidReturnType() throws Exception { + RpcResponse response = successfulResponse(new byte[0]); + TransportClient transportClient = (request, address) -> CompletableFuture.completedFuture(response); + ServiceDiscovery serviceDiscovery = serviceName -> new InetSocketAddress("127.0.0.1", 8080); + RpcClient rpcClient = new RpcClient(transportClient, serviceDiscovery); + + Object result = rpcClient.sendRequest(minimalRequest(), void.class).get(1, TimeUnit.SECONDS); + assertNull(result); + } + + private static RpcResponse successfulResponse(byte[] body) { + return RpcResponse.newBuilder() + .setRequestId("req-1") + .setMessage("Success") + .setData(ByteString.copyFrom(body)) + .build(); } private static RpcRequest minimalRequest() { From e06974b74075ed4dfcf122b849ccfafc94ea29c7 Mon Sep 17 00:00:00 2001 From: Xiaoyumuxi <3075514079@qq.com> Date: Mon, 14 Sep 2026 14:34:50 +0800 Subject: [PATCH 14/18] test: verify JDK proxy reuses RpcClient --- .../rpc/core/client/JdkProxyFactoryTest.java | 73 +++++++++++++++++++ 1 file changed, 73 insertions(+) create mode 100644 rpc-core/src/test/java/com/xiaoyu/rpc/core/client/JdkProxyFactoryTest.java diff --git a/rpc-core/src/test/java/com/xiaoyu/rpc/core/client/JdkProxyFactoryTest.java b/rpc-core/src/test/java/com/xiaoyu/rpc/core/client/JdkProxyFactoryTest.java new file mode 100644 index 0000000..2e0c5a1 --- /dev/null +++ b/rpc-core/src/test/java/com/xiaoyu/rpc/core/client/JdkProxyFactoryTest.java @@ -0,0 +1,73 @@ +package com.xiaoyu.rpc.core.client; + +import com.google.protobuf.ByteString; +import com.xiaoyu.rpc.common.extension.ExtensionLoader; +import com.xiaoyu.rpc.common.serialization.Serializer; +import com.xiaoyu.rpc.common.vo.RpcResponse; +import com.xiaoyu.rpc.core.config.RpcConfig; +import com.xiaoyu.rpc.core.registry.ServiceDiscovery; +import com.xiaoyu.rpc.core.transport.TransportClient; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Test; + +import java.lang.reflect.Field; +import java.net.InetSocketAddress; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.atomic.AtomicInteger; + +import static org.junit.jupiter.api.Assertions.assertEquals; + +@DisplayName("JdkProxyFactory 客户端复用测试") +class JdkProxyFactoryTest { + + @BeforeEach + void setUp() throws Exception { + System.setProperty("rpc.registry", "local"); + System.setProperty("rpc.serializer", "java"); + resetRpcConfigSingleton(); + } + + @AfterEach + void tearDown() throws Exception { + System.clearProperty("rpc.registry"); + System.clearProperty("rpc.serializer"); + resetRpcConfigSingleton(); + } + + @Test + @DisplayName("同一个代理的多次调用复用注入的 RpcClient") + void testReuseRpcClientAcrossInvocations() throws Exception { + Serializer serializer = ExtensionLoader.getExtensionLoader(Serializer.class).getExtension("java"); + AtomicInteger requestCount = new AtomicInteger(); + + TransportClient transportClient = (request, address) -> { + requestCount.incrementAndGet(); + RpcResponse response = RpcResponse.newBuilder() + .setRequestId(request.getRequestId()) + .setMessage("Success") + .setData(ByteString.copyFrom(serializer.serialize("ok"))) + .build(); + return CompletableFuture.completedFuture(response); + }; + ServiceDiscovery serviceDiscovery = serviceName -> new InetSocketAddress("127.0.0.1", 8080); + RpcClient rpcClient = new RpcClient(transportClient, serviceDiscovery); + + EchoService proxy = new JdkProxyFactory(rpcClient).getProxy(EchoService.class); + + assertEquals("ok", proxy.echo("first")); + assertEquals("ok", proxy.echo("second")); + assertEquals(2, requestCount.get()); + } + + interface EchoService { + String echo(String value); + } + + private static void resetRpcConfigSingleton() throws Exception { + Field field = RpcConfig.class.getDeclaredField("instance"); + field.setAccessible(true); + field.set(null, null); + } +} From ac89f15ee6c44081bc85fd828c23e7bd58cb9226 Mon Sep 17 00:00:00 2001 From: Xiaoyumuxi <3075514079@qq.com> Date: Mon, 14 Sep 2026 14:35:13 +0800 Subject: [PATCH 15/18] test: cover primitive parameters and business executor --- .../rpc/core/server/NettyRpcHandlerTest.java | 133 ++++++++++++++++++ 1 file changed, 133 insertions(+) create mode 100644 rpc-transport-netty/src/test/java/com/xiaoyu/rpc/core/server/NettyRpcHandlerTest.java diff --git a/rpc-transport-netty/src/test/java/com/xiaoyu/rpc/core/server/NettyRpcHandlerTest.java b/rpc-transport-netty/src/test/java/com/xiaoyu/rpc/core/server/NettyRpcHandlerTest.java new file mode 100644 index 0000000..f7307e4 --- /dev/null +++ b/rpc-transport-netty/src/test/java/com/xiaoyu/rpc/core/server/NettyRpcHandlerTest.java @@ -0,0 +1,133 @@ +package com.xiaoyu.rpc.core.server; + +import com.google.protobuf.ByteString; +import com.xiaoyu.rpc.common.extension.ExtensionLoader; +import com.xiaoyu.rpc.common.serialization.Serializer; +import com.xiaoyu.rpc.common.vo.RpcRequest; +import com.xiaoyu.rpc.common.vo.RpcResponse; +import com.xiaoyu.rpc.core.config.RpcConfig; +import io.netty.channel.embedded.EmbeddedChannel; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Test; + +import java.lang.reflect.Field; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicReference; + +import static org.junit.jupiter.api.Assertions.*; + +@DisplayName("NettyRpcHandler 业务执行测试") +class NettyRpcHandlerTest { + + private Serializer serializer; + + @BeforeEach + void setUp() throws Exception { + System.setProperty("rpc.registry", "local"); + System.setProperty("rpc.serializer", "java"); + resetRpcConfigSingleton(); + serializer = ExtensionLoader.getExtensionLoader(Serializer.class).getExtension("java"); + } + + @AfterEach + void tearDown() throws Exception { + System.clearProperty("rpc.registry"); + System.clearProperty("rpc.serializer"); + resetRpcConfigSingleton(); + } + + @Test + @DisplayName("基本类型参数可以正确解析并调用服务") + void testPrimitiveParameterInvocation() { + ServiceRepository.registerService(PrimitiveService.class.getName(), new PrimitiveServiceImpl()); + EmbeddedChannel channel = new EmbeddedChannel(new NettyRpcHandler(Runnable::run)); + + RpcRequest request = RpcRequest.newBuilder() + .setRequestId("primitive-1") + .setInterfaceName(PrimitiveService.class.getName()) + .setMethodName("add") + .addParamTypes("int") + .addParamTypes("int") + .addParameters(ByteString.copyFrom(serializer.serialize(20))) + .addParameters(ByteString.copyFrom(serializer.serialize(22))) + .build(); + + channel.writeInbound(request); + RpcResponse response = channel.readOutbound(); + + assertNotNull(response); + assertEquals("Success", response.getMessage()); + Integer result = serializer.deserialize(response.getData().toByteArray(), Integer.class); + assertEquals(42, result); + channel.finishAndReleaseAll(); + } + + @Test + @DisplayName("业务方法在注入的业务线程中执行") + void testBusinessInvocationRunsOffEventLoop() throws Exception { + CountDownLatch invoked = new CountDownLatch(1); + AtomicReference threadName = new AtomicReference<>(); + ServiceRepository.registerService(ThreadService.class.getName(), new ThreadServiceImpl(invoked, threadName)); + + ExecutorService executor = Executors.newSingleThreadExecutor(r -> new Thread(r, "rpc-business-test")); + EmbeddedChannel channel = new EmbeddedChannel(new NettyRpcHandler(executor)); + try { + RpcRequest request = RpcRequest.newBuilder() + .setRequestId("thread-1") + .setInterfaceName(ThreadService.class.getName()) + .setMethodName("currentThread") + .build(); + + channel.writeInbound(request); + + assertTrue(invoked.await(1, TimeUnit.SECONDS), "Business method should be invoked"); + assertEquals("rpc-business-test", threadName.get()); + } finally { + executor.shutdownNow(); + channel.finishAndReleaseAll(); + } + } + + public interface PrimitiveService { + int add(int left, int right); + } + + public static class PrimitiveServiceImpl implements PrimitiveService { + @Override + public int add(int left, int right) { + return left + right; + } + } + + public interface ThreadService { + String currentThread(); + } + + public static class ThreadServiceImpl implements ThreadService { + private final CountDownLatch invoked; + private final AtomicReference threadName; + + public ThreadServiceImpl(CountDownLatch invoked, AtomicReference threadName) { + this.invoked = invoked; + this.threadName = threadName; + } + + @Override + public String currentThread() { + threadName.set(Thread.currentThread().getName()); + invoked.countDown(); + return threadName.get(); + } + } + + private static void resetRpcConfigSingleton() throws Exception { + Field field = RpcConfig.class.getDeclaredField("instance"); + field.setAccessible(true); + field.set(null, null); + } +} From a4ccacadb728a82d3b1cdbf2303855d0617dbdd4 Mon Sep 17 00:00:00 2001 From: Xiaoyumuxi <3075514079@qq.com> Date: Mon, 14 Sep 2026 14:37:51 +0800 Subject: [PATCH 16/18] ci: avoid host cgroup namespace for Nacos --- .github/workflows/ci.yml | 1 - 1 file changed, 1 deletion(-) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 1d7e392..3d41f7d 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -25,7 +25,6 @@ jobs: --health-timeout=5s --health-retries=15 --health-start-period=60s - --cgroupns=host steps: - name: Checkout code From cb81f7e8353ad55e6d57004cc62879b321a1664c Mon Sep 17 00:00:00 2001 From: Xiaoyumuxi <3075514079@qq.com> Date: Mon, 14 Sep 2026 14:38:46 +0800 Subject: [PATCH 17/18] ci: decouple tests from external Nacos service --- .github/workflows/ci.yml | 18 ------------------ 1 file changed, 18 deletions(-) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 3d41f7d..ca8e7ff 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -10,22 +10,6 @@ jobs: build-and-test: runs-on: ubuntu-24.04 - services: - nacos: - image: nacos/nacos-server:v2.4.3-slim - env: - MODE: standalone - JAVA_OPT: "-Djdk.attach.allowAttachSelf=true" - ports: - - 8848:8848 - - 9848:9848 - options: >- - --health-cmd="curl -f http://localhost:8848/nacos/v1/console/health/readiness" - --health-interval=10s - --health-timeout=5s - --health-retries=15 - --health-start-period=60s - steps: - name: Checkout code uses: actions/checkout@v4 @@ -45,8 +29,6 @@ jobs: - name: Run integration tests run: mvn test -pl rpc-consumer -am -Dtest=FullIntegrationTest - env: - RPC_REGISTRY: local - name: Check test coverage run: mvn jacoco:report From e4756b31af200a0189576915ef49c6bc0de6287f Mon Sep 17 00:00:00 2001 From: Xiaoyumuxi <3075514079@qq.com> Date: Mon, 14 Sep 2026 14:39:49 +0800 Subject: [PATCH 18/18] ci: include reactor dependencies in unit tests --- .github/workflows/ci.yml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index ca8e7ff..7b6f813 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -25,7 +25,7 @@ jobs: run: mvn clean package -DskipTests - name: Run unit tests - run: mvn test -pl rpc-core,rpc-transport-netty + run: mvn test -pl rpc-core,rpc-transport-netty -am - name: Run integration tests run: mvn test -pl rpc-consumer -am -Dtest=FullIntegrationTest