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]

Reply via email to