From 445cc826420adc1af9ccacd4a32e7ce4c2b36c72 Mon Sep 17 00:00:00 2001 From: GuoHao Date: Mon, 1 Mar 2021 11:25:58 +0800 Subject: [PATCH] Support triple echo and fix conn bug (#7292) * Add sync connect to avoid connect failed error * Support echo service --- .../java/org/apache/dubbo/remoting/api/Connection.java | 8 ++++++-- .../org/apache/dubbo/rpc/protocol/tri/ServerStream.java | 6 +++++- .../org/apache/dubbo/rpc/protocol/tri/TripleInvoker.java | 2 ++ 3 files changed, 13 insertions(+), 3 deletions(-) diff --git a/dubbo-remoting/dubbo-remoting-api/src/main/java/org/apache/dubbo/remoting/api/Connection.java b/dubbo-remoting/dubbo-remoting-api/src/main/java/org/apache/dubbo/remoting/api/Connection.java index 49ba1ae419..021cdaff99 100644 --- a/dubbo-remoting/dubbo-remoting-api/src/main/java/org/apache/dubbo/remoting/api/Connection.java +++ b/dubbo-remoting/dubbo-remoting-api/src/main/java/org/apache/dubbo/remoting/api/Connection.java @@ -66,6 +66,7 @@ public class Connection extends AbstractReferenceCounted implements ReferenceCou private final InetSocketAddress remote; private final AtomicBoolean closed = new AtomicBoolean(false); private final AtomicReference channel = new AtomicReference<>(); + private final ChannelFuture initConnectFuture; public Connection(URL url) { url = ExecutorUtil.setThreadName(url, "DubboClientHandler"); @@ -76,7 +77,7 @@ public class Connection extends AbstractReferenceCounted implements ReferenceCou this.closeFuture = new DefaultPromise<>(GlobalEventExecutor.INSTANCE); this.remote = getConnectAddress(); final Bootstrap bootstrap = open(); - bootstrap.connect(); + this.initConnectFuture = bootstrap.connect(); } public static Connection getConnectionFromChannel(Channel channel) { @@ -124,7 +125,7 @@ public class Connection extends AbstractReferenceCounted implements ReferenceCou @Override public String toString() { - return "(Ref=" + ReferenceCountUtil.refCnt(this) + ",local="+(getChannel()==null?null:getChannel().localAddress())+",remote="+getRemote(); + return "(Ref=" + ReferenceCountUtil.refCnt(this) + ",local=" + (getChannel() == null ? null : getChannel().localAddress()) + ",remote=" + getRemote(); } public void onGoaway(Channel channel) { @@ -143,6 +144,9 @@ public class Connection extends AbstractReferenceCounted implements ReferenceCou } } + public void connectSync(){ + this.initConnectFuture.awaitUninterruptibly(); + } public boolean isAvailable() { final Channel channel = getChannel(); return channel != null && channel.isActive(); diff --git a/dubbo-rpc/dubbo-rpc-triple/src/main/java/org/apache/dubbo/rpc/protocol/tri/ServerStream.java b/dubbo-rpc/dubbo-rpc-triple/src/main/java/org/apache/dubbo/rpc/protocol/tri/ServerStream.java index b924f4627f..3a18c4b913 100644 --- a/dubbo-rpc/dubbo-rpc-triple/src/main/java/org/apache/dubbo/rpc/protocol/tri/ServerStream.java +++ b/dubbo-rpc/dubbo-rpc-triple/src/main/java/org/apache/dubbo/rpc/protocol/tri/ServerStream.java @@ -34,6 +34,7 @@ import org.apache.dubbo.rpc.model.ProviderModel; import org.apache.dubbo.rpc.model.ServiceDescriptor; import org.apache.dubbo.rpc.model.ServiceRepository; import org.apache.dubbo.rpc.protocol.tri.GrpcStatus.Code; +import org.apache.dubbo.rpc.service.EchoService; import org.apache.dubbo.rpc.service.GenericService; import org.apache.dubbo.triple.TripleWrapper; @@ -216,7 +217,10 @@ public class ServerStream extends AbstractStream implements Stream { if (CommonConstants.$INVOKE.equals(methodName) || CommonConstants.$INVOKE_ASYNC.equals(methodName)) { this.methodDescriptor = repo.lookupMethod(GenericService.class.getName(), methodName); setNeedWrap(true); - } else { + } else if("$echo".equals(methodName)) { + this.methodDescriptor=repo.lookupMethod(EchoService.class.getName(),methodName); + setNeedWrap(true); + }else{ if (methods == null || methods.isEmpty()) { responseErr(ctx, GrpcStatus.fromCode(Code.UNIMPLEMENTED) .withDescription("Method not found:" + methodName + " of service:" + serviceDescriptor.getServiceName())); diff --git a/dubbo-rpc/dubbo-rpc-triple/src/main/java/org/apache/dubbo/rpc/protocol/tri/TripleInvoker.java b/dubbo-rpc/dubbo-rpc-triple/src/main/java/org/apache/dubbo/rpc/protocol/tri/TripleInvoker.java index 037a74db57..2b00edd730 100644 --- a/dubbo-rpc/dubbo-rpc-triple/src/main/java/org/apache/dubbo/rpc/protocol/tri/TripleInvoker.java +++ b/dubbo-rpc/dubbo-rpc-triple/src/main/java/org/apache/dubbo/rpc/protocol/tri/TripleInvoker.java @@ -100,6 +100,8 @@ public class TripleInvoker extends AbstractInvoker { AsyncRpcResult result = new AsyncRpcResult(respFuture, inv); result.setExecutor(executor); + connection.connectSync(); + if (!connection.isAvailable()) { Response response = new Response(req.getId(), req.getVersion()); response.setStatus(Response.CHANNEL_INACTIVE);