When a header frame with an end_stream flag is received, the close method of the streaming decoder is called (#14313)
This commit is contained in:
parent
e4fa369b29
commit
e49b8b2411
|
|
@ -138,6 +138,11 @@ public abstract class AbstractServerHttpChannelObserver implements CustomizableH
|
||||||
if (httpMetadata == null) {
|
if (httpMetadata == null) {
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
if (!headerSent) {
|
||||||
|
HttpHeaders headers = httpMetadata.headers();
|
||||||
|
headers.set(HttpHeaderNames.STATUS.getName(), resolveStatusCode(throwable));
|
||||||
|
headers.set(HttpHeaderNames.CONTENT_TYPE.getName(), responseEncoder.contentType());
|
||||||
|
}
|
||||||
trailersCustomizer.accept(httpMetadata.headers(), throwable);
|
trailersCustomizer.accept(httpMetadata.headers(), throwable);
|
||||||
getHttpChannel().writeHeader(httpMetadata);
|
getHttpChannel().writeHeader(httpMetadata);
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -92,8 +92,8 @@ public class GrpcHttp2ServerTransportListener extends GenericHttp2ServerTranspor
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
protected void onMetadataCompletion(Http2Header metadata) {
|
protected void onMetadataCompletion(Http2Header metadata) {
|
||||||
super.onMetadataCompletion(metadata);
|
|
||||||
processGrpcHeaders(metadata);
|
processGrpcHeaders(metadata);
|
||||||
|
super.onMetadataCompletion(metadata);
|
||||||
}
|
}
|
||||||
|
|
||||||
private void processGrpcHeaders(Http2Header metadata) {
|
private void processGrpcHeaders(Http2Header metadata) {
|
||||||
|
|
|
||||||
|
|
@ -19,12 +19,10 @@ package org.apache.dubbo.rpc.protocol.tri.h12.http2;
|
||||||
import org.apache.dubbo.common.URL;
|
import org.apache.dubbo.common.URL;
|
||||||
import org.apache.dubbo.common.threadpool.manager.ExecutorRepository;
|
import org.apache.dubbo.common.threadpool.manager.ExecutorRepository;
|
||||||
import org.apache.dubbo.common.threadpool.serial.SerializingExecutor;
|
import org.apache.dubbo.common.threadpool.serial.SerializingExecutor;
|
||||||
import org.apache.dubbo.remoting.http12.HttpMethods;
|
|
||||||
import org.apache.dubbo.remoting.http12.h2.CancelStreamException;
|
import org.apache.dubbo.remoting.http12.h2.CancelStreamException;
|
||||||
import org.apache.dubbo.remoting.http12.h2.H2StreamChannel;
|
import org.apache.dubbo.remoting.http12.h2.H2StreamChannel;
|
||||||
import org.apache.dubbo.remoting.http12.h2.Http2Header;
|
import org.apache.dubbo.remoting.http12.h2.Http2Header;
|
||||||
import org.apache.dubbo.remoting.http12.h2.Http2InputMessage;
|
import org.apache.dubbo.remoting.http12.h2.Http2InputMessage;
|
||||||
import org.apache.dubbo.remoting.http12.h2.Http2InputMessageFrame;
|
|
||||||
import org.apache.dubbo.remoting.http12.h2.Http2ServerChannelObserver;
|
import org.apache.dubbo.remoting.http12.h2.Http2ServerChannelObserver;
|
||||||
import org.apache.dubbo.remoting.http12.h2.Http2TransportListener;
|
import org.apache.dubbo.remoting.http12.h2.Http2TransportListener;
|
||||||
import org.apache.dubbo.remoting.http12.message.DefaultListeningDecoder;
|
import org.apache.dubbo.remoting.http12.message.DefaultListeningDecoder;
|
||||||
|
|
@ -50,15 +48,11 @@ import org.apache.dubbo.rpc.protocol.tri.h12.ServerStreamServerCallListener;
|
||||||
import org.apache.dubbo.rpc.protocol.tri.h12.StreamingHttpMessageListener;
|
import org.apache.dubbo.rpc.protocol.tri.h12.StreamingHttpMessageListener;
|
||||||
import org.apache.dubbo.rpc.protocol.tri.h12.UnaryServerCallListener;
|
import org.apache.dubbo.rpc.protocol.tri.h12.UnaryServerCallListener;
|
||||||
|
|
||||||
import java.io.ByteArrayInputStream;
|
|
||||||
import java.util.concurrent.Executor;
|
import java.util.concurrent.Executor;
|
||||||
|
|
||||||
public class GenericHttp2ServerTransportListener extends AbstractServerTransportListener<Http2Header, Http2InputMessage>
|
public class GenericHttp2ServerTransportListener extends AbstractServerTransportListener<Http2Header, Http2InputMessage>
|
||||||
implements Http2TransportListener {
|
implements Http2TransportListener {
|
||||||
|
|
||||||
private static final Http2InputMessage EMPTY_MESSAGE =
|
|
||||||
new Http2InputMessageFrame(new ByteArrayInputStream(new byte[0]), true);
|
|
||||||
|
|
||||||
private final ExecutorSupport executorSupport;
|
private final ExecutorSupport executorSupport;
|
||||||
private final StreamingDecoder streamingDecoder;
|
private final StreamingDecoder streamingDecoder;
|
||||||
private final FrameworkModel frameworkModel;
|
private final FrameworkModel frameworkModel;
|
||||||
|
|
@ -93,18 +87,6 @@ public class GenericHttp2ServerTransportListener extends AbstractServerTransport
|
||||||
return new SerializingExecutor(executorSupport.getExecutor(metadata));
|
return new SerializingExecutor(executorSupport.getExecutor(metadata));
|
||||||
}
|
}
|
||||||
|
|
||||||
@Override
|
|
||||||
protected void doOnMetadata(Http2Header metadata) {
|
|
||||||
if (metadata.isEndStream()) {
|
|
||||||
if (!HttpMethods.supportBody(metadata.method())) {
|
|
||||||
super.doOnMetadata(metadata);
|
|
||||||
doOnData(EMPTY_MESSAGE);
|
|
||||||
}
|
|
||||||
return;
|
|
||||||
}
|
|
||||||
super.doOnMetadata(metadata);
|
|
||||||
}
|
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
protected HttpMessageListener buildHttpMessageListener() {
|
protected HttpMessageListener buildHttpMessageListener() {
|
||||||
RpcInvocationBuildContext context = getContext();
|
RpcInvocationBuildContext context = getContext();
|
||||||
|
|
@ -180,6 +162,9 @@ public class GenericHttp2ServerTransportListener extends AbstractServerTransport
|
||||||
protected void onMetadataCompletion(Http2Header metadata) {
|
protected void onMetadataCompletion(Http2Header metadata) {
|
||||||
serverChannelObserver.setResponseEncoder(getContext().getHttpMessageEncoder());
|
serverChannelObserver.setResponseEncoder(getContext().getHttpMessageEncoder());
|
||||||
serverChannelObserver.request(1);
|
serverChannelObserver.request(1);
|
||||||
|
if (metadata.isEndStream()) {
|
||||||
|
getStreamingDecoder().close();
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
|
|
|
||||||
Loading…
Reference in New Issue