From 7b1199a98eca7da92f9773e048eaa30eded8b7e4 Mon Sep 17 00:00:00 2001 From: Albumen Kevin Date: Fri, 24 Mar 2023 18:15:51 +0800 Subject: [PATCH] Reject if response do not match any request (#11882) --- .../exchange/codec/ExchangeCodec.java | 28 ++++++++++++------- .../remoting/codec/ExchangeCodecTest.java | 21 ++++++++++++++ .../dubbo/rpc/protocol/dubbo/DubboCodec.java | 4 +-- .../protocol/dubbo/DubboCountCodecTest.java | 13 +++++++-- 4 files changed, 51 insertions(+), 15 deletions(-) diff --git a/dubbo-remoting/dubbo-remoting-api/src/main/java/org/apache/dubbo/remoting/exchange/codec/ExchangeCodec.java b/dubbo-remoting/dubbo-remoting-api/src/main/java/org/apache/dubbo/remoting/exchange/codec/ExchangeCodec.java index d631030637..58a2eaeea8 100644 --- a/dubbo-remoting/dubbo-remoting-api/src/main/java/org/apache/dubbo/remoting/exchange/codec/ExchangeCodec.java +++ b/dubbo-remoting/dubbo-remoting-api/src/main/java/org/apache/dubbo/remoting/exchange/codec/ExchangeCodec.java @@ -42,7 +42,10 @@ import org.apache.dubbo.remoting.transport.ExceedPayloadLimitException; import java.io.ByteArrayInputStream; import java.io.IOException; import java.io.InputStream; +import java.text.SimpleDateFormat; +import java.util.Date; +import static org.apache.dubbo.common.constants.LoggerCodeConstants.PROTOCOL_TIMEOUT_SERVER; import static org.apache.dubbo.common.constants.LoggerCodeConstants.TRANSPORT_EXCEED_PAYLOAD_LIMIT; import static org.apache.dubbo.common.constants.LoggerCodeConstants.TRANSPORT_FAILED_RESPONSE; import static org.apache.dubbo.common.constants.LoggerCodeConstants.TRANSPORT_SKIP_UNUSED_STREAM; @@ -171,7 +174,7 @@ public class ExchangeCodec extends TelnetCodec { data = decodeEventData(channel, CodecSupport.deserialize(channel.getUrl(), new ByteArrayInputStream(eventPayload), proto), eventPayload); } } else { - data = decodeResponseData(channel, CodecSupport.deserialize(channel.getUrl(), is, proto), getRequestData(id)); + data = decodeResponseData(channel, CodecSupport.deserialize(channel.getUrl(), is, proto), getRequestData(channel, res, id)); } res.setResult(data); } else { @@ -213,16 +216,21 @@ public class ExchangeCodec extends TelnetCodec { } } - protected Object getRequestData(long id) { + protected Object getRequestData(Channel channel, Response response, long id) { DefaultFuture future = DefaultFuture.getFuture(id); - if (future == null) { - return null; + if (future != null) { + Request req = future.getRequest(); + if (req != null) { + return req.getData(); + } } - Request req = future.getRequest(); - if (req == null) { - return null; - } - return req.getData(); + + logger.warn(PROTOCOL_TIMEOUT_SERVER, "", "", "The timeout response finally returned at " + + (new SimpleDateFormat("yyyy-MM-dd HH:mm:ss.SSS").format(new Date())) + + ", response status is " + response.getStatus() + ", response id is " + response.getId() + + (channel == null ? "" : ", channel: " + channel.getLocalAddress() + + " -> " + channel.getRemoteAddress()) + ", please check provider side for detailed result."); + throw new IllegalArgumentException("Failed to find any request match the response, response id: " + id); } protected void encodeRequest(Channel channel, ChannelBuffer buffer, Request req) throws IOException { @@ -431,7 +439,7 @@ public class ExchangeCodec extends TelnetCodec { try { if (eventBytes != null) { int dataLen = eventBytes.length; - int threshold = ConfigurationUtils.getSystemConfiguration(channel.getUrl().getScopeModel()).getInt("deserialization.event.size", 50); + int threshold = ConfigurationUtils.getSystemConfiguration(channel.getUrl().getScopeModel()).getInt("deserialization.event.size", 15); if (dataLen > threshold) { throw new IllegalArgumentException("Event data too long, actual size " + threshold + ", threshold " + threshold + " rejected for security consideration."); } diff --git a/dubbo-remoting/dubbo-remoting-api/src/test/java/org/apache/dubbo/remoting/codec/ExchangeCodecTest.java b/dubbo-remoting/dubbo-remoting-api/src/test/java/org/apache/dubbo/remoting/codec/ExchangeCodecTest.java index bc40fae681..0951739c82 100644 --- a/dubbo-remoting/dubbo-remoting-api/src/test/java/org/apache/dubbo/remoting/codec/ExchangeCodecTest.java +++ b/dubbo-remoting/dubbo-remoting-api/src/test/java/org/apache/dubbo/remoting/codec/ExchangeCodecTest.java @@ -31,12 +31,14 @@ import org.apache.dubbo.remoting.buffer.ChannelBuffers; import org.apache.dubbo.remoting.exchange.Request; import org.apache.dubbo.remoting.exchange.Response; import org.apache.dubbo.remoting.exchange.codec.ExchangeCodec; +import org.apache.dubbo.remoting.exchange.support.DefaultFuture; import org.apache.dubbo.remoting.telnet.codec.TelnetCodec; import org.apache.dubbo.rpc.model.FrameworkModel; import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; +import org.mockito.Mockito; import java.io.ByteArrayOutputStream; import java.io.IOException; @@ -140,6 +142,8 @@ class ExchangeCodecTest extends TelnetCodecTest { @Test void test_Decode_Error_Length() throws IOException { + DefaultFuture future = DefaultFuture.newFuture(Mockito.mock(Channel.class), new Request(0), 100000, null); + byte[] header = new byte[]{MAGIC_HIGH, MAGIC_LOW, SERIALIZATION_BYTE, 20, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0}; Person person = new Person(); byte[] request = getRequestBytes(person, header); @@ -151,6 +155,8 @@ class ExchangeCodecTest extends TelnetCodecTest { Assertions.assertEquals(person, obj.getResult()); //only decode necessary bytes Assertions.assertEquals(request.length, buffer.readerIndex()); + + future.cancel(); } @Test @@ -229,6 +235,8 @@ class ExchangeCodecTest extends TelnetCodecTest { @Test void test_Decode_Return_Response_Person() throws IOException { + DefaultFuture future = DefaultFuture.newFuture(Mockito.mock(Channel.class), new Request(0), 100000, null); + //00000010-response/oneway/hearbeat=false/hessian |20-stats=ok|id=0|length=0 byte[] header = new byte[]{MAGIC_HIGH, MAGIC_LOW, SERIALIZATION_BYTE, 20, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0}; Person person = new Person(); @@ -238,6 +246,8 @@ class ExchangeCodecTest extends TelnetCodecTest { Assertions.assertEquals(20, obj.getStatus()); Assertions.assertEquals(person, obj.getResult()); System.out.println(obj); + + future.cancel(); } @Test //The status input has a problem, and the read information is wrong when the serialization is serialized. @@ -329,6 +339,8 @@ class ExchangeCodecTest extends TelnetCodecTest { @Test void test_Header_Response_NoSerializationFlag() throws IOException { + DefaultFuture future = DefaultFuture.newFuture(Mockito.mock(Channel.class), new Request(0), 100000, null); + //00000010-response/oneway/hearbeat=false/noset |20-stats=ok|id=0|length=0 byte[] header = new byte[]{MAGIC_HIGH, MAGIC_LOW, SERIALIZATION_BYTE, 20, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0}; Person person = new Person(); @@ -338,10 +350,14 @@ class ExchangeCodecTest extends TelnetCodecTest { Assertions.assertEquals(20, obj.getStatus()); Assertions.assertEquals(person, obj.getResult()); System.out.println(obj); + + future.cancel(); } @Test void test_Header_Response_Heartbeat() throws IOException { + DefaultFuture future = DefaultFuture.newFuture(Mockito.mock(Channel.class), new Request(0), 100000, null); + //00000010-response/oneway/hearbeat=true |20-stats=ok|id=0|length=0 byte[] header = new byte[]{MAGIC_HIGH, MAGIC_LOW, SERIALIZATION_BYTE, 20, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0}; Person person = new Person(); @@ -351,6 +367,8 @@ class ExchangeCodecTest extends TelnetCodecTest { Assertions.assertEquals(20, obj.getStatus()); Assertions.assertEquals(person, obj.getResult()); System.out.println(obj); + + future.cancel(); } @Test @@ -376,6 +394,8 @@ class ExchangeCodecTest extends TelnetCodecTest { @Test void test_Encode_Response() throws IOException { + DefaultFuture future = DefaultFuture.newFuture(Mockito.mock(Channel.class), new Request(1001), 100000, null); + ChannelBuffer encodeBuffer = ChannelBuffers.dynamicBuffer(1024); Channel channel = getClientSideChannel(url); Response response = new Response(); @@ -401,6 +421,7 @@ class ExchangeCodecTest extends TelnetCodecTest { // encode response verson ?? // Assertions.assertEquals(response.getProtocolVersion(), obj.getVersion()); + future.cancel(); } @Test diff --git a/dubbo-rpc/dubbo-rpc-dubbo/src/main/java/org/apache/dubbo/rpc/protocol/dubbo/DubboCodec.java b/dubbo-rpc/dubbo-rpc-dubbo/src/main/java/org/apache/dubbo/rpc/protocol/dubbo/DubboCodec.java index a71444ab16..e744f969ca 100644 --- a/dubbo-rpc/dubbo-rpc-dubbo/src/main/java/org/apache/dubbo/rpc/protocol/dubbo/DubboCodec.java +++ b/dubbo-rpc/dubbo-rpc-dubbo/src/main/java/org/apache/dubbo/rpc/protocol/dubbo/DubboCodec.java @@ -102,12 +102,12 @@ public class DubboCodec extends ExchangeCodec { DecodeableRpcResult result; if (channel.getUrl().getParameter(DECODE_IN_IO_THREAD_KEY, DEFAULT_DECODE_IN_IO_THREAD)) { result = new DecodeableRpcResult(channel, res, is, - (Invocation) getRequestData(id), proto); + (Invocation) getRequestData(channel, res, id), proto); result.decode(); } else { result = new DecodeableRpcResult(channel, res, new UnsafeByteArrayInputStream(readMessageData(is)), - (Invocation) getRequestData(id), proto); + (Invocation) getRequestData(channel, res, id), proto); } data = result; } diff --git a/dubbo-rpc/dubbo-rpc-dubbo/src/test/java/org/apache/dubbo/rpc/protocol/dubbo/DubboCountCodecTest.java b/dubbo-rpc/dubbo-rpc-dubbo/src/test/java/org/apache/dubbo/rpc/protocol/dubbo/DubboCountCodecTest.java index 50593274d5..59e1f4c5a2 100644 --- a/dubbo-rpc/dubbo-rpc-dubbo/src/test/java/org/apache/dubbo/rpc/protocol/dubbo/DubboCountCodecTest.java +++ b/dubbo-rpc/dubbo-rpc-dubbo/src/test/java/org/apache/dubbo/rpc/protocol/dubbo/DubboCountCodecTest.java @@ -22,6 +22,7 @@ import org.apache.dubbo.remoting.buffer.ChannelBuffer; import org.apache.dubbo.remoting.buffer.ChannelBuffers; import org.apache.dubbo.remoting.exchange.Request; import org.apache.dubbo.remoting.exchange.Response; +import org.apache.dubbo.remoting.exchange.support.DefaultFuture; import org.apache.dubbo.remoting.exchange.support.MultiMessage; import org.apache.dubbo.rpc.AppResponse; import org.apache.dubbo.rpc.RpcInvocation; @@ -32,7 +33,9 @@ import org.apache.dubbo.rpc.protocol.dubbo.support.DemoService; import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Test; +import java.util.ArrayList; import java.util.Iterator; +import java.util.List; import static org.apache.dubbo.rpc.Constants.INPUT_KEY; import static org.apache.dubbo.rpc.Constants.OUTPUT_KEY; @@ -45,23 +48,25 @@ class DubboCountCodecTest { ChannelBuffer buffer = ChannelBuffers.buffer(2048); Channel channel = new MockChannel(); Assertions.assertEquals(Codec2.DecodeResult.NEED_MORE_INPUT, dubboCountCodec.decode(channel, buffer)); + List futures = new ArrayList<>(); for (int i = 0; i < 10; i++) { - Request request = new Request(1); + Request request = new Request(i); + futures.add(DefaultFuture.newFuture(channel, request, 1000, null)); RpcInvocation rpcInvocation = new RpcInvocation(null, "echo", DemoService.class.getName(), "", new Class[]{String.class}, new String[]{"yug"}); request.setData(rpcInvocation); dubboCountCodec.encode(channel, buffer, request); } for (int i = 0; i < 10; i++) { - Response response = new Response(1); + Response response = new Response(i); AppResponse appResponse = new AppResponse(i); response.setResult(appResponse); dubboCountCodec.encode(channel, buffer, response); } MultiMessage multiMessage = (MultiMessage) dubboCountCodec.decode(channel, buffer); - Assertions.assertEquals(multiMessage.size(), 20); + Assertions.assertEquals(20, multiMessage.size()); int requestCount = 0; int responseCount = 0; Iterator iterator = multiMessage.iterator(); @@ -79,6 +84,8 @@ class DubboCountCodecTest { } Assertions.assertEquals(requestCount, 10); Assertions.assertEquals(responseCount, 10); + + futures.forEach(DefaultFuture::cancel); } }