This is an automated email from the ASF dual-hosted git repository.

gosonzhang pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/inlong.git


The following commit(s) were added to refs/heads/master by this push:
     new 2b64674d6a [INLONG-8332][DataProxy] Return original content for 
MSG_ORIGINAL_RETURN type messages (#8333)
2b64674d6a is described below

commit 2b64674d6a61e4438ed864911380994324c2d3dc
Author: Goson Zhang <[email protected]>
AuthorDate: Tue Jun 27 17:00:54 2023 +0800

    [INLONG-8332][DataProxy] Return original content for MSG_ORIGINAL_RETURN 
type messages (#8333)
---
 .../sink/mq/MessageQueueZoneSinkContext.java       | 12 +++----
 .../dataproxy/source2/InLongMessageHandler.java    | 39 +++++++++++++++++++++-
 .../dataproxy/source2/v0msg/AbsV0MsgCodec.java     |  5 +++
 .../dataproxy/source2/v0msg/CodecTextMsg.java      |  4 +++
 4 files changed, 53 insertions(+), 7 deletions(-)

diff --git 
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/sink/mq/MessageQueueZoneSinkContext.java
 
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/sink/mq/MessageQueueZoneSinkContext.java
index becfade7bd..6989c94ea0 100644
--- 
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/sink/mq/MessageQueueZoneSinkContext.java
+++ 
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/sink/mq/MessageQueueZoneSinkContext.java
@@ -302,12 +302,12 @@ public class MessageQueueZoneSinkContext extends 
SinkContext {
         if (packProfile instanceof SimplePackProfile) {
             SimplePackProfile simpleProfile = (SimplePackProfile) packProfile;
             StringBuilder statsKey = new StringBuilder(512)
-                    .append(sinkName).append(AttrConstants.SEPARATOR)
-                    
.append(simpleProfile.getInlongGroupId()).append(AttrConstants.SEPARATOR)
-                    
.append(simpleProfile.getInlongStreamId()).append(AttrConstants.SEPARATOR)
-                    .append(topic).append(AttrConstants.SEPARATOR)
-                    
.append(NetworkUtils.getLocalIp()).append(AttrConstants.SEPARATOR)
-                    .append(remoteId).append(AttrConstants.SEPARATOR)
+                    .append(sinkName).append(AttrConstants.SEP_HASHTAG)
+                    
.append(simpleProfile.getInlongGroupId()).append(AttrConstants.SEP_HASHTAG)
+                    
.append(simpleProfile.getInlongStreamId()).append(AttrConstants.SEP_HASHTAG)
+                    .append(topic).append(AttrConstants.SEP_HASHTAG)
+                    
.append(NetworkUtils.getLocalIp()).append(AttrConstants.SEP_HASHTAG)
+                    .append(remoteId).append(AttrConstants.SEP_HASHTAG)
                     
.append(simpleProfile.getProperties().get(ConfigConstants.PKG_TIME_KEY));
             monitorIndex.addFailStats(statsKey.toString(), 1);
         }
diff --git 
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/source2/InLongMessageHandler.java
 
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/source2/InLongMessageHandler.java
index ad34d1e3b9..0a2d27c1c8 100644
--- 
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/source2/InLongMessageHandler.java
+++ 
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/source2/InLongMessageHandler.java
@@ -479,7 +479,7 @@ public class InLongMessageHandler extends 
ChannelInboundHandlerAdapter {
         if (MsgType.MSG_BIN_MULTI_BODY.equals(msgType)) {
             retData = buildBinMsgRspPackage(strBuff.toString(), 
msgObj.getUniq());
         } else {
-            retData = buildTxtMsgRspPackage(msgType, strBuff.toString());
+            retData = buildTxtMsgRspPackage(msgType, strBuff.toString(), 
msgObj);
         }
         strBuff.delete(0, strBuff.length());
         flushV0MsgPackage(source, channel, retData, msgObj.getAttr());
@@ -598,6 +598,43 @@ public class InLongMessageHandler extends 
ChannelInboundHandlerAdapter {
         return buffer;
     }
 
+    /**
+     * Build default-msg response message ByteBuf
+     *
+     * @param msgType  the message type
+     * @param attrs    the return attribute
+     * @param msgObj   the request message object
+     * @return ByteBuf
+     */
+    private ByteBuf buildTxtMsgRspPackage(MsgType msgType, String attrs, 
AbsV0MsgCodec msgObj) {
+        int attrsLen = 0;
+        int bodyLen = 0;
+        byte[] backBody = null;
+        if (attrs != null) {
+            attrsLen = attrs.length();
+        }
+        if (MsgType.MSG_ORIGINAL_RETURN.equals(msgType)) {
+            backBody = msgObj.getOrigBody();
+            if (backBody != null) {
+                bodyLen = backBody.length;
+            }
+        }
+        // backTotalLen = mstType + bodyLen + body + attrsLen + attrs
+        int backTotalLen = 1 + 4 + bodyLen + 4 + attrsLen;
+        ByteBuf buffer = ByteBufAllocator.DEFAULT.buffer(4 + backTotalLen);
+        buffer.writeInt(backTotalLen);
+        buffer.writeByte(msgType.getValue());
+        buffer.writeInt(bodyLen);
+        if (bodyLen > 0) {
+            buffer.writeBytes(backBody);
+        }
+        buffer.writeInt(attrsLen);
+        if (attrsLen > 0) {
+            buffer.writeBytes(attrs.getBytes(StandardCharsets.UTF_8));
+        }
+        return buffer;
+    }
+
     /**
      * Build heartbeat response message ByteBuf
      *
diff --git 
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/source2/v0msg/AbsV0MsgCodec.java
 
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/source2/v0msg/AbsV0MsgCodec.java
index 0764f7e599..fd367a3522 100644
--- 
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/source2/v0msg/AbsV0MsgCodec.java
+++ 
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/source2/v0msg/AbsV0MsgCodec.java
@@ -56,6 +56,7 @@ public abstract class AbsV0MsgCodec {
     protected int msgCount;
     protected String origAttr = "";
     protected byte[] bodyData;
+    protected byte[] origBody = null;
     protected long dataTimeMs;
     protected String groupId;
     protected String streamId = "";
@@ -144,6 +145,10 @@ public abstract class AbsV0MsgCodec {
         return strRemoteIP;
     }
 
+    public byte[] getOrigBody() {
+        return origBody;
+    }
+
     public long getMsgRcvTime() {
         return msgRcvTime;
     }
diff --git 
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/source2/v0msg/CodecTextMsg.java
 
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/source2/v0msg/CodecTextMsg.java
index 07a8f1aac5..c7189dbe80 100644
--- 
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/source2/v0msg/CodecTextMsg.java
+++ 
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/source2/v0msg/CodecTextMsg.java
@@ -71,6 +71,10 @@ public class CodecTextMsg extends AbsV0MsgCodec {
         // extract body bytes
         this.bodyData = new byte[bodyLen];
         cb.getBytes(msgHeadPos + TXT_MSG_BODY_OFFSET, this.bodyData, 0, 
bodyLen);
+        if (MsgType.MSG_ORIGINAL_RETURN.equals(MsgType.valueOf(msgType))) {
+            this.origBody = new byte[bodyLen];
+            System.arraycopy(this.bodyData, 0, this.origBody, 0, bodyLen);
+        }
         // get attribute length
         int attrLen = cb.getInt(msgHeadPos + TXT_MSG_BODY_OFFSET + bodyLen);
         if (attrLen < 0) {

Reply via email to