Skip to content

Commit ea7dd34

Browse files
authored
fix: harden RPC timeout, threading and configuration (#6)
* fix: support primitive RPC parameter and return types * fix: deserialize primitive RPC return types * fix: reuse RpcClient in JDK proxies * fix: unify RPC configuration overrides * config: add timeout and business executor settings * fix: expose core reliability settings in Spring Boot starter * fix: propagate Spring Boot RPC settings to core * fix: add RPC request timeout and pending cleanup * fix: offload RPC business calls from Netty event loop * fix: reuse server handler in auto protocol mode * fix: add bounded business executor to Netty server * test: cover RPC configuration overrides * test: cover primitive RPC return values * test: verify JDK proxy reuses RpcClient * test: cover primitive parameters and business executor * ci: avoid host cgroup namespace for Nacos * ci: decouple tests from external Nacos service * ci: include reactor dependencies in unit tests
2 parents 060e75a + e4756b3 commit ea7dd34

16 files changed

Lines changed: 797 additions & 314 deletions

File tree

‎.github/workflows/ci.yml‎

Lines changed: 1 addition & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -10,23 +10,6 @@ jobs:
1010
build-and-test:
1111
runs-on: ubuntu-24.04
1212

13-
services:
14-
nacos:
15-
image: nacos/nacos-server:v2.4.3-slim
16-
env:
17-
MODE: standalone
18-
JAVA_OPT: "-Djdk.attach.allowAttachSelf=true"
19-
ports:
20-
- 8848:8848
21-
- 9848:9848
22-
options: >-
23-
--health-cmd="curl -f http://localhost:8848/nacos/v1/console/health/readiness"
24-
--health-interval=10s
25-
--health-timeout=5s
26-
--health-retries=15
27-
--health-start-period=60s
28-
--cgroupns=host
29-
3013
steps:
3114
- name: Checkout code
3215
uses: actions/checkout@v4
@@ -42,12 +25,10 @@ jobs:
4225
run: mvn clean package -DskipTests
4326

4427
- name: Run unit tests
45-
run: mvn test -pl rpc-core,rpc-transport-netty
28+
run: mvn test -pl rpc-core,rpc-transport-netty -am
4629

4730
- name: Run integration tests
4831
run: mvn test -pl rpc-consumer -am -Dtest=FullIntegrationTest
49-
env:
50-
RPC_REGISTRY: local
5132

5233
- name: Check test coverage
5334
run: mvn jacoco:report

‎rpc-core/src/main/java/com/xiaoyu/rpc/core/client/JdkProxyFactory.java‎

Lines changed: 13 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,18 +1,29 @@
11
package com.xiaoyu.rpc.core.client;
22

3+
import com.google.protobuf.ByteString;
34
import com.xiaoyu.rpc.common.serialization.Serializer;
45
import com.xiaoyu.rpc.common.serialization.SerializerCode;
56
import com.xiaoyu.rpc.common.vo.RpcRequest;
6-
import com.google.protobuf.ByteString;
77
import com.xiaoyu.rpc.core.config.RpcConfig;
88

99
import java.lang.reflect.InvocationHandler;
1010
import java.lang.reflect.Method;
1111
import java.lang.reflect.Proxy;
12+
import java.util.Objects;
1213
import java.util.concurrent.CompletableFuture;
1314

1415
public class JdkProxyFactory implements ProxyFactory {
1516

17+
private final RpcClient rpcClient;
18+
19+
public JdkProxyFactory() {
20+
this(new RpcClient());
21+
}
22+
23+
JdkProxyFactory(RpcClient rpcClient) {
24+
this.rpcClient = Objects.requireNonNull(rpcClient, "rpcClient");
25+
}
26+
1627
@Override
1728
@SuppressWarnings("unchecked")
1829
public <T> T getProxy(Class<T> clazz) {
@@ -34,7 +45,6 @@ public Object invoke(Object proxy, Method method, Object[] args) throws Throwabl
3445
}
3546

3647
if (args != null) {
37-
// 获取配置的序列化器
3848
Serializer serializer = SerializerCode
3949
.getSerializerByCode(RpcConfig.getInstance().getSerializerCode());
4050
for (Object arg : args) {
@@ -44,7 +54,7 @@ public Object invoke(Object proxy, Method method, Object[] args) throws Throwabl
4454
}
4555

4656
RpcRequest request = builder.build();
47-
CompletableFuture<Object> future = new RpcClient().sendRequest(request, method.getReturnType());
57+
CompletableFuture<Object> future = rpcClient.sendRequest(request, method.getReturnType());
4858
// 如果业务接口声明的返回类型是异步的,直接返回 Future;否则阻塞等待结果
4959
if (CompletableFuture.class.isAssignableFrom(method.getReturnType())) {
5060
return future;

‎rpc-core/src/main/java/com/xiaoyu/rpc/core/client/RpcClient.java‎

Lines changed: 21 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -1,14 +1,19 @@
11
package com.xiaoyu.rpc.core.client;
22

33
import com.xiaoyu.rpc.common.extension.ExtensionLoader;
4+
import com.xiaoyu.rpc.common.serialization.Serializer;
5+
import com.xiaoyu.rpc.common.serialization.SerializerCode;
46
import com.xiaoyu.rpc.common.vo.RpcRequest;
7+
import com.xiaoyu.rpc.common.vo.RpcResponse;
58
import com.xiaoyu.rpc.core.config.RpcConfig;
69
import com.xiaoyu.rpc.core.registry.ServiceDiscovery;
710
import com.xiaoyu.rpc.core.transport.Transport;
811
import com.xiaoyu.rpc.core.transport.TransportClient;
12+
import com.xiaoyu.rpc.core.util.TypeUtils;
913

1014
import java.net.InetSocketAddress;
1115
import java.util.Objects;
16+
import java.util.concurrent.CompletableFuture;
1217

1318
public class RpcClient {
1419

@@ -31,37 +36,40 @@ public RpcClient() {
3136
this.serviceDiscovery = Objects.requireNonNull(serviceDiscovery, "serviceDiscovery");
3237
}
3338

34-
public java.util.concurrent.CompletableFuture<Object> sendRequest(RpcRequest request, Class<?> returnType) {
39+
public CompletableFuture<Object> sendRequest(RpcRequest request, Class<?> returnType) {
3540
try {
3641
// 先做一次服务发现(同步查找,通常会命中本地缓存)
3742
InetSocketAddress address = serviceDiscovery.lookupService(request.getInterfaceName());
3843

3944
if (address == null) {
40-
java.util.concurrent.CompletableFuture<Object> future = new java.util.concurrent.CompletableFuture<>();
45+
CompletableFuture<Object> future = new CompletableFuture<>();
4146
future.completeExceptionally(new RuntimeException("未发现服务: " + request.getInterfaceName()));
4247
return future;
4348
}
4449

4550
// 交给传输层发送,返回异步 Future
46-
java.util.concurrent.CompletableFuture<Object> transportFuture = transportClient.sendRequest(request,
47-
address);
51+
CompletableFuture<Object> transportFuture = transportClient.sendRequest(request, address);
4852

4953
// 在回调里把响应体反序列化成目标返回类型
5054
return transportFuture.thenApply(result -> {
51-
if (result instanceof com.xiaoyu.rpc.common.vo.RpcResponse) {
52-
com.xiaoyu.rpc.common.vo.RpcResponse response = (com.xiaoyu.rpc.common.vo.RpcResponse) result;
53-
54-
byte[] data = response.getData().toByteArray();
55+
if (!(result instanceof RpcResponse)) {
56+
throw new RuntimeException("Unexpected response type: " + result.getClass());
57+
}
5558

56-
com.xiaoyu.rpc.common.serialization.Serializer serializer = com.xiaoyu.rpc.common.serialization.SerializerCode
57-
.getSerializerByCode(RpcConfig.getInstance().getSerializerCode());
58-
return serializer.deserialize(data, returnType);
59+
RpcResponse response = (RpcResponse) result;
60+
if (returnType == void.class || returnType == Void.class) {
61+
return null;
5962
}
60-
throw new RuntimeException("Unexpected response type: " + result.getClass());
63+
64+
byte[] data = response.getData().toByteArray();
65+
Serializer serializer = SerializerCode
66+
.getSerializerByCode(RpcConfig.getInstance().getSerializerCode());
67+
Class<?> deserializeType = TypeUtils.wrapPrimitive(returnType);
68+
return serializer.deserialize(data, deserializeType);
6169
});
6270

6371
} catch (Exception e) {
64-
java.util.concurrent.CompletableFuture<Object> future = new java.util.concurrent.CompletableFuture<>();
72+
CompletableFuture<Object> future = new CompletableFuture<>();
6573
future.completeExceptionally(e);
6674
return future;
6775
}

0 commit comments

Comments
 (0)