BitoAgent commented on code in PR #13786:
URL: https://github.com/apache/dubbo/pull/13786#discussion_r1573858702
##########
dubbo-remoting/dubbo-remoting-http12/src/main/java/org/apache/dubbo/remoting/http12/message/codec/JsonCodec.java:
##########
@@ -40,6 +40,7 @@ public void encode(OutputStream os, Object data, Charset
charset) throws EncodeE
}
}
+ @Override
Review Comment:
**Issue**: The @Override annotation is added to the encode method which
enhances readability and ensures proper overriding. However, the method lacks
input validation which might lead to EncodeException if null data is passed.
<br> **Fix**: Add null checks for the data parameter to avoid potential
NullPointerExceptions or EncodeExceptions during runtime. <br> **Code
Suggestion**:
```
if (data == null) {
throw new IllegalArgumentException("Data cannot be null.");
}
// Continue with encoding process
```
##########
dubbo-rpc/dubbo-rpc-triple/src/main/java/org/apache/dubbo/rpc/protocol/tri/h12/grpc/GrpcRequestHandlerMapping.java:
##########
@@ -42,9 +43,16 @@ protected boolean supportContentType(String contentType) {
@Override
protected void determineHttpMessageCodec(RequestHandler handler, URL url,
HttpRequest request) {
- HttpMessageCodec codec = CODEC_FACTORY.createCodec(url,
getFrameworkModel(), request.contentType());
- handler.setHttpMessageDecoder(codec);
- handler.setHttpMessageEncoder(codec);
+ GrpcCompositeCodec grpcCompositeCodec =
+ (GrpcCompositeCodec) CODEC_FACTORY.createCodec(url,
getFrameworkModel(), request.contentType());
+ MethodDescriptor methodDescriptor =
DescriptorUtils.findMethodDescriptor(
+ handler.getServiceDescriptor(), handler.getMethodName(),
handler.isHasStub());
+ if (methodDescriptor != null) {
+ handler.setMethodDescriptor(methodDescriptor);
+ grpcCompositeCodec.loadPackableMethod(methodDescriptor);
+ }
+ handler.setHttpMessageDecoder(grpcCompositeCodec);
+ handler.setHttpMessageEncoder(grpcCompositeCodec);
Review Comment:
**Security Issue**: The method determineHttpMessageCodec dynamically loads
a codec based on the request content type without validating the content type,
which could lead to loading of unexpected or malicious codecs. <br> **Fix**:
Implement validation for the content type against a whitelist of allowed
content types before loading the codec. <br> **Code Suggestion**:
```
protected void determineHttpMessageCodec(RequestHandler handler, URL url,
HttpRequest request) {
if (!isValidContentType(request.contentType())) {
throw new IllegalArgumentException("Invalid content type");
}
GrpcCompositeCodec grpcCompositeCodec =
(GrpcCompositeCodec) CODEC_FACTORY.createCodec(url,
getFrameworkModel(), request.contentType());
MethodDescriptor methodDescriptor =
DescriptorUtils.findMethodDescriptor(
handler.getServiceDescriptor(), handler.getMethodName(),
handler.isHasStub());
if (methodDescriptor != null) {
handler.setMethodDescriptor(methodDescriptor);
grpcCompositeCodec.loadPackableMethod(methodDescriptor);
}
handler.setHttpMessageDecoder(grpcCompositeCodec);
handler.setHttpMessageEncoder(grpcCompositeCodec);
}
private boolean isValidContentType(String contentType) {
// Implementation of content type validation against a whitelist
return true; // Simplified for this example
}
```
##########
dubbo-remoting/dubbo-remoting-http12/src/main/java/org/apache/dubbo/remoting/http12/message/LengthFieldStreamingDecoder.java:
##########
@@ -130,16 +127,12 @@ private void deliver() {
}
private void processHeader() throws IOException {
- ByteArrayOutputStream bos = new
ByteArrayOutputStream(lengthFieldOffset + lengthFieldLength);
byte[] offsetData = new byte[lengthFieldOffset];
int ignore = accumulate.read(offsetData);
- bos.write(offsetData);
processOffset(new ByteArrayInputStream(offsetData), lengthFieldOffset);
byte[] lengthBytes = new byte[lengthFieldLength];
ignore = accumulate.read(lengthBytes);
- bos.write(lengthBytes);
requiredLength = bytesToInt(lengthBytes);
- this.dataHeader = new ByteArrayInputStream(bos.toByteArray());
// Continue reading the frame body.
state = DecodeState.PAYLOAD;
Review Comment:
**Security Issue**: The method 'processHeader' lacks proper validation of
the header length, potentially leading to buffer overflow or denial of service
if an attacker sends crafted input with a negative size or size that's
excessively large. <br> **Fix**: Implement checks to ensure that the header
length (both 'lengthFieldOffset' and 'lengthFieldLength') is within expected
bounds before processing it. <br> **Code Suggestion**:
```
private void processHeader() throws IOException {
byte[] offsetData = new byte[lengthFieldOffset];
int ignore = accumulate.read(offsetData);
if (lengthFieldOffset < 0 || lengthFieldLength < 0 || lengthFieldOffset
+ lengthFieldLength > accumulate.available()) {
throw new IllegalArgumentException("Invalid header size");
}
processOffset(new ByteArrayInputStream(offsetData), lengthFieldOffset);
byte[] lengthBytes = new byte[lengthFieldLength];
ignore = accumulate.read(lengthBytes);
requiredLength = bytesToInt(lengthBytes);
}
```
##########
dubbo-rpc/dubbo-rpc-triple/src/main/java/org/apache/dubbo/rpc/protocol/tri/h12/grpc/GrpcHttp2ServerTransportListener.java:
##########
@@ -145,39 +146,31 @@ public void onMessage(InputStream inputStream) {
private class DetermineMethodDescriptorListener implements
StreamingDecoder.FragmentListener {
- @Override
- public void onFragmentMessage(InputStream rawMessage) {}
-
@Override
public void onClose() {
getStreamingDecoder().close();
}
@Override
- public void onFragmentMessage(InputStream dataHeader, InputStream
rawMessage) {
+ public void onFragmentMessage(InputStream rawMessage) {
try {
- ByteArrayOutputStream merged =
- new ByteArrayOutputStream(dataHeader.available() +
rawMessage.available());
- StreamUtils.copy(dataHeader, merged);
- byte[] data = StreamUtils.readBytes(rawMessage);
-
RpcInvocationBuildContext context = getContext();
if (null == context.getMethodDescriptor()) {
-
context.setMethodDescriptor(DescriptorUtils.findTripleMethodDescriptor(
- context.getServiceDescriptor(),
context.getMethodName(), data));
+ byte[] data = StreamUtils.readBytes(rawMessage);
+ MethodDescriptor methodDescriptor =
DescriptorUtils.findTripleMethodDescriptor(
+ context.getServiceDescriptor(),
context.getMethodName(), data);
+ context.setMethodDescriptor(methodDescriptor);
setHttpMessageListener(GrpcHttp2ServerTransportListener.super.buildHttpMessageListener());
// replace decoder
GrpcCompositeCodec grpcCompositeCodec =
(GrpcCompositeCodec) context.getHttpMessageDecoder();
- MethodMetadata methodMetadata =
context.getMethodMetadata();
-
grpcCompositeCodec.setDecodeTypes(methodMetadata.getActualRequestTypes());
- grpcCompositeCodec.setEncodeTypes(new Class[]
{methodMetadata.getActualResponseType()});
+ grpcCompositeCodec.loadPackableMethod(methodDescriptor);
getServerChannelObserver().setResponseEncoder(grpcCompositeCodec);
+ rawMessage = new ByteArrayInputStream(data);
}
- merged.write(data);
- getHttpMessageListener().onMessage(new
ByteArrayInputStream(merged.toByteArray()));
+ getStreamingDecoder().invokeListener(rawMessage);
} catch (IOException e) {
throw new DecodeException(e);
}
Review Comment:
**Security Issue**: The implementation of onFragmentMessage does not
validate the input before processing, which could lead to security
vulnerabilities such as injection attacks or buffer overflows. <br> **Fix**:
Implement input validation for the rawMessage InputStream before processing.
Ensure that the data conforms to expected formats and sizes to prevent
injection attacks or other forms of malicious input. <br> **Code Suggestion**:
```
Implement input validation for the rawMessage InputStream before
processing. Ensure that the data conforms to expected formats and sizes to
prevent injection attacks or other forms of malicious input.
```
##########
dubbo-rpc/dubbo-rpc-triple/src/main/java/org/apache/dubbo/rpc/protocol/tri/h12/grpc/GrpcCompositeCodecFactory.java:
##########
@@ -22,19 +22,14 @@
import org.apache.dubbo.remoting.http12.message.HttpMessageDecoderFactory;
import org.apache.dubbo.remoting.http12.message.HttpMessageEncoderFactory;
import org.apache.dubbo.remoting.http12.message.MediaType;
-import org.apache.dubbo.remoting.utils.UrlUtils;
import org.apache.dubbo.rpc.model.FrameworkModel;
@Activate
public class GrpcCompositeCodecFactory implements HttpMessageEncoderFactory,
HttpMessageDecoderFactory {
@Override
public HttpMessageCodec createCodec(URL url, FrameworkModel
frameworkModel, String mediaType) {
- String serializeName = UrlUtils.serializationOrDefault(url);
- WrapperHttpMessageCodec wrapperHttpMessageCodec = new
WrapperHttpMessageCodec(url, frameworkModel);
- wrapperHttpMessageCodec.setSerializeType(serializeName);
- ProtobufHttpMessageCodec protobufHttpMessageCodec = new
ProtobufHttpMessageCodec();
- return new GrpcCompositeCodec(protobufHttpMessageCodec,
wrapperHttpMessageCodec);
+ return new GrpcCompositeCodec(url, frameworkModel, mediaType);
Review Comment:
**Security Issue**: Direct instantiation of GrpcCompositeCodec with URL,
frameworkModel, and mediaType without validation or sanitization may introduce
security risks, such as information leakage or remote code execution if the
inputs are controlled by an attacker. <br> **Fix**: Validate and sanitize the
URL, frameworkModel, and mediaType parameters before using them to instantiate
GrpcCompositeCodec. <br> **Code Suggestion**:
```
+ Validate and sanitize the URL, frameworkModel, and mediaType parameters
before using them to instantiate GrpcCompositeCodec.
```
##########
dubbo-rpc/dubbo-rpc-triple/src/main/java/org/apache/dubbo/rpc/protocol/tri/h12/grpc/GrpcCompositeCodec.java:
##########
@@ -16,123 +16,103 @@
*/
package org.apache.dubbo.rpc.protocol.tri.h12.grpc;
+import org.apache.dubbo.common.URL;
+import org.apache.dubbo.common.config.ConfigurationUtils;
+import org.apache.dubbo.common.io.StreamUtils;
+import org.apache.dubbo.common.utils.ArrayUtils;
import org.apache.dubbo.remoting.http12.exception.DecodeException;
import org.apache.dubbo.remoting.http12.exception.EncodeException;
import org.apache.dubbo.remoting.http12.message.HttpMessageCodec;
import org.apache.dubbo.remoting.http12.message.MediaType;
+import org.apache.dubbo.rpc.model.FrameworkModel;
+import org.apache.dubbo.rpc.model.MethodDescriptor;
+import org.apache.dubbo.rpc.model.PackableMethod;
+import org.apache.dubbo.rpc.model.PackableMethodFactory;
import java.io.IOException;
import java.io.InputStream;
import java.io.OutputStream;
import java.nio.charset.Charset;
+import java.util.Map;
+import java.util.concurrent.ConcurrentHashMap;
-import com.google.protobuf.Message;
-
-import static
org.apache.dubbo.common.constants.CommonConstants.PROTOBUF_MESSAGE_CLASS_NAME;
+import static org.apache.dubbo.common.constants.CommonConstants.DEFAULT_KEY;
+import static
org.apache.dubbo.common.constants.CommonConstants.DUBBO_PACKABLE_METHOD_FACTORY;
public class GrpcCompositeCodec implements HttpMessageCodec {
- private final ProtobufHttpMessageCodec protobufHttpMessageCodec;
+ private static final String PACKABLE_METHOD_CACHE =
"PACKABLE_METHOD_CACHE";
- private final WrapperHttpMessageCodec wrapperHttpMessageCodec;
+ private final URL url;
- public GrpcCompositeCodec(
- ProtobufHttpMessageCodec protobufHttpMessageCodec,
WrapperHttpMessageCodec wrapperHttpMessageCodec) {
- this.protobufHttpMessageCodec = protobufHttpMessageCodec;
- this.wrapperHttpMessageCodec = wrapperHttpMessageCodec;
- }
+ private final FrameworkModel frameworkModel;
+
+ private final String mediaType;
- public void setEncodeTypes(Class<?>[] encodeTypes) {
- this.wrapperHttpMessageCodec.setEncodeTypes(encodeTypes);
+ private PackableMethod packableMethod;
+
+ public GrpcCompositeCodec(URL url, FrameworkModel frameworkModel, String
mediaType) {
+ this.url = url;
+ this.frameworkModel = frameworkModel;
+ this.mediaType = mediaType;
}
- public void setDecodeTypes(Class<?>[] decodeTypes) {
- this.wrapperHttpMessageCodec.setDecodeTypes(decodeTypes);
+ public void loadPackableMethod(MethodDescriptor methodDescriptor) {
+ if (methodDescriptor instanceof PackableMethod) {
+ packableMethod = (PackableMethod) methodDescriptor;
+ return;
+ }
+ Map<MethodDescriptor, PackableMethod> cacheMap =
(Map<MethodDescriptor, PackableMethod>) url.getServiceModel()
+ .getServiceMetadata()
+ .getAttributeMap()
+ .computeIfAbsent(PACKABLE_METHOD_CACHE, k -> new
ConcurrentHashMap<>());
+ packableMethod = cacheMap.computeIfAbsent(methodDescriptor, md ->
frameworkModel
+ .getExtensionLoader(PackableMethodFactory.class)
+
.getExtension(ConfigurationUtils.getGlobalConfiguration(url.getApplicationModel())
+ .getString(DUBBO_PACKABLE_METHOD_FACTORY, DEFAULT_KEY))
+ .create(methodDescriptor, url, mediaType));
Review Comment:
**Security Issue**: The implementation of loading and caching
PackableMethod instances does not validate the MethodDescriptor objects,
potentially leading to arbitrary method execution if the object can be
manipulated by an attacker. <br> **Fix**: Implement strict validation of
MethodDescriptor objects before they are processed. Ensure that the
MethodDescriptor is from a trusted source and has not been tampered with. <br>
**Code Suggestion**:
```
+ if (!isValidMethodDescriptor(methodDescriptor)) {
+ throw new IllegalArgumentException("Invalid MethodDescriptor
detected.");
+ }
+ Map<MethodDescriptor, PackableMethod> cacheMap =
(Map<MethodDescriptor, PackableMethod>) url.getServiceModel()
.getServiceMetadata()
.getAttributeMap()
.computeIfAbsent(PACKABLE_METHOD_CACHE, k -> new
ConcurrentHashMap<>());
```
##########
dubbo-rpc/dubbo-rpc-triple/src/main/java/org/apache/dubbo/rpc/protocol/tri/h12/grpc/GrpcHttp2ServerTransportListener.java:
##########
@@ -145,39 +146,31 @@ public void onMessage(InputStream inputStream) {
private class DetermineMethodDescriptorListener implements
StreamingDecoder.FragmentListener {
- @Override
- public void onFragmentMessage(InputStream rawMessage) {}
-
@Override
public void onClose() {
getStreamingDecoder().close();
}
@Override
- public void onFragmentMessage(InputStream dataHeader, InputStream
rawMessage) {
+ public void onFragmentMessage(InputStream rawMessage) {
try {
- ByteArrayOutputStream merged =
- new ByteArrayOutputStream(dataHeader.available() +
rawMessage.available());
- StreamUtils.copy(dataHeader, merged);
- byte[] data = StreamUtils.readBytes(rawMessage);
-
RpcInvocationBuildContext context = getContext();
if (null == context.getMethodDescriptor()) {
-
context.setMethodDescriptor(DescriptorUtils.findTripleMethodDescriptor(
- context.getServiceDescriptor(),
context.getMethodName(), data));
+ byte[] data = StreamUtils.readBytes(rawMessage);
+ MethodDescriptor methodDescriptor =
DescriptorUtils.findTripleMethodDescriptor(
+ context.getServiceDescriptor(),
context.getMethodName(), data);
+ context.setMethodDescriptor(methodDescriptor);
setHttpMessageListener(GrpcHttp2ServerTransportListener.super.buildHttpMessageListener());
// replace decoder
GrpcCompositeCodec grpcCompositeCodec =
(GrpcCompositeCodec) context.getHttpMessageDecoder();
- MethodMetadata methodMetadata =
context.getMethodMetadata();
-
grpcCompositeCodec.setDecodeTypes(methodMetadata.getActualRequestTypes());
- grpcCompositeCodec.setEncodeTypes(new Class[]
{methodMetadata.getActualResponseType()});
+ grpcCompositeCodec.loadPackableMethod(methodDescriptor);
getServerChannelObserver().setResponseEncoder(grpcCompositeCodec);
+ rawMessage = new ByteArrayInputStream(data);
}
- merged.write(data);
- getHttpMessageListener().onMessage(new
ByteArrayInputStream(merged.toByteArray()));
+ getStreamingDecoder().invokeListener(rawMessage);
} catch (IOException e) {
throw new DecodeException(e);
}
Review Comment:
**Performance Issue**: The method 'onFragmentMessage' performs a rawMessage
read operation inside a try block without checking the size of the InputStream,
which could lead to inefficient memory usage and potential performance
degradation if the InputStream is significantly large. <br> **Fix**: Implement
a mechanism to check the size of the InputStream before reading it. If the size
exceeds a certain threshold, consider processing the stream in chunks or using
a more efficient way to handle large streams to optimize memory usage and
performance. <br> **Code Suggestion**:
```
Implement a mechanism to check the size of the InputStream before reading
it. If the size exceeds a certain threshold, consider processing the stream in
chunks or using a more efficient way to handle large streams to optimize memory
usage and performance.
```
##########
dubbo-rpc/dubbo-rpc-triple/src/main/java/org/apache/dubbo/rpc/protocol/tri/h12/grpc/GrpcCompositeCodec.java:
##########
@@ -16,123 +16,103 @@
*/
package org.apache.dubbo.rpc.protocol.tri.h12.grpc;
+import org.apache.dubbo.common.URL;
+import org.apache.dubbo.common.config.ConfigurationUtils;
+import org.apache.dubbo.common.io.StreamUtils;
+import org.apache.dubbo.common.utils.ArrayUtils;
import org.apache.dubbo.remoting.http12.exception.DecodeException;
import org.apache.dubbo.remoting.http12.exception.EncodeException;
import org.apache.dubbo.remoting.http12.message.HttpMessageCodec;
import org.apache.dubbo.remoting.http12.message.MediaType;
+import org.apache.dubbo.rpc.model.FrameworkModel;
+import org.apache.dubbo.rpc.model.MethodDescriptor;
+import org.apache.dubbo.rpc.model.PackableMethod;
+import org.apache.dubbo.rpc.model.PackableMethodFactory;
import java.io.IOException;
import java.io.InputStream;
import java.io.OutputStream;
import java.nio.charset.Charset;
+import java.util.Map;
+import java.util.concurrent.ConcurrentHashMap;
-import com.google.protobuf.Message;
-
-import static
org.apache.dubbo.common.constants.CommonConstants.PROTOBUF_MESSAGE_CLASS_NAME;
+import static org.apache.dubbo.common.constants.CommonConstants.DEFAULT_KEY;
+import static
org.apache.dubbo.common.constants.CommonConstants.DUBBO_PACKABLE_METHOD_FACTORY;
public class GrpcCompositeCodec implements HttpMessageCodec {
- private final ProtobufHttpMessageCodec protobufHttpMessageCodec;
+ private static final String PACKABLE_METHOD_CACHE =
"PACKABLE_METHOD_CACHE";
- private final WrapperHttpMessageCodec wrapperHttpMessageCodec;
+ private final URL url;
- public GrpcCompositeCodec(
- ProtobufHttpMessageCodec protobufHttpMessageCodec,
WrapperHttpMessageCodec wrapperHttpMessageCodec) {
- this.protobufHttpMessageCodec = protobufHttpMessageCodec;
- this.wrapperHttpMessageCodec = wrapperHttpMessageCodec;
- }
+ private final FrameworkModel frameworkModel;
+
+ private final String mediaType;
- public void setEncodeTypes(Class<?>[] encodeTypes) {
- this.wrapperHttpMessageCodec.setEncodeTypes(encodeTypes);
+ private PackableMethod packableMethod;
+
+ public GrpcCompositeCodec(URL url, FrameworkModel frameworkModel, String
mediaType) {
+ this.url = url;
+ this.frameworkModel = frameworkModel;
+ this.mediaType = mediaType;
}
- public void setDecodeTypes(Class<?>[] decodeTypes) {
- this.wrapperHttpMessageCodec.setDecodeTypes(decodeTypes);
+ public void loadPackableMethod(MethodDescriptor methodDescriptor) {
+ if (methodDescriptor instanceof PackableMethod) {
+ packableMethod = (PackableMethod) methodDescriptor;
+ return;
+ }
+ Map<MethodDescriptor, PackableMethod> cacheMap =
(Map<MethodDescriptor, PackableMethod>) url.getServiceModel()
+ .getServiceMetadata()
+ .getAttributeMap()
+ .computeIfAbsent(PACKABLE_METHOD_CACHE, k -> new
ConcurrentHashMap<>());
+ packableMethod = cacheMap.computeIfAbsent(methodDescriptor, md ->
frameworkModel
+ .getExtensionLoader(PackableMethodFactory.class)
+
.getExtension(ConfigurationUtils.getGlobalConfiguration(url.getApplicationModel())
+ .getString(DUBBO_PACKABLE_METHOD_FACTORY, DEFAULT_KEY))
+ .create(methodDescriptor, url, mediaType));
Review Comment:
**Optimization Issue**: The implementation of loading and caching
PackableMethod instances could be optimized for concurrent access patterns. The
current use of ConcurrentHashMap.computeIfAbsent is efficient but still
involves multiple steps and lambda expressions which could be streamlined
further for high-throughput scenarios. <br> **Fix**: Consider using a more
direct caching strategy that minimizes lambda creation for each computeIfAbsent
call. One approach could be to pre-load PackableMethods for known
MethodDescriptors during the initialization phase of GrpcCompositeCodec, thus
avoiding the need for computeIfAbsent during the critical path of method
invocation. <br> **Code Suggestion**:
```
+ Map<MethodDescriptor, PackableMethod> cacheMap =
preloadedPackableMethods;
+ packableMethod = cacheMap.get(methodDescriptor);
+ if (packableMethod == null) {
+ packableMethod = loadPackableMethod(methodDescriptor);
+ cacheMap.put(methodDescriptor, packableMethod);
+ }
```
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]