This commit is contained in:
icodening 2024-06-24 18:00:15 +08:00 committed by GitHub
commit e81a2981fc
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194
8 changed files with 99 additions and 22 deletions

View File

@ -166,6 +166,7 @@ public class TripleHttp2Protocol extends AbstractWireProtocol implements ScopeMo
codec.connection().local().flowController().frameWriter(codec.encoder().frameWriter());
List<ChannelHandler> handlers = new ArrayList<>();
handlers.add(new ChannelHandlerPretender(codec));
handlers.add(new ChannelHandlerPretender(new FlushConsolidationHandler(64, true)));
handlers.add(new ChannelHandlerPretender(new Http2MultiplexHandler(new ChannelDuplexHandler())));
handlers.add(new ChannelHandlerPretender(new TriplePingPongHandler(UrlUtils.getCloseTimeout(url))));
handlers.add(new ChannelHandlerPretender(new TripleGoAwayHandler()));

View File

@ -19,7 +19,10 @@ package org.apache.dubbo.rpc.protocol.tri.call;
import org.apache.dubbo.common.logger.ErrorTypeAwareLogger;
import org.apache.dubbo.common.logger.LoggerFactory;
import org.apache.dubbo.common.stream.StreamObserver;
import org.apache.dubbo.common.threadpool.serial.SerializingExecutor;
import org.apache.dubbo.remoting.RemotingException;
import org.apache.dubbo.remoting.api.connection.AbstractConnectionClient;
import org.apache.dubbo.rpc.RpcException;
import org.apache.dubbo.rpc.TriRpcStatus;
import org.apache.dubbo.rpc.model.FrameworkModel;
import org.apache.dubbo.rpc.protocol.tri.RequestMetadata;
@ -32,10 +35,14 @@ import org.apache.dubbo.rpc.protocol.tri.stream.TripleClientStream;
import org.apache.dubbo.rpc.protocol.tri.transport.TripleWriteQueue;
import java.util.Map;
import java.util.Queue;
import java.util.concurrent.ConcurrentLinkedQueue;
import java.util.concurrent.Executor;
import io.netty.channel.Channel;
import io.netty.handler.codec.http2.Http2Exception;
import io.netty.handler.codec.http2.Http2NoMoreStreamIdsException;
import io.netty.util.concurrent.Future;
import static io.netty.handler.codec.http2.Http2Error.FLOW_CONTROL_ERROR;
import static org.apache.dubbo.common.constants.LoggerCodeConstants.PROTOCOL_FAILED_RESPONSE;
@ -43,6 +50,7 @@ import static org.apache.dubbo.common.constants.LoggerCodeConstants.PROTOCOL_FAI
import static org.apache.dubbo.common.constants.LoggerCodeConstants.PROTOCOL_STREAM_LISTENER;
public class TripleClientCall implements ClientCall, ClientStream.Listener {
private static final Object EMPTY = new Object();
private static final ErrorTypeAwareLogger LOGGER = LoggerFactory.getErrorTypeAwareLogger(TripleClientCall.class);
private final AbstractConnectionClient connectionClient;
private final Executor executor;
@ -52,10 +60,12 @@ public class TripleClientCall implements ClientCall, ClientStream.Listener {
private ClientStream stream;
private ClientCall.Listener listener;
private boolean canceled;
private boolean headerSent;
private boolean autoRequest = true;
private boolean done;
private Http2Exception.StreamException streamException;
private volatile boolean sendingHeader = false;
private volatile boolean streamCreated = false;
private final Queue<Runnable> pendingTasks = new ConcurrentLinkedQueue<>();
public TripleClientCall(
AbstractConnectionClient connectionClient,
@ -63,7 +73,7 @@ public class TripleClientCall implements ClientCall, ClientStream.Listener {
FrameworkModel frameworkModel,
TripleWriteQueue writeQueue) {
this.connectionClient = connectionClient;
this.executor = executor;
this.executor = new SerializingExecutor(executor);
this.frameworkModel = frameworkModel;
this.writeQueue = writeQueue;
}
@ -147,7 +157,7 @@ public class TripleClientCall implements ClientCall, ClientStream.Listener {
return;
}
// did not create stream
if (!headerSent) {
if (!streamCreated) {
return;
}
canceled = true;
@ -184,9 +194,70 @@ public class TripleClientCall implements ClientCall, ClientStream.Listener {
} else if (canceled) {
throw new IllegalStateException("Call already canceled");
}
if (!headerSent) {
headerSent = true;
stream.sendHeader(requestMetadata.toHeaders());
doSendMessage(message);
}
private void doSendMessage(Object message) {
pendingTasks.offer(() -> sendData(message));
if (streamCreated) {
executor.execute(this::drainTasks);
}
}
private void sendHeader() {
if (sendingHeader) {
return;
}
if (!streamCreated) {
sendingHeader = true;
Future<?> sendHeaderFuture = stream.sendHeader(requestMetadata.toHeaders());
sendHeaderFuture.addListener(f -> {
sendingHeader = false;
if (!f.isSuccess()) {
Throwable cause = f.cause();
if (cause instanceof Http2NoMoreStreamIdsException) {
// stream id used up and channel will be removed
Channel lastChannel = (Channel) connectionClient.getChannel(true);
if (lastChannel != null && lastChannel.isActive() && lastChannel.isOpen()) {
// already reconnected
start(requestMetadata, listener);
return;
}
// blocking reconnect
executor.execute(() -> {
try {
synchronized (connectionClient) {
Channel channel = (Channel) connectionClient.getChannel(true);
if (channel == null) {
connectionClient.reconnect();
}
}
start(requestMetadata, listener);
} catch (RemotingException e) {
pendingTasks.clear();
throw new RpcException(e);
}
});
}
} else {
streamCreated = true;
executor.execute(this::drainTasks);
}
});
}
}
private void drainTasks() {
Runnable r;
while ((r = pendingTasks.poll()) != null) {
r.run();
}
}
private void sendData(Object message) {
if (EMPTY == message) {
doHalfClose();
return;
}
final byte[] data;
try {
@ -220,7 +291,14 @@ public class TripleClientCall implements ClientCall, ClientStream.Listener {
@Override
public void halfClose() {
if (!headerSent) {
pendingTasks.offer(this::doHalfClose);
if (streamCreated) {
executor.execute(this::drainTasks);
}
}
private void doHalfClose() {
if (!streamCreated) {
return;
}
if (canceled) {
@ -242,8 +320,9 @@ public class TripleClientCall implements ClientCall, ClientStream.Listener {
public StreamObserver<Object> start(RequestMetadata metadata, ClientCall.Listener responseListener) {
this.requestMetadata = metadata;
this.listener = responseListener;
this.stream = new TripleClientStream(
frameworkModel, executor, (Channel) connectionClient.getChannel(true), this, writeQueue);
this.stream =
new TripleClientStream(frameworkModel, executor, (Channel) connectionClient.getChannel(true), this);
this.sendHeader();
return new ClientCallToObserverAdapter<>(this);
}

View File

@ -47,13 +47,13 @@ public class DataQueueCommand extends StreamQueueCommand {
@Override
public void doSend(ChannelHandlerContext ctx, ChannelPromise promise) {
if (data == null) {
ctx.write(new DefaultHttp2DataFrame(endStream), promise);
ctx.writeAndFlush(new DefaultHttp2DataFrame(endStream), promise);
} else {
ByteBuf buf = ctx.alloc().buffer();
buf.writeByte(compressFlag);
buf.writeInt(data.length);
buf.writeBytes(data);
ctx.write(new DefaultHttp2DataFrame(buf, endStream), promise);
ctx.writeAndFlush(new DefaultHttp2DataFrame(buf, endStream), promise);
}
}

View File

@ -34,6 +34,6 @@ public class EndStreamQueueCommand extends StreamQueueCommand {
@Override
public void doSend(ChannelHandlerContext ctx, ChannelPromise promise) {
ctx.write(new DefaultHttp2DataFrame(true), promise);
ctx.writeAndFlush(new DefaultHttp2DataFrame(true), promise);
}
}

View File

@ -55,6 +55,6 @@ public class HeaderQueueCommand extends StreamQueueCommand {
@Override
public void doSend(ChannelHandlerContext ctx, ChannelPromise promise) {
ctx.write(new DefaultHttp2HeadersFrame(headers, endStream), promise);
ctx.writeAndFlush(new DefaultHttp2HeadersFrame(headers, endStream), promise);
}
}

View File

@ -40,7 +40,7 @@ public abstract class QueuedCommand {
public void run(Channel channel) {
if (channel.isActive()) {
channel.write(this).addListener(future -> {
channel.writeAndFlush(this).addListener(future -> {
if (future.isSuccess()) {
promise.setSuccess();
} else {

View File

@ -44,6 +44,6 @@ public class TextDataQueueCommand extends StreamQueueCommand {
@Override
public void doSend(ChannelHandlerContext ctx, ChannelPromise promise) {
ByteBuf buf = ByteBufUtil.writeUtf8(ctx.alloc(), data);
ctx.write(new DefaultHttp2DataFrame(buf, endStream), promise);
ctx.writeAndFlush(new DefaultHttp2DataFrame(buf, endStream), promise);
}
}

View File

@ -59,6 +59,7 @@ import io.netty.channel.ChannelHandlerContext;
import io.netty.channel.ChannelInboundHandlerAdapter;
import io.netty.handler.codec.http2.Http2Error;
import io.netty.handler.codec.http2.Http2Headers;
import io.netty.handler.codec.http2.Http2NoMoreStreamIdsException;
import io.netty.handler.codec.http2.Http2StreamChannel;
import io.netty.handler.codec.http2.Http2StreamChannelBootstrap;
import io.netty.util.ReferenceCountUtil;
@ -99,15 +100,11 @@ public class TripleClientStream extends AbstractStream implements ClientStream {
}
public TripleClientStream(
FrameworkModel frameworkModel,
Executor executor,
Channel parent,
ClientStream.Listener listener,
TripleWriteQueue writeQueue) {
FrameworkModel frameworkModel, Executor executor, Channel parent, ClientStream.Listener listener) {
super(executor, frameworkModel);
this.parent = parent;
this.listener = listener;
this.writeQueue = writeQueue;
this.writeQueue = new TripleWriteQueue();
this.streamChannelFuture = initHttp2StreamChannel(parent);
}
@ -138,7 +135,7 @@ public class TripleClientStream extends AbstractStream implements ClientStream {
}
final HeaderQueueCommand headerCmd = HeaderQueueCommand.createHeaders(streamChannelFuture, headers);
return writeQueue.enqueueFuture(headerCmd, parent.eventLoop()).addListener(future -> {
if (!future.isSuccess()) {
if (!future.isSuccess() && !(future.cause() instanceof Http2NoMoreStreamIdsException)) {
transportException(future.cause());
}
});