Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
18 commits
Select commit Hold shift + click to select a range
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
21 changes: 1 addition & 20 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -10,23 +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
--cgroupns=host

steps:
- name: Checkout code
uses: actions/checkout@v4
Expand All @@ -42,12 +25,10 @@ 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
env:
RPC_REGISTRY: local

- name: Check test coverage
run: mvn jacoco:report
Expand Down
Original file line number Diff line number Diff line change
@@ -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> T getProxy(Class<T> clazz) {
Expand All @@ -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) {
Expand All @@ -44,7 +54,7 @@ public Object invoke(Object proxy, Method method, Object[] args) throws Throwabl
}

RpcRequest request = builder.build();
CompletableFuture<Object> future = new RpcClient().sendRequest(request, method.getReturnType());
CompletableFuture<Object> future = rpcClient.sendRequest(request, method.getReturnType());
// 如果业务接口声明的返回类型是异步的,直接返回 Future;否则阻塞等待结果
if (CompletableFuture.class.isAssignableFrom(method.getReturnType())) {
return future;
Expand Down
34 changes: 21 additions & 13 deletions rpc-core/src/main/java/com/xiaoyu/rpc/core/client/RpcClient.java
Original file line number Diff line number Diff line change
@@ -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 {

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

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

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

// 交给传输层发送,返回异步 Future
java.util.concurrent.CompletableFuture<Object> transportFuture = transportClient.sendRequest(request,
address);
CompletableFuture<Object> 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<Object> future = new java.util.concurrent.CompletableFuture<>();
CompletableFuture<Object> future = new CompletableFuture<>();
future.completeExceptionally(e);
return future;
}
Expand Down
Loading
Loading