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

mikexue pushed a commit to branch cloudevents
in repository https://gitbox.apache.org/repos/asf/incubator-eventmesh.git


The following commit(s) were added to refs/heads/cloudevents by this push:
     new 6828f57  [Feature #562] Implement CloudEvents protocol adaptor (#595)
6828f57 is described below

commit 6828f57d130efbce5872ee5c4cfef5a0b7d5d283
Author: mike_xwm <[email protected]>
AuthorDate: Thu Nov 18 09:46:28 2021 +0800

    [Feature #562] Implement CloudEvents protocol adaptor (#595)
    
    * [Feature #564] Support CloudEvents protocols for pub/sub in 
EventMesh-feature design
    
    * support cloudevents api in eventmesh-connector-api module
    
    * fix checkStyle
    
    * fix checkStyle
    
    * fix checkStyle
    
    * 1.support LifeCycle.java
    2.update Consumer and Producer
    
    * fix remove the extra blank line
    
    * support cloudEvents
    
    * Add files via upload
    
    * Update README.md
    
    * support cloudEvents
    
    * support cloudEvents
    
    * [ISSUE #580] Add checkstyle gradle plugin (#581)
    
    * Add checkstyle gradle plugin, change plugin package
    
    * skip check in ci
    
    * support cloudEvents
    
    * support cloudevents
    
    * update wechat-official qr code
    
    * update mesh-helper qr code
    
    * Add files via upload
    
    * update README.md
    
    * update README.md
    
    * Update .asf.yaml
    
    * support cloudEvents
    
    * support cloudEvents
    
    * [ISSUE #588] Fix typo in README.md (#589)
    
    close #588
    
    * support cloudEvents
    
    * [Bug #590] Consumer subscription topic is invalid (#590) (#592)
    
    * [Bug #590] Consumer subscription topic is invalid (#590)
    
    * [Bug #590] Consumer subscription topic is invalid (#590)
    
    close #590
    
    * support cloudEvents adaptor
    
    * [Feature #562] Implement CloudEvents adaptor
    
    Co-authored-by: Eason Chen <[email protected]>
    Co-authored-by: Wenjun Ruan <[email protected]>
    Co-authored-by: Nicholas Zhan <[email protected]>
    Co-authored-by: hagsyn <[email protected]>
---
 .asf.yaml                                          |   3 +
 .github/workflows/ci.yml                           |   3 +-
 CONTRIBUTING.md                                    |   2 +
 CONTRIBUTING.zh-CN.md                              |   1 +
 README.md                                          |  34 ++++-----
 build.gradle                                       |  73 +++++++++++--------
 docs/images/Wechat-helper.jpeg                     | Bin 0 -> 41112 bytes
 docs/images/eventmesh-runtime2.png                 | Bin 0 -> 2041725 bytes
 docs/images/mesh-helper.jpg                        | Bin 0 -> 18197 bytes
 docs/images/wechat-official.png                    | Bin 21067 -> 16935 bytes
 .../common/protocol/http/common/ProtocolKey.java   |   6 ++
 .../header/message/ReplyMessageRequestHeader.java  |  36 ++++++++++
 .../message/SendMessageBatchRequestHeader.java     |  36 ++++++++++
 .../message/SendMessageBatchV2RequestHeader.java   |  37 ++++++++++
 .../header/message/SendMessageRequestHeader.java   |  37 ++++++++++
 .../eventmesh-connector-rocketmq/gradle.properties |   5 +-
 .../gradle.properties                              |   5 +-
 .../eventmesh/protocol/api/ProtocolAdaptor.java    |   2 +-
 .../gradle.properties                              |   3 +-
 .../cloudevents/CloudEventsProtocolAdaptor.java    |  65 ++++++++++++++---
 .../http/SendMessageBatchProtocolResolver.java     |  11 +++
 .../http/SendMessageBatchV2ProtocolResolver.java   |  78 ++++++++++++++++++++
 .../http/SendMessageRequestProtocolResolver.java   |  79 +++++++++++++++++++++
 .../resolver/tcp/TcpMessageProtocolResolver.java   |  51 +++++++++++++
 .../gradle.properties                              |   3 +-
 .../openmessage/OpenMessageProtocolAdaptor.java    |   2 +-
 .../eventmesh-registry-namesrv}/gradle.properties  |   3 +-
 .../protocol/http/consumer/EventMeshConsumer.java  |   6 +-
 .../http/processor/BatchSendMessageProcessor.java  |  11 +--
 .../processor/BatchSendMessageV2Processor.java     |  16 +++--
 .../http/processor/ReplyMessageProcessor.java      |  16 +++--
 .../http/processor/SendAsyncMessageProcessor.java  |  18 ++---
 .../http/processor/SendSyncMessageProcessor.java   |  16 +++--
 .../http/processor/SubscribeProcessor.java         |  13 +++-
 .../protocol/http/push/AsyncHTTPPushRequest.java   |   4 +-
 .../tcp/client/group/ClientGroupWrapper.java       |   8 +--
 .../tcp/client/session/push/SessionPusher.java     |   4 +-
 .../tcp/client/session/send/SessionSender.java     |   8 +--
 .../tcp/client/task/MessageTransferTask.java       |   8 +--
 .../eventmesh-security-acl}/gradle.properties      |   3 +-
 40 files changed, 587 insertions(+), 119 deletions(-)

diff --git a/.asf.yaml b/.asf.yaml
index 6428938..66dc6d7 100644
--- a/.asf.yaml
+++ b/.asf.yaml
@@ -19,6 +19,7 @@ github:
   description: EventMesh is a dynamic event-driven application runtime used to 
decouple the application and backend middleware layer, which supports a wide 
range of use cases that encompass complex multi-cloud, widely distributed 
topologies using diverse technology stacks.
   homepage: https://eventmesh.apache.org/
   labels:
+    - pubsub
     - event-mesh
     - event-gateway
     - event-driven
@@ -33,6 +34,8 @@ github:
     - message-bus
     - cqrs
     - multi-runtime
+    - microservice
+    - state-management
   enabled_merge_buttons:
     squash:  true
     merge:   false
diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml
index ff1c87e..f2b0aa2 100644
--- a/.github/workflows/ci.yml
+++ b/.github/workflows/ci.yml
@@ -63,7 +63,8 @@ jobs:
           java-version: ${{ matrix.java }}
 
       - name: Build
-        run: ./gradlew clean build jacocoTestReport checkLicense
+        # skip check here, since we use Checkstyle task to check the added file
+        run: ./gradlew clean build jacocoTestReport checkLicense -x check
 
       - name: Perform CodeQL analysis
         uses: github/codeql-action/analyze@v1
diff --git a/CONTRIBUTING.md b/CONTRIBUTING.md
index 659122c..16b6b66 100644
--- a/CONTRIBUTING.md
+++ b/CONTRIBUTING.md
@@ -19,6 +19,8 @@ Editor -> Code Style -> Java -> Scheme -> Import Scheme -> 
CheckStyle Configurat
 ```
 If you can't see CheckStyle Configuration section under Import Scheme, you can 
install CheckStyle-IDEA plugin first, and you will see it.
 
+You can also use `./gradlew check` to check the code style.
+(NOTE: this command will check all file in project, when you submit a pr, the 
ci will only check the file has been changed in this pr).
 ## Contributing
 
 We are always very happy to have contributions, whether for typo fix, bug fix 
or big new features. Please do not ever
diff --git a/CONTRIBUTING.zh-CN.md b/CONTRIBUTING.zh-CN.md
index 7d87ec8..0fede22 100644
--- a/CONTRIBUTING.zh-CN.md
+++ b/CONTRIBUTING.zh-CN.md
@@ -16,6 +16,7 @@ Editor -> Code Style -> Java -> Scheme -> Import Scheme -> 
CheckStyle Configurat
 ```
 如果你在Import Scheme下看不到CheckStyle 
Configuration选项,你可以先安装CheckStyle-IDEA插件,然后你就可以看到这个选项了。
 
+你也可以通过执行`./gradlew check`来检查代码格式。(NOTE: 这个命令将会检查整个项目中的代码格式, 
当你提交一个PR时,CI只会检查在此次PR中被被修改的文件的代码格式)
 ## 贡献
 
 无论是对于拼写错误,BUG修复还是重要的新功能,我们总是很乐意接受您的贡献。请不要犹豫,在Github Issue上提出或者通过邮件列表进行讨论。
diff --git a/README.md b/README.md
index 2a25858..cce12ce 100644
--- a/README.md
+++ b/README.md
@@ -10,22 +10,13 @@
 
 ![logo](docs/images/logo2.png)
 ## What is EventMesh?
-EventMesh(incubating) is a dynamic cloud-native eventing infrastruture used to 
decouple the application and backend middleware layer, which supports a wide 
range of use cases that encompass complex multi-cloud, widely distributed 
topologies using diverse technology stacks.
+EventMesh(incubating) is a dynamic event-driven application runtime used to 
decouple the application and backend middleware layer, which supports a wide 
range of use cases that encompass complex multi-cloud, widely distributed 
topologies using diverse technology stacks.
 
 ![architecture1](docs/images/eventmesh-multi-runtime.png)
 
-**EventMesh Ecosystem:**
-
-![architecture1](docs/images/eventmesh-define.png)
-
 **EventMesh Architecture:**
 
-![architecture1](docs/images/eventmesh-runtime.png)
-
-**EventMesh Cloud Native:**
-
-![architecture2](docs/images/eventmesh-panels.png)
-
+![architecture1](docs/images/eventmesh-runtime2.png)
 
 **Components:**
 
@@ -37,9 +28,11 @@ EventMesh(incubating) is a dynamic cloud-native eventing 
infrastruture used to d
 * **eventmesh-connector-rocketmq** : an implementation of 
eventmesh-connector-api, pub event to or sub event from RocketMQ as EventStore.
 * **eventmesh-connector-kafka(WIP)** : an implementation of 
eventmesh-connector-api, pub event to or sub event from Kafka as EventStore.
 * **eventmesh-connector-redis(WIP)** : an implementation of 
eventmesh-connector-api, pub event to or sub event from Redis as EventStore.
+* **eventmesh-connector-defibus(WIP)** : an implementation of 
eventmesh-connector-api, pub event to or sub event from 
[DeFiBus](https://github.com/webankfintech/defibus) as EventStore
 * **eventmesh-admin** : clients,topics,subscriptions and other management.
-* **eventmesh-registry-plugin** : plugins for registry.
-* **eventmesh-security-plugin** : plugins for security.
+* **eventmesh-registry-plugin** : plugins for registry adapter.
+* **eventmesh-security-plugin** : plugins for security adpater.
+* **eventmesh-protocol-plugin** : plugins for protocol adapter.
 
 **Protocol:**
 
@@ -54,9 +47,10 @@ Event & Service
 - [ ] Event transaction
 - [ ] At-least-once/at-most-once delivery guarantees
 
-Store
+Connector
 - [x] RocketMQ
 - [x] InMemory
+- [ ] Federated
 - [ ] Kafka
 - [ ] Redis
 - [ ] Pulsar
@@ -96,7 +90,7 @@ Governance
 - [x] Client management
 - [ ] Topic management
 - [ ] Metadata registry
-- [ ] Schema registry
+- [x] Schema registry
 - [ ] Dynamic config
 
 Choreography
@@ -107,11 +101,13 @@ Security
 - [ ] Auth
 - [ ] ACL
 
+Runtime
+- [ ] WebAssembly runtime
 
 ## Quick Start
-1. [Event-store](https://rocketmq.apache.org/docs/quick-start/) (RocketMQ, 
ignore this step if use standalone).
+1. [Connector quickstart](https://rocketmq.apache.org/docs/quick-start/) 
(RocketMQ, ignore this step if use standalone).
 2. [Runtime quickstart](docs/en/instructions/eventmesh-runtime-quickstart.md) 
or [Runtime quickstart with 
docker](docs/en/instructions/eventmesh-runtime-quickstart-with-docker.md).
-3. [Java examples ](docs/en/instructions/eventmesh-sdk-java-quickstart.md).
+3. [Java SDK examples](docs/en/instructions/eventmesh-sdk-java-quickstart.md).
 
 ## Contributing
 Contributions are always welcomed! Please see [CONTRIBUTING](CONTRIBUTING.md) 
for detailed guidelines.
@@ -131,9 +127,9 @@ EventMesh enriches the <a 
href="https://landscape.cncf.io/serverless?license=apa
 [Apache License, Version 2.0](http://www.apache.org/licenses/LICENSE-2.0.html) 
Copyright (C) Apache Software Foundation.
 
 ## Community
-|               WeChat group                |                WeChat official 
account                 |
+|               WeChat group                |                WeChat public 
account                 |
 | :---------------------------------------: | 
:----------------------------------------------------: |
-| ![wechat_qr](docs/images/mesh-helper.png) | 
![wechat_official_qr](docs/images/wechat-official.png) |
+| ![wechat_qr](docs/images/mesh-helper.jpg) | 
![wechat_official_qr](docs/images/wechat-official.png) |
 
 
 
diff --git a/build.gradle b/build.gradle
index 2559b74..2aa26fb 100644
--- a/build.gradle
+++ b/build.gradle
@@ -47,6 +47,7 @@ allprojects {
     apply plugin: "pmd"
     apply plugin: "java-library"
     apply plugin: 'signing'
+    apply plugin: 'checkstyle'
     apply plugin: 'com.github.jk1.dependency-license-report'
 
     [compileJava, compileTestJava, javadoc]*.options*.encoding = 'UTF-8'
@@ -78,6 +79,14 @@ allprojects {
         writer.flush()
         writer.close()
     }
+
+    checkstyle {
+        toolVersion = '9.0'
+        ignoreFailures = false
+        showViolations = true
+        maxWarnings = 0
+        configFile = new File("${rootDir}/style/checkStyle.xml")
+    }
 }
 
 task tar(type: Tar) {
@@ -103,38 +112,42 @@ task installPlugin() {
     if (!new File("${rootDir}/dist").exists()) {
         return
     }
-    // pluginType -> [pluginInstanceName -> moduleName]
-    Map<String, Map<String, String>> pluginTypeMap = [
-            "connector": ["rocketmq": "eventmesh-connector-rocketmq", 
"standalone": "eventmesh-connector-standalone",],
-            "security" : ["acl": "eventmesh-security-acl",],
-            "registry" : ["namesrv": "eventmesh-registry-namesrv",]
-    ]
     String[] libJars = java.util.Optional.ofNullable(new 
File("${rootDir}/dist/lib").list()).orElseGet(() -> new String[0])
     getAllprojects().forEach(subProject -> {
-        pluginTypeMap.forEach((pluginType, pluginInstanceMap) -> {
-            pluginInstanceMap.forEach((pluginInstanceName, moduleName) -> {
-                if (moduleName == subProject.name) {
-                    println String.format("install plugin, pluginType: %s, 
pluginInstanceName: %s, module: %s",
-                            pluginType, pluginInstanceName, moduleName)
-
-                    new 
File("${rootDir}/dist/plugin/${pluginType}/${pluginInstanceName}").mkdirs()
-                    copy {
-                        into 
"${rootDir}/dist/plugin/${pluginType}/${pluginInstanceName}"
-                        from "${subProject.getProjectDir()}/dist/apps"
-                    }
-                    copy {
-                        into 
"${rootDir}/dist/plugin/${pluginType}/${pluginInstanceName}"
-                        from "${subProject.getProjectDir()}/dist/lib/"
-                        exclude(libJars)
-                    }
-                    copy {
-                        into "${rootDir}/dist/conf"
-                        from "${subProject.getProjectDir()}/dist/conf"
-                        exclude 'META-INF'
-                    }
-                }
-            })
-        })
+        var file = new File("${subProject.projectDir}/gradle.properties")
+        if (!file.exists()) {
+            return
+        }
+        var properties = new Properties()
+        properties.load(new FileInputStream(file))
+        var pluginType = properties.getProperty("pluginType")
+        var pluginName = properties.getProperty("pluginName")
+        if (pluginType == null || pluginName == null) {
+            return
+        }
+        var pluginFile = new 
File("${rootDir}/dist/plugin/${pluginType}/${pluginName}")
+        if (pluginFile.exists()) {
+            return
+        }
+        pluginFile.mkdirs()
+        println String.format(
+                "install plugin, pluginType: %s, pluginInstanceName: %s, 
module: %s", pluginType, pluginName, subProject.getName()
+        )
+
+        copy {
+            into "${rootDir}/dist/plugin/${pluginType}/${pluginName}"
+            from "${subProject.getProjectDir()}/dist/apps"
+        }
+        copy {
+            into "${rootDir}/dist/plugin/${pluginType}/${pluginName}"
+            from "${subProject.getProjectDir()}/dist/lib/"
+            exclude(libJars)
+        }
+        copy {
+            into "${rootDir}/dist/conf"
+            from "${subProject.getProjectDir()}/dist/conf"
+            exclude 'META-INF'
+        }
     })
 }
 
diff --git a/docs/images/Wechat-helper.jpeg b/docs/images/Wechat-helper.jpeg
new file mode 100644
index 0000000..194599d
Binary files /dev/null and b/docs/images/Wechat-helper.jpeg differ
diff --git a/docs/images/eventmesh-runtime2.png 
b/docs/images/eventmesh-runtime2.png
new file mode 100644
index 0000000..73593f0
Binary files /dev/null and b/docs/images/eventmesh-runtime2.png differ
diff --git a/docs/images/mesh-helper.jpg b/docs/images/mesh-helper.jpg
new file mode 100644
index 0000000..1ab0263
Binary files /dev/null and b/docs/images/mesh-helper.jpg differ
diff --git a/docs/images/wechat-official.png b/docs/images/wechat-official.png
index 32f2896..ecb03a2 100644
Binary files a/docs/images/wechat-official.png and 
b/docs/images/wechat-official.png differ
diff --git 
a/eventmesh-common/src/main/java/org/apache/eventmesh/common/protocol/http/common/ProtocolKey.java
 
b/eventmesh-common/src/main/java/org/apache/eventmesh/common/protocol/http/common/ProtocolKey.java
index 28e944e..a41ab76 100644
--- 
a/eventmesh-common/src/main/java/org/apache/eventmesh/common/protocol/http/common/ProtocolKey.java
+++ 
b/eventmesh-common/src/main/java/org/apache/eventmesh/common/protocol/http/common/ProtocolKey.java
@@ -23,6 +23,12 @@ public class ProtocolKey {
     public static final String LANGUAGE = "Language";
     public static final String VERSION = "Version";
 
+    public static final String PROTOCOL_TYPE = "protocol_type";
+
+    public static final String PROTOCOL_VERSION = "protocol_version";
+
+    public static final String PROTOCOL_DESC = "protocol_desc";
+
     public static class ClientInstanceKey {
         ////////////////////////////////////Protocol layer requester 
description///////////
         public static final String ENV = "Env";
diff --git 
a/eventmesh-common/src/main/java/org/apache/eventmesh/common/protocol/http/header/message/ReplyMessageRequestHeader.java
 
b/eventmesh-common/src/main/java/org/apache/eventmesh/common/protocol/http/header/message/ReplyMessageRequestHeader.java
index d25936e..496b190 100644
--- 
a/eventmesh-common/src/main/java/org/apache/eventmesh/common/protocol/http/header/message/ReplyMessageRequestHeader.java
+++ 
b/eventmesh-common/src/main/java/org/apache/eventmesh/common/protocol/http/header/message/ReplyMessageRequestHeader.java
@@ -38,6 +38,15 @@ public class ReplyMessageRequestHeader extends Header {
     //protocol version adopted by requester, default:1.0
     private ProtocolVersion version;
 
+    //protocol type, cloudevents or eventmeshMessage
+    private String protocolType;
+
+    //protocol version, cloudevents:1.0 or 0.3
+    private String protocolVersion;
+
+    //protocol desc
+    private String protocolDesc;
+
     //the environment number of the requester
     private String env;
 
@@ -139,10 +148,37 @@ public class ReplyMessageRequestHeader extends Header {
         this.ip = ip;
     }
 
+    public String getProtocolType() {
+        return protocolType;
+    }
+
+    public void setProtocolType(String protocolType) {
+        this.protocolType = protocolType;
+    }
+
+    public String getProtocolVersion() {
+        return protocolVersion;
+    }
+
+    public void setProtocolVersion(String protocolVersion) {
+        this.protocolVersion = protocolVersion;
+    }
+
+    public String getProtocolDesc() {
+        return protocolDesc;
+    }
+
+    public void setProtocolDesc(String protocolDesc) {
+        this.protocolDesc = protocolDesc;
+    }
+
     public static ReplyMessageRequestHeader buildHeader(Map<String, Object> 
headerParam) {
         ReplyMessageRequestHeader header = new ReplyMessageRequestHeader();
         header.setCode(MapUtils.getString(headerParam, 
ProtocolKey.REQUEST_CODE));
         header.setVersion(ProtocolVersion.get(MapUtils.getString(headerParam, 
ProtocolKey.VERSION)));
+        header.setProtocolType(MapUtils.getString(headerParam, 
ProtocolKey.PROTOCOL_TYPE));
+        header.setProtocolVersion(MapUtils.getString(headerParam, 
ProtocolKey.PROTOCOL_VERSION));
+        header.setProtocolDesc(MapUtils.getString(headerParam, 
ProtocolKey.PROTOCOL_DESC));
         String lan = StringUtils.isBlank(MapUtils.getString(headerParam, 
ProtocolKey.LANGUAGE))
                 ? Constants.LANGUAGE_JAVA : MapUtils.getString(headerParam, 
ProtocolKey.LANGUAGE);
         header.setLanguage(lan);
diff --git 
a/eventmesh-common/src/main/java/org/apache/eventmesh/common/protocol/http/header/message/SendMessageBatchRequestHeader.java
 
b/eventmesh-common/src/main/java/org/apache/eventmesh/common/protocol/http/header/message/SendMessageBatchRequestHeader.java
index 62cafa4..85492cf 100644
--- 
a/eventmesh-common/src/main/java/org/apache/eventmesh/common/protocol/http/header/message/SendMessageBatchRequestHeader.java
+++ 
b/eventmesh-common/src/main/java/org/apache/eventmesh/common/protocol/http/header/message/SendMessageBatchRequestHeader.java
@@ -39,6 +39,15 @@ public class SendMessageBatchRequestHeader extends Header {
     //protocol version adopted by requester, default:1.0
     private ProtocolVersion version;
 
+    //protocol type, cloudevents or eventmeshMessage
+    private String protocolType;
+
+    //protocol version, cloudevents:1.0 or 0.3
+    private String protocolVersion;
+
+    //protocol desc
+    private String protocolDesc;
+
     //the environment number of the requester
     private String env;
 
@@ -140,10 +149,37 @@ public class SendMessageBatchRequestHeader extends Header 
{
         this.ip = ip;
     }
 
+    public String getProtocolType() {
+        return protocolType;
+    }
+
+    public void setProtocolType(String protocolType) {
+        this.protocolType = protocolType;
+    }
+
+    public String getProtocolVersion() {
+        return protocolVersion;
+    }
+
+    public void setProtocolVersion(String protocolVersion) {
+        this.protocolVersion = protocolVersion;
+    }
+
+    public String getProtocolDesc() {
+        return protocolDesc;
+    }
+
+    public void setProtocolDesc(String protocolDesc) {
+        this.protocolDesc = protocolDesc;
+    }
+
     public static SendMessageBatchRequestHeader buildHeader(final Map<String, 
Object> headerParam) {
         SendMessageBatchRequestHeader header = new 
SendMessageBatchRequestHeader();
         header.setCode(MapUtils.getString(headerParam, 
ProtocolKey.REQUEST_CODE));
         header.setVersion(ProtocolVersion.get(MapUtils.getString(headerParam, 
ProtocolKey.VERSION)));
+        header.setProtocolType(MapUtils.getString(headerParam, 
ProtocolKey.PROTOCOL_TYPE));
+        header.setProtocolVersion(MapUtils.getString(headerParam, 
ProtocolKey.PROTOCOL_VERSION));
+        header.setProtocolDesc(MapUtils.getString(headerParam, 
ProtocolKey.PROTOCOL_DESC));
         String lan = StringUtils.isBlank(MapUtils.getString(headerParam, 
ProtocolKey.LANGUAGE))
                 ? Constants.LANGUAGE_JAVA : MapUtils.getString(headerParam, 
ProtocolKey.LANGUAGE);
         header.setLanguage(lan);
diff --git 
a/eventmesh-common/src/main/java/org/apache/eventmesh/common/protocol/http/header/message/SendMessageBatchV2RequestHeader.java
 
b/eventmesh-common/src/main/java/org/apache/eventmesh/common/protocol/http/header/message/SendMessageBatchV2RequestHeader.java
index b5d2e14..15ff19a 100644
--- 
a/eventmesh-common/src/main/java/org/apache/eventmesh/common/protocol/http/header/message/SendMessageBatchV2RequestHeader.java
+++ 
b/eventmesh-common/src/main/java/org/apache/eventmesh/common/protocol/http/header/message/SendMessageBatchV2RequestHeader.java
@@ -38,6 +38,15 @@ public class SendMessageBatchV2RequestHeader extends Header {
     //protocol version adopted by requester, default:1.0
     private ProtocolVersion version;
 
+    //protocol type, cloudevents or eventmeshMessage
+    private String protocolType;
+
+    //protocol version, cloudevents:1.0 or 0.3
+    private String protocolVersion;
+
+    //protocol desc
+    private String protocolDesc;
+
     //the environment number of the requester
     private String env;
 
@@ -139,10 +148,38 @@ public class SendMessageBatchV2RequestHeader extends 
Header {
         this.ip = ip;
     }
 
+    public String getProtocolType() {
+        return protocolType;
+    }
+
+    public void setProtocolType(String protocolType) {
+        this.protocolType = protocolType;
+    }
+
+    public String getProtocolVersion() {
+        return protocolVersion;
+    }
+
+    public void setProtocolVersion(String protocolVersion) {
+        this.protocolVersion = protocolVersion;
+    }
+
+    public String getProtocolDesc() {
+        return protocolDesc;
+    }
+
+    public void setProtocolDesc(String protocolDesc) {
+        this.protocolDesc = protocolDesc;
+    }
+
     public static SendMessageBatchV2RequestHeader buildHeader(final 
Map<String, Object> headerParam) {
         SendMessageBatchV2RequestHeader header = new 
SendMessageBatchV2RequestHeader();
         header.setCode(MapUtils.getString(headerParam, 
ProtocolKey.REQUEST_CODE));
         header.setVersion(ProtocolVersion.get(MapUtils.getString(headerParam, 
ProtocolKey.VERSION)));
+        header.setProtocolType(MapUtils.getString(headerParam, 
ProtocolKey.PROTOCOL_TYPE));
+        header.setProtocolVersion(MapUtils.getString(headerParam, 
ProtocolKey.PROTOCOL_VERSION));
+        header.setProtocolDesc(MapUtils.getString(headerParam, 
ProtocolKey.PROTOCOL_DESC));
+
         String lan = StringUtils.isBlank(MapUtils.getString(headerParam, 
ProtocolKey.LANGUAGE))
                 ? Constants.LANGUAGE_JAVA : MapUtils.getString(headerParam, 
ProtocolKey.LANGUAGE);
         header.setLanguage(lan);
diff --git 
a/eventmesh-common/src/main/java/org/apache/eventmesh/common/protocol/http/header/message/SendMessageRequestHeader.java
 
b/eventmesh-common/src/main/java/org/apache/eventmesh/common/protocol/http/header/message/SendMessageRequestHeader.java
index 0eeea9a..f637563 100644
--- 
a/eventmesh-common/src/main/java/org/apache/eventmesh/common/protocol/http/header/message/SendMessageRequestHeader.java
+++ 
b/eventmesh-common/src/main/java/org/apache/eventmesh/common/protocol/http/header/message/SendMessageRequestHeader.java
@@ -38,6 +38,15 @@ public class SendMessageRequestHeader extends Header {
     //protocol version adopted by requester, default:1.0
     private ProtocolVersion version;
 
+    //protocol type, cloudevents or eventmeshMessage
+    private String protocolType;
+
+    //protocol version, cloudevents:1.0 or 0.3
+    private String protocolVersion;
+
+    //protocol desc
+    private String protocolDesc;
+
     //the environment number of the requester
     private String env;
 
@@ -139,10 +148,38 @@ public class SendMessageRequestHeader extends Header {
         this.ip = ip;
     }
 
+    public String getProtocolType() {
+        return protocolType;
+    }
+
+    public void setProtocolType(String protocolType) {
+        this.protocolType = protocolType;
+    }
+
+    public String getProtocolVersion() {
+        return protocolVersion;
+    }
+
+    public void setProtocolVersion(String protocolVersion) {
+        this.protocolVersion = protocolVersion;
+    }
+
+    public String getProtocolDesc() {
+        return protocolDesc;
+    }
+
+    public void setProtocolDesc(String protocolDesc) {
+        this.protocolDesc = protocolDesc;
+    }
+
     public static SendMessageRequestHeader buildHeader(Map<String, Object> 
headerParam) {
         SendMessageRequestHeader header = new SendMessageRequestHeader();
         header.setCode(MapUtils.getString(headerParam, 
ProtocolKey.REQUEST_CODE));
         header.setVersion(ProtocolVersion.get(MapUtils.getString(headerParam, 
ProtocolKey.VERSION)));
+        header.setProtocolType(MapUtils.getString(headerParam, 
ProtocolKey.PROTOCOL_TYPE));
+        header.setProtocolVersion(MapUtils.getString(headerParam, 
ProtocolKey.PROTOCOL_VERSION));
+        header.setProtocolDesc(MapUtils.getString(headerParam, 
ProtocolKey.PROTOCOL_DESC));
+
         String lan = StringUtils.isBlank(MapUtils.getString(headerParam, 
ProtocolKey.LANGUAGE))
                 ? Constants.LANGUAGE_JAVA : MapUtils.getString(headerParam, 
ProtocolKey.LANGUAGE);
         header.setLanguage(lan);
diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-rocketmq/gradle.properties 
b/eventmesh-connector-plugin/eventmesh-connector-rocketmq/gradle.properties
index 3d49f4c..4bcaa62 100644
--- a/eventmesh-connector-plugin/eventmesh-connector-rocketmq/gradle.properties
+++ b/eventmesh-connector-plugin/eventmesh-connector-rocketmq/gradle.properties
@@ -14,4 +14,7 @@
 # limitations under the License.
 #
 
-rocketmq_version=4.7.1
\ No newline at end of file
+rocketmq_version=4.7.1
+
+pluginType=connector
+pluginName=rocketmq
\ No newline at end of file
diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-standalone/gradle.properties 
b/eventmesh-connector-plugin/eventmesh-connector-standalone/gradle.properties
index 5e32028..9499e38 100644
--- 
a/eventmesh-connector-plugin/eventmesh-connector-standalone/gradle.properties
+++ 
b/eventmesh-connector-plugin/eventmesh-connector-standalone/gradle.properties
@@ -12,4 +12,7 @@
 # WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
 # See the License for the specific language governing permissions and
 # limitations under the License.
-#
\ No newline at end of file
+#
+
+pluginType=connector
+pluginName=standalone
\ No newline at end of file
diff --git 
a/eventmesh-protocol-plugin/eventmesh-protocol-api/src/main/java/org/apache/eventmesh/protocol/api/ProtocolAdaptor.java
 
b/eventmesh-protocol-plugin/eventmesh-protocol-api/src/main/java/org/apache/eventmesh/protocol/api/ProtocolAdaptor.java
index 82b88d7..9505df9 100644
--- 
a/eventmesh-protocol-plugin/eventmesh-protocol-api/src/main/java/org/apache/eventmesh/protocol/api/ProtocolAdaptor.java
+++ 
b/eventmesh-protocol-plugin/eventmesh-protocol-api/src/main/java/org/apache/eventmesh/protocol/api/ProtocolAdaptor.java
@@ -59,7 +59,7 @@ public interface ProtocolAdaptor<T> {
      * @param cloudEvent clout event
      * @return target protocol
      */
-    T fromCloudEvent(CloudEvent cloudEvent) throws ProtocolHandleException;
+    Object fromCloudEvent(CloudEvent cloudEvent) throws 
ProtocolHandleException;
 
     /**
      * Get protocol type.
diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-rocketmq/gradle.properties 
b/eventmesh-protocol-plugin/eventmesh-protocol-cloudevents/gradle.properties
similarity index 94%
copy from 
eventmesh-connector-plugin/eventmesh-connector-rocketmq/gradle.properties
copy to 
eventmesh-protocol-plugin/eventmesh-protocol-cloudevents/gradle.properties
index 3d49f4c..b3e1265 100644
--- a/eventmesh-connector-plugin/eventmesh-connector-rocketmq/gradle.properties
+++ b/eventmesh-protocol-plugin/eventmesh-protocol-cloudevents/gradle.properties
@@ -14,4 +14,5 @@
 # limitations under the License.
 #
 
-rocketmq_version=4.7.1
\ No newline at end of file
+pluginType=protocol
+pluginName=cloudevents
\ No newline at end of file
diff --git 
a/eventmesh-protocol-plugin/eventmesh-protocol-cloudevents/src/main/java/org/apache/eventmesh/protocol/cloudevents/CloudEventsProtocolAdaptor.java
 
b/eventmesh-protocol-plugin/eventmesh-protocol-cloudevents/src/main/java/org/apache/eventmesh/protocol/cloudevents/CloudEventsProtocolAdaptor.java
index fb4ca45..52c5321 100644
--- 
a/eventmesh-protocol-plugin/eventmesh-protocol-cloudevents/src/main/java/org/apache/eventmesh/protocol/cloudevents/CloudEventsProtocolAdaptor.java
+++ 
b/eventmesh-protocol-plugin/eventmesh-protocol-cloudevents/src/main/java/org/apache/eventmesh/protocol/cloudevents/CloudEventsProtocolAdaptor.java
@@ -18,12 +18,22 @@
 package org.apache.eventmesh.protocol.cloudevents;
 
 import io.cloudevents.CloudEvent;
+import org.apache.commons.lang3.StringUtils;
+import org.apache.eventmesh.common.Constants;
 import org.apache.eventmesh.common.command.HttpCommand;
+import org.apache.eventmesh.common.protocol.http.body.Body;
+import org.apache.eventmesh.common.protocol.http.common.RequestCode;
+import org.apache.eventmesh.common.protocol.tcp.Header;
 import org.apache.eventmesh.common.protocol.tcp.Package;
 import org.apache.eventmesh.protocol.api.ProtocolAdaptor;
 
 import org.apache.eventmesh.protocol.api.exception.ProtocolHandleException;
+import 
org.apache.eventmesh.protocol.cloudevents.resolver.http.SendMessageBatchProtocolResolver;
+import 
org.apache.eventmesh.protocol.cloudevents.resolver.http.SendMessageBatchV2ProtocolResolver;
+import 
org.apache.eventmesh.protocol.cloudevents.resolver.http.SendMessageRequestProtocolResolver;
+import 
org.apache.eventmesh.protocol.cloudevents.resolver.tcp.TcpMessageProtocolResolver;
 
+import java.nio.charset.StandardCharsets;
 import java.util.List;
 
 /**
@@ -34,14 +44,43 @@ import java.util.List;
 public class CloudEventsProtocolAdaptor<T> implements ProtocolAdaptor<T> {
 
     @Override
-    public CloudEvent toCloudEvent(T cloudEvent) {
+    public CloudEvent toCloudEvent(T cloudEvent) throws 
ProtocolHandleException {
 
-        if (cloudEvent instanceof Package){
-            //todo:convert package to cloudevents
-        }else if (cloudEvent instanceof HttpCommand){
-            //todo:convert httpCommand to cloudevents
+        if (cloudEvent instanceof Package) {
+            Header header = ((Package) cloudEvent).getHeader();
+            Object body = ((Package) cloudEvent).getBody();
+
+            return deserializeTcpProtocol(header, body);
+
+        } else if (cloudEvent instanceof HttpCommand) {
+            org.apache.eventmesh.common.protocol.http.header.Header header = 
((HttpCommand) cloudEvent).getHeader();
+            Body body = ((HttpCommand) cloudEvent).getBody();
+            String requestCode = ((HttpCommand) cloudEvent).getRequestCode();
+
+            return deserializeHttpProtocol(requestCode, header, body);
+        } else {
+            throw new ProtocolHandleException(String.format("protocol class: 
%s", cloudEvent.getClass()));
         }
-        return null;
+    }
+
+    private CloudEvent deserializeTcpProtocol(Header header, Object body) 
throws ProtocolHandleException {
+        return TcpMessageProtocolResolver.buildEvent(header, body);
+    }
+
+    private CloudEvent deserializeHttpProtocol(String requestCode, 
org.apache.eventmesh.common.protocol.http.header.Header header, Body body) 
throws ProtocolHandleException {
+
+        if 
(String.valueOf(RequestCode.MSG_BATCH_SEND.getRequestCode()).equals(requestCode))
 {
+            return SendMessageBatchProtocolResolver.buildEvent(header, body);
+        } else if 
(String.valueOf(RequestCode.MSG_BATCH_SEND_V2.getRequestCode()).equals(requestCode))
 {
+            return SendMessageBatchV2ProtocolResolver.buildEvent(header, body);
+        } else if 
(String.valueOf(RequestCode.MSG_SEND_SYNC.getRequestCode()).equals(requestCode))
 {
+            return SendMessageRequestProtocolResolver.buildEvent(header, body);
+        } else if 
(String.valueOf(RequestCode.MSG_SEND_ASYNC.getRequestCode()).equals(requestCode))
 {
+            return SendMessageRequestProtocolResolver.buildEvent(header, body);
+        } else {
+            throw new ProtocolHandleException(String.format("unsupported 
requestCode: %s", requestCode));
+        }
+
     }
 
     @Override
@@ -50,8 +89,18 @@ public class CloudEventsProtocolAdaptor<T> implements 
ProtocolAdaptor<T> {
     }
 
     @Override
-    public T fromCloudEvent(CloudEvent cloudEvent) {
-        return null;
+    public Object fromCloudEvent(CloudEvent cloudEvent) throws 
ProtocolHandleException {
+        String protocolDesc = 
cloudEvent.getExtension(Constants.PROTOCOL_DESC).toString();
+        if (StringUtils.equals("http", protocolDesc)) {
+            return new String(cloudEvent.getData().toBytes(), 
StandardCharsets.UTF_8);
+        } else if (StringUtils.equals("tcp", protocolDesc)) {
+            Package pkg = new Package();
+            pkg.setBody(cloudEvent);
+            return pkg;
+        } else {
+            throw new ProtocolHandleException(String.format("Unsupported 
protocolDesc: %s", protocolDesc));
+        }
+
     }
 
     @Override
diff --git 
a/eventmesh-protocol-plugin/eventmesh-protocol-cloudevents/src/main/java/org/apache/eventmesh/protocol/cloudevents/resolver/http/SendMessageBatchProtocolResolver.java
 
b/eventmesh-protocol-plugin/eventmesh-protocol-cloudevents/src/main/java/org/apache/eventmesh/protocol/cloudevents/resolver/http/SendMessageBatchProtocolResolver.java
new file mode 100644
index 0000000..6ea8d26
--- /dev/null
+++ 
b/eventmesh-protocol-plugin/eventmesh-protocol-cloudevents/src/main/java/org/apache/eventmesh/protocol/cloudevents/resolver/http/SendMessageBatchProtocolResolver.java
@@ -0,0 +1,11 @@
+package org.apache.eventmesh.protocol.cloudevents.resolver.http;
+
+import io.cloudevents.CloudEvent;
+import org.apache.eventmesh.common.protocol.http.body.Body;
+import org.apache.eventmesh.common.protocol.http.header.Header;
+
+public class SendMessageBatchProtocolResolver {
+    public static CloudEvent buildEvent(Header header, Body body) {
+        return null;
+    }
+}
diff --git 
a/eventmesh-protocol-plugin/eventmesh-protocol-cloudevents/src/main/java/org/apache/eventmesh/protocol/cloudevents/resolver/http/SendMessageBatchV2ProtocolResolver.java
 
b/eventmesh-protocol-plugin/eventmesh-protocol-cloudevents/src/main/java/org/apache/eventmesh/protocol/cloudevents/resolver/http/SendMessageBatchV2ProtocolResolver.java
new file mode 100644
index 0000000..d5ea9a5
--- /dev/null
+++ 
b/eventmesh-protocol-plugin/eventmesh-protocol-cloudevents/src/main/java/org/apache/eventmesh/protocol/cloudevents/resolver/http/SendMessageBatchV2ProtocolResolver.java
@@ -0,0 +1,78 @@
+package org.apache.eventmesh.protocol.cloudevents.resolver.http;
+
+import io.cloudevents.CloudEvent;
+import io.cloudevents.SpecVersion;
+import io.cloudevents.core.builder.CloudEventBuilder;
+import io.cloudevents.core.v03.CloudEventV03;
+import io.cloudevents.core.v1.CloudEventV1;
+import org.apache.commons.lang3.StringUtils;
+import org.apache.eventmesh.common.protocol.http.body.Body;
+import 
org.apache.eventmesh.common.protocol.http.body.message.SendMessageBatchV2RequestBody;
+import org.apache.eventmesh.common.protocol.http.common.ProtocolVersion;
+import org.apache.eventmesh.common.protocol.http.header.Header;
+import 
org.apache.eventmesh.common.protocol.http.header.message.SendMessageBatchV2RequestHeader;
+import org.apache.eventmesh.common.utils.JsonUtils;
+import org.apache.eventmesh.protocol.api.exception.ProtocolHandleException;
+
+public class SendMessageBatchV2ProtocolResolver {
+    public static CloudEvent buildEvent(Header header, Body body) throws 
ProtocolHandleException {
+        try {
+            SendMessageBatchV2RequestHeader sendMessageBatchV2RequestHeader = 
(SendMessageBatchV2RequestHeader) header;
+            SendMessageBatchV2RequestBody sendMessageBatchV2RequestBody = 
(SendMessageBatchV2RequestBody) body;
+
+            String protocolType = 
sendMessageBatchV2RequestHeader.getProtocolType();
+            String protocolDesc = 
sendMessageBatchV2RequestHeader.getProtocolDesc();
+            String protocolVersion = 
sendMessageBatchV2RequestHeader.getProtocolVersion();
+
+            String code = sendMessageBatchV2RequestHeader.getCode();
+            String env = sendMessageBatchV2RequestHeader.getEnv();
+            String idc = sendMessageBatchV2RequestHeader.getIdc();
+            String ip = sendMessageBatchV2RequestHeader.getIp();
+            String pid = sendMessageBatchV2RequestHeader.getPid();
+            String sys = sendMessageBatchV2RequestHeader.getSys();
+            String username = sendMessageBatchV2RequestHeader.getUsername();
+            String passwd = sendMessageBatchV2RequestHeader.getPasswd();
+            ProtocolVersion version = 
sendMessageBatchV2RequestHeader.getVersion();
+            String language = sendMessageBatchV2RequestHeader.getLanguage();
+
+            String content = sendMessageBatchV2RequestBody.getMsg();
+
+            CloudEvent event = null;
+            if (StringUtils.equals(SpecVersion.V1.toString(), 
protocolVersion)) {
+                event = JsonUtils.deserialize(content, CloudEventV1.class);
+                event = CloudEventBuilder.from(event)
+                        .withExtension("code", code)
+                        .withExtension("env", env)
+                        .withExtension("idc", idc)
+                        .withExtension("ip", ip)
+                        .withExtension("pid", pid)
+                        .withExtension("sys", sys)
+                        .withExtension("username", username)
+                        .withExtension("passwd", passwd)
+                        .withExtension("version", version.getVersion())
+                        .withExtension("language", language)
+                        .withExtension("protocolType", protocolType)
+                        .withExtension("protocolDesc", protocolDesc)
+                        .withExtension("protocolVersion", protocolVersion)
+                        .build();
+            } else if (StringUtils.equals(SpecVersion.V03.toString(), 
protocolVersion)) {
+                event = JsonUtils.deserialize(content, CloudEventV03.class);
+                event = CloudEventBuilder.from(event)
+                        .withExtension("code", code)
+                        .withExtension("env", env)
+                        .withExtension("idc", idc)
+                        .withExtension("ip", ip)
+                        .withExtension("pid", pid)
+                        .withExtension("sys", sys)
+                        .withExtension("username", username)
+                        .withExtension("passwd", passwd)
+                        .withExtension("version", version.getVersion())
+                        .withExtension("language", language)
+                        .build();
+            }
+            return event;
+        } catch (Exception e) {
+            throw new ProtocolHandleException(e.getMessage(), e.getCause());
+        }
+    }
+}
diff --git 
a/eventmesh-protocol-plugin/eventmesh-protocol-cloudevents/src/main/java/org/apache/eventmesh/protocol/cloudevents/resolver/http/SendMessageRequestProtocolResolver.java
 
b/eventmesh-protocol-plugin/eventmesh-protocol-cloudevents/src/main/java/org/apache/eventmesh/protocol/cloudevents/resolver/http/SendMessageRequestProtocolResolver.java
new file mode 100644
index 0000000..8055625
--- /dev/null
+++ 
b/eventmesh-protocol-plugin/eventmesh-protocol-cloudevents/src/main/java/org/apache/eventmesh/protocol/cloudevents/resolver/http/SendMessageRequestProtocolResolver.java
@@ -0,0 +1,79 @@
+package org.apache.eventmesh.protocol.cloudevents.resolver.http;
+
+import io.cloudevents.CloudEvent;
+import io.cloudevents.SpecVersion;
+import io.cloudevents.core.builder.CloudEventBuilder;
+import io.cloudevents.core.v03.CloudEventV03;
+import io.cloudevents.core.v1.CloudEventV1;
+import org.apache.commons.lang3.StringUtils;
+import org.apache.eventmesh.common.protocol.http.body.Body;
+import 
org.apache.eventmesh.common.protocol.http.body.message.SendMessageRequestBody;
+import org.apache.eventmesh.common.protocol.http.common.ProtocolVersion;
+import org.apache.eventmesh.common.protocol.http.header.Header;
+import 
org.apache.eventmesh.common.protocol.http.header.message.SendMessageRequestHeader;
+import org.apache.eventmesh.common.utils.JsonUtils;
+import org.apache.eventmesh.protocol.api.exception.ProtocolHandleException;
+
+public class SendMessageRequestProtocolResolver {
+
+    public static CloudEvent buildEvent(Header header, Body body) throws 
ProtocolHandleException {
+        try {
+            SendMessageRequestHeader sendMessageRequestHeader = 
(SendMessageRequestHeader) header;
+            SendMessageRequestBody sendMessageRequestBody = 
(SendMessageRequestBody) body;
+
+            String protocolType = sendMessageRequestHeader.getProtocolType();
+            String protocolDesc = sendMessageRequestHeader.getProtocolDesc();
+            String protocolVersion = 
sendMessageRequestHeader.getProtocolVersion();
+
+            String code = sendMessageRequestHeader.getCode();
+            String env = sendMessageRequestHeader.getEnv();
+            String idc = sendMessageRequestHeader.getIdc();
+            String ip = sendMessageRequestHeader.getIp();
+            String pid = sendMessageRequestHeader.getPid();
+            String sys = sendMessageRequestHeader.getSys();
+            String username = sendMessageRequestHeader.getUsername();
+            String passwd = sendMessageRequestHeader.getPasswd();
+            ProtocolVersion version = sendMessageRequestHeader.getVersion();
+            String language = sendMessageRequestHeader.getLanguage();
+
+            String content = sendMessageRequestBody.getContent();
+
+            CloudEvent event = null;
+            if (StringUtils.equals(SpecVersion.V1.toString(), 
protocolVersion)) {
+                event = JsonUtils.deserialize(content, CloudEventV1.class);
+                event = CloudEventBuilder.from(event)
+                        .withExtension("code", code)
+                        .withExtension("env", env)
+                        .withExtension("idc", idc)
+                        .withExtension("ip", ip)
+                        .withExtension("pid", pid)
+                        .withExtension("sys", sys)
+                        .withExtension("username", username)
+                        .withExtension("passwd", passwd)
+                        .withExtension("version", version.getVersion())
+                        .withExtension("language", language)
+                        .withExtension("protocolType", protocolType)
+                        .withExtension("protocolDesc", protocolDesc)
+                        .withExtension("protocolVersion", protocolVersion)
+                        .build();
+            } else if (StringUtils.equals(SpecVersion.V03.toString(), 
protocolVersion)) {
+                event = JsonUtils.deserialize(content, CloudEventV03.class);
+                event = CloudEventBuilder.from(event)
+                        .withExtension("code", code)
+                        .withExtension("env", env)
+                        .withExtension("idc", idc)
+                        .withExtension("ip", ip)
+                        .withExtension("pid", pid)
+                        .withExtension("sys", sys)
+                        .withExtension("username", username)
+                        .withExtension("passwd", passwd)
+                        .withExtension("version", version.getVersion())
+                        .withExtension("language", language)
+                        .build();
+            }
+            return event;
+        } catch (Exception e) {
+            throw new ProtocolHandleException(e.getMessage(), e.getCause());
+        }
+    }
+}
diff --git 
a/eventmesh-protocol-plugin/eventmesh-protocol-cloudevents/src/main/java/org/apache/eventmesh/protocol/cloudevents/resolver/tcp/TcpMessageProtocolResolver.java
 
b/eventmesh-protocol-plugin/eventmesh-protocol-cloudevents/src/main/java/org/apache/eventmesh/protocol/cloudevents/resolver/tcp/TcpMessageProtocolResolver.java
new file mode 100644
index 0000000..3403fb5
--- /dev/null
+++ 
b/eventmesh-protocol-plugin/eventmesh-protocol-cloudevents/src/main/java/org/apache/eventmesh/protocol/cloudevents/resolver/tcp/TcpMessageProtocolResolver.java
@@ -0,0 +1,51 @@
+package org.apache.eventmesh.protocol.cloudevents.resolver.tcp;
+
+import io.cloudevents.CloudEvent;
+import io.cloudevents.SpecVersion;
+import io.cloudevents.core.builder.CloudEventBuilder;
+import org.apache.commons.lang3.StringUtils;
+import org.apache.eventmesh.common.Constants;
+import org.apache.eventmesh.common.protocol.tcp.Header;
+import org.apache.eventmesh.protocol.api.exception.ProtocolHandleException;
+
+public class TcpMessageProtocolResolver {
+
+    public static CloudEvent buildEvent(Header header, Object body) throws 
ProtocolHandleException {
+        CloudEventBuilder cloudEventBuilder;
+
+        String protocolType = 
header.getProperty(Constants.PROTOCOL_TYPE).toString();
+        String protocolVersion = 
header.getProperty(Constants.PROTOCOL_VERSION).toString();
+        String protocolDesc = 
header.getProperty(Constants.PROTOCOL_DESC).toString();
+
+        if (StringUtils.isBlank(protocolType)
+                || StringUtils.isBlank(protocolVersion)
+                || StringUtils.isBlank(protocolDesc)) {
+            throw new ProtocolHandleException(String.format("invalid protocol 
params protocolType %s|protocolVersion %s|protocolDesc %s",
+                    protocolType, protocolVersion, protocolDesc));
+        }
+
+        if (!StringUtils.equals("cloudevents", protocolType)) {
+            throw new ProtocolHandleException(String.format("Unsupported 
protocolType: %s", protocolType));
+        }
+        if (StringUtils.equals(SpecVersion.V1.toString(), protocolVersion)) {
+            cloudEventBuilder = CloudEventBuilder.v1((CloudEvent) body);
+
+            for (String propKey : header.getProperties().keySet()) {
+                cloudEventBuilder.withExtension(propKey, 
header.getProperty(propKey).toString());
+            }
+
+            return cloudEventBuilder.build();
+
+        } else if (StringUtils.equals(SpecVersion.V03.toString(), 
protocolVersion)) {
+            cloudEventBuilder = CloudEventBuilder.v03((CloudEvent) body);
+
+            for (String propKey : header.getProperties().keySet()) {
+                cloudEventBuilder.withExtension(propKey, 
header.getProperty(propKey).toString());
+            }
+
+            return cloudEventBuilder.build();
+        } else {
+            throw new ProtocolHandleException(String.format("Unsupported 
protocolVersion: %s", protocolVersion));
+        }
+    }
+}
diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-rocketmq/gradle.properties 
b/eventmesh-protocol-plugin/eventmesh-protocol-openmessage/gradle.properties
similarity index 94%
copy from 
eventmesh-connector-plugin/eventmesh-connector-rocketmq/gradle.properties
copy to 
eventmesh-protocol-plugin/eventmesh-protocol-openmessage/gradle.properties
index 3d49f4c..c414dfb 100644
--- a/eventmesh-connector-plugin/eventmesh-connector-rocketmq/gradle.properties
+++ b/eventmesh-protocol-plugin/eventmesh-protocol-openmessage/gradle.properties
@@ -14,4 +14,5 @@
 # limitations under the License.
 #
 
-rocketmq_version=4.7.1
\ No newline at end of file
+pluginType=protocol
+pluginName=openmessage
\ No newline at end of file
diff --git 
a/eventmesh-protocol-plugin/eventmesh-protocol-openmessage/src/main/java/org/apache/eventmesh/protocol/openmessage/OpenMessageProtocolAdaptor.java
 
b/eventmesh-protocol-plugin/eventmesh-protocol-openmessage/src/main/java/org/apache/eventmesh/protocol/openmessage/OpenMessageProtocolAdaptor.java
index 55b7c6f..6e00275 100644
--- 
a/eventmesh-protocol-plugin/eventmesh-protocol-openmessage/src/main/java/org/apache/eventmesh/protocol/openmessage/OpenMessageProtocolAdaptor.java
+++ 
b/eventmesh-protocol-plugin/eventmesh-protocol-openmessage/src/main/java/org/apache/eventmesh/protocol/openmessage/OpenMessageProtocolAdaptor.java
@@ -46,7 +46,7 @@ public class OpenMessageProtocolAdaptor<T> implements 
ProtocolAdaptor<T> {
     }
 
     @Override
-    public T fromCloudEvent(CloudEvent cloudEvent) {
+    public Object fromCloudEvent(CloudEvent cloudEvent) {
         return null;
     }
 
diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-rocketmq/gradle.properties 
b/eventmesh-registry-plugin/eventmesh-registry-namesrv/gradle.properties
similarity index 95%
copy from 
eventmesh-connector-plugin/eventmesh-connector-rocketmq/gradle.properties
copy to eventmesh-registry-plugin/eventmesh-registry-namesrv/gradle.properties
index 3d49f4c..ace94b4 100644
--- a/eventmesh-connector-plugin/eventmesh-connector-rocketmq/gradle.properties
+++ b/eventmesh-registry-plugin/eventmesh-registry-namesrv/gradle.properties
@@ -14,4 +14,5 @@
 # limitations under the License.
 #
 
-rocketmq_version=4.7.1
\ No newline at end of file
+pluginType=registry
+pluginName=namesrv
\ No newline at end of file
diff --git 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/consumer/EventMeshConsumer.java
 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/consumer/EventMeshConsumer.java
index 1cb62bf..fd4a704 100644
--- 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/consumer/EventMeshConsumer.java
+++ 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/consumer/EventMeshConsumer.java
@@ -23,7 +23,7 @@ import java.util.concurrent.atomic.AtomicBoolean;
 
 import io.cloudevents.CloudEvent;
 
-import io.cloudevents.core.v1.CloudEventBuilder;
+import io.cloudevents.core.builder.CloudEventBuilder;
 import org.apache.commons.collections4.MapUtils;
 import org.apache.eventmesh.api.*;
 import org.apache.eventmesh.api.exception.OnExceptionContext;
@@ -115,7 +115,7 @@ public class EventMeshConsumer {
                     String bizSeqNo = (String) 
event.getExtension(Constants.PROPERTY_MESSAGE_SEARCH_KEYS);
                     String uniqueId = (String) 
event.getExtension(Constants.RMB_UNIQ_ID);
 
-                    event = new CloudEventBuilder(event)
+                    event = CloudEventBuilder.from(event)
                             
.withExtension(EventMeshConstants.REQ_MQ2EVENTMESH_TIMESTAMP, 
String.valueOf(System.currentTimeMillis()))
                             .build();
 //                    
message.getUserProperties().put(EventMeshConstants.REQ_MQ2EVENTMESH_TIMESTAMP, 
String.valueOf(System.currentTimeMillis()));
@@ -164,7 +164,7 @@ public class EventMeshConsumer {
                 @Override
                 public void consume(CloudEvent event, AsyncConsumeContext 
context) {
 
-                    event = new CloudEventBuilder(event)
+                    event = CloudEventBuilder.from(event)
                             
.withExtension(EventMeshConstants.REQ_MQ2EVENTMESH_TIMESTAMP, 
String.valueOf(System.currentTimeMillis()))
                             .build();
 
diff --git 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/BatchSendMessageProcessor.java
 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/BatchSendMessageProcessor.java
index 73650f8..fe18a31 100644
--- 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/BatchSendMessageProcessor.java
+++ 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/BatchSendMessageProcessor.java
@@ -19,8 +19,7 @@ package 
org.apache.eventmesh.runtime.core.protocol.http.processor;
 
 import io.cloudevents.CloudEvent;
 import io.cloudevents.CloudEventData;
-import io.cloudevents.core.v1.CloudEventBuilder;
-import io.cloudevents.core.v1.CloudEventV1;
+import io.cloudevents.core.builder.CloudEventBuilder;
 import io.netty.channel.ChannelHandlerContext;
 import org.apache.commons.collections4.CollectionUtils;
 import org.apache.commons.lang3.StringUtils;
@@ -82,13 +81,15 @@ public class BatchSendMessageProcessor implements 
HttpRequestProcessor {
                 EventMeshConstants.PROTOCOL_HTTP,
                 RemotingHelper.parseChannelRemoteAddr(ctx.channel()), 
IPUtil.getLocalAddress());
 
+        SendMessageBatchRequestHeader sendMessageBatchRequestHeader = 
(SendMessageBatchRequestHeader) asyncContext.getRequest().getHeader();
+
         SendMessageBatchResponseHeader sendMessageBatchResponseHeader =
                 
SendMessageBatchResponseHeader.buildHeader(Integer.valueOf(asyncContext.getRequest().getRequestCode()),
 eventMeshHTTPServer.getEventMeshHttpConfiguration().eventMeshCluster,
                         IPUtil.getLocalAddress(), 
eventMeshHTTPServer.getEventMeshHttpConfiguration().eventMeshEnv,
                         
eventMeshHTTPServer.getEventMeshHttpConfiguration().eventMeshIDC);
 
-
-        ProtocolAdaptor httpCommandProtocolAdaptor = 
ProtocolPluginFactory.getProtocolAdaptor("cloudevents");
+        String protocolType = sendMessageBatchRequestHeader.getProtocolType();
+        ProtocolAdaptor httpCommandProtocolAdaptor = 
ProtocolPluginFactory.getProtocolAdaptor(protocolType);
         List<CloudEvent> eventList = 
httpCommandProtocolAdaptor.toBatchCloudEvent(asyncContext.getRequest());
 
         if (CollectionUtils.isEmpty(eventList)) {
@@ -212,7 +213,7 @@ public class BatchSendMessageProcessor implements 
HttpRequestProcessor {
                 String ttl = 
Objects.requireNonNull(cloudEvent.getExtension(SendMessageRequestBody.TTL)).toString();
 
                 if (StringUtils.isBlank(ttl) || !StringUtils.isNumeric(ttl)) {
-                    cloudEvent = new CloudEventBuilder(cloudEvent)
+                    cloudEvent = CloudEventBuilder.from(cloudEvent)
                             .withExtension(SendMessageRequestBody.TTL, 
String.valueOf(EventMeshConstants.DEFAULT_MSG_TTL_MILLS))
                             .withExtension("msgType", "persistent")
                             .build();
diff --git 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/BatchSendMessageV2Processor.java
 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/BatchSendMessageV2Processor.java
index 1c38933..c38ff5d 100644
--- 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/BatchSendMessageV2Processor.java
+++ 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/BatchSendMessageV2Processor.java
@@ -18,7 +18,7 @@
 package org.apache.eventmesh.runtime.core.protocol.http.processor;
 
 import io.cloudevents.CloudEvent;
-import io.cloudevents.core.v1.CloudEventBuilder;
+import io.cloudevents.core.builder.CloudEventBuilder;
 import io.netty.channel.ChannelHandlerContext;
 import org.apache.commons.lang3.StringUtils;
 import org.apache.eventmesh.api.SendCallback;
@@ -74,8 +74,11 @@ public class BatchSendMessageV2Processor implements 
HttpRequestProcessor {
                 EventMeshConstants.PROTOCOL_HTTP,
                 RemotingHelper.parseChannelRemoteAddr(ctx.channel()), 
IPUtil.getLocalAddress());
 
-        ProtocolAdaptor httpCommandProtocolAdaptor = 
ProtocolPluginFactory.getProtocolAdaptor("cloudevents");
-        CloudEvent event = 
httpCommandProtocolAdaptor.toCloudEventV1(asyncContext.getRequest());
+        SendMessageBatchV2RequestHeader sendMessageBatchV2RequestHeader = 
(SendMessageBatchV2RequestHeader) asyncContext.getRequest().getHeader();
+
+        String protocolType = 
sendMessageBatchV2RequestHeader.getProtocolType();
+        ProtocolAdaptor httpCommandProtocolAdaptor = 
ProtocolPluginFactory.getProtocolAdaptor(protocolType);
+        CloudEvent event = 
httpCommandProtocolAdaptor.toCloudEvent(asyncContext.getRequest());
 
         SendMessageBatchV2ResponseHeader sendMessageBatchV2ResponseHeader =
                 
SendMessageBatchV2ResponseHeader.buildHeader(Integer.valueOf(asyncContext.getRequest().getRequestCode()),
 eventMeshHTTPServer.getEventMeshHttpConfiguration().eventMeshCluster,
@@ -83,7 +86,8 @@ public class BatchSendMessageV2Processor implements 
HttpRequestProcessor {
                         
eventMeshHTTPServer.getEventMeshHttpConfiguration().eventMeshIDC);
 
         //validate event
-        if (StringUtils.isBlank(event.getId())
+        if (event != null
+                || StringUtils.isBlank(event.getId())
                 || event.getSource() != null
                 || event.getSpecVersion() != null
                 || StringUtils.isBlank(event.getType())
@@ -173,12 +177,12 @@ public class BatchSendMessageV2Processor implements 
HttpRequestProcessor {
         String ttl = String.valueOf(EventMeshConstants.DEFAULT_MSG_TTL_MILLS);
         if 
(StringUtils.isBlank(event.getExtension(SendMessageRequestBody.TTL).toString())
                 && 
!StringUtils.isNumeric(event.getExtension(SendMessageRequestBody.TTL).toString()))
 {
-            event = new 
CloudEventBuilder(event).withExtension(SendMessageRequestBody.TTL, ttl).build();
+            event = 
CloudEventBuilder.from(event).withExtension(SendMessageRequestBody.TTL, 
ttl).build();
         }
 
 
         try {
-            event = new CloudEventBuilder(event)
+            event = CloudEventBuilder.from(event)
                     .withExtension("msgType", "persistent")
                     
.withExtension(EventMeshConstants.REQ_C2EVENTMESH_TIMESTAMP, 
String.valueOf(System.currentTimeMillis()))
                     
.withExtension(EventMeshConstants.REQ_EVENTMESH2MQ_TIMESTAMP, 
String.valueOf(System.currentTimeMillis()))
diff --git 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/ReplyMessageProcessor.java
 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/ReplyMessageProcessor.java
index 3050a8b..7630315 100644
--- 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/ReplyMessageProcessor.java
+++ 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/ReplyMessageProcessor.java
@@ -18,7 +18,7 @@
 package org.apache.eventmesh.runtime.core.protocol.http.processor;
 
 import io.cloudevents.CloudEvent;
-import io.cloudevents.core.v1.CloudEventBuilder;
+import io.cloudevents.core.builder.CloudEventBuilder;
 import io.netty.channel.ChannelHandlerContext;
 import org.apache.commons.collections4.MapUtils;
 import org.apache.commons.lang3.StringUtils;
@@ -77,11 +77,12 @@ public class ReplyMessageProcessor implements 
HttpRequestProcessor {
                 EventMeshConstants.PROTOCOL_HTTP,
                 RemotingHelper.parseChannelRemoteAddr(ctx.channel()), 
IPUtil.getLocalAddress());
 
-//        ReplyMessageRequestHeader replyMessageRequestHeader = 
(ReplyMessageRequestHeader) asyncContext.getRequest().getHeader();
+        ReplyMessageRequestHeader replyMessageRequestHeader = 
(ReplyMessageRequestHeader) asyncContext.getRequest().getHeader();
 //        ReplyMessageRequestBody replyMessageRequestBody = 
(ReplyMessageRequestBody) asyncContext.getRequest().getBody();
 
-        ProtocolAdaptor httpCommandProtocolAdaptor = 
ProtocolPluginFactory.getProtocolAdaptor("cloudevents");
-        CloudEvent event = 
httpCommandProtocolAdaptor.toCloudEventV1(asyncContext.getRequest());
+        String protocolType = replyMessageRequestHeader.getProtocolType();
+        ProtocolAdaptor httpCommandProtocolAdaptor = 
ProtocolPluginFactory.getProtocolAdaptor(protocolType);
+        CloudEvent event = 
httpCommandProtocolAdaptor.toCloudEvent(asyncContext.getRequest());
 
         ReplyMessageResponseHeader replyMessageResponseHeader =
                 
ReplyMessageResponseHeader.buildHeader(Integer.valueOf(asyncContext.getRequest().getRequestCode()),
 eventMeshHTTPServer.getEventMeshHttpConfiguration().eventMeshCluster,
@@ -89,7 +90,8 @@ public class ReplyMessageProcessor implements 
HttpRequestProcessor {
                         
eventMeshHTTPServer.getEventMeshHttpConfiguration().eventMeshIDC);
 
         //validate event
-        if (StringUtils.isBlank(event.getId())
+        if (event != null
+                || StringUtils.isBlank(event.getId())
                 || event.getSource() != null
                 || event.getSpecVersion() != null
                 || StringUtils.isBlank(event.getType())
@@ -179,7 +181,7 @@ public class ReplyMessageProcessor implements 
HttpRequestProcessor {
         try {
             // body
 //            
omsMsg.setBody(replyMessageRequestBody.getContent().getBytes(EventMeshConstants.DEFAULT_CHARSET));
-            event = new CloudEventBuilder(event)
+            event = CloudEventBuilder.from(event)
                     .withSubject(replyTopic)
                     .withExtension("msgType", "persistent")
                     .withExtension(Constants.PROPERTY_MESSAGE_TIMEOUT, 
String.valueOf(EventMeshConstants.DEFAULT_TIMEOUT_IN_MILLISECONDS))
@@ -239,7 +241,7 @@ public class ReplyMessageProcessor implements 
HttpRequestProcessor {
 
 
         try {
-            CloudEvent clone = new 
CloudEventBuilder(sendMessageContext.getEvent())
+            CloudEvent clone = 
CloudEventBuilder.from(sendMessageContext.getEvent())
                     
.withExtension(EventMeshConstants.REQ_EVENTMESH2MQ_TIMESTAMP, 
String.valueOf(System.currentTimeMillis()))
                     .build();
             sendMessageContext.setEvent(clone);
diff --git 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/SendAsyncMessageProcessor.java
 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/SendAsyncMessageProcessor.java
index c6d6cc1..768b0cc 100644
--- 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/SendAsyncMessageProcessor.java
+++ 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/SendAsyncMessageProcessor.java
@@ -17,6 +17,7 @@
 
 package org.apache.eventmesh.runtime.core.protocol.http.processor;
 
+import io.cloudevents.core.builder.CloudEventBuilder;
 import jdk.nashorn.internal.runtime.URIUtils;
 import org.apache.commons.lang3.StringUtils;
 import org.apache.eventmesh.api.SendCallback;
@@ -53,7 +54,6 @@ import java.util.Objects;
 import java.util.concurrent.TimeUnit;
 
 import io.cloudevents.CloudEvent;
-import io.cloudevents.core.v1.CloudEventBuilder;
 import io.netty.channel.ChannelHandlerContext;
 
 public class SendAsyncMessageProcessor implements HttpRequestProcessor {
@@ -81,7 +81,7 @@ public class SendAsyncMessageProcessor implements 
HttpRequestProcessor {
                 EventMeshConstants.PROTOCOL_HTTP,
                 RemotingHelper.parseChannelRemoteAddr(ctx.channel()), 
IPUtil.getLocalAddress());
 
-//        SendMessageRequestHeader sendMessageRequestHeader = 
(SendMessageRequestHeader) asyncContext.getRequest().getHeader();
+        SendMessageRequestHeader sendMessageRequestHeader = 
(SendMessageRequestHeader) asyncContext.getRequest().getHeader();
 //        SendMessageRequestBody sendMessageRequestBody = 
(SendMessageRequestBody) asyncContext.getRequest().getBody();
 
         SendMessageResponseHeader sendMessageResponseHeader =
@@ -89,11 +89,13 @@ public class SendAsyncMessageProcessor implements 
HttpRequestProcessor {
                         IPUtil.getLocalAddress(), 
eventMeshHTTPServer.getEventMeshHttpConfiguration().eventMeshEnv,
                         
eventMeshHTTPServer.getEventMeshHttpConfiguration().eventMeshIDC);
 
-        ProtocolAdaptor httpCommandProtocolAdaptor = 
ProtocolPluginFactory.getProtocolAdaptor("cloudevents");
-        CloudEvent event = 
httpCommandProtocolAdaptor.toCloudEventV1(asyncContext.getRequest());
+        String protocolType = sendMessageRequestHeader.getProtocolType();
+        ProtocolAdaptor httpCommandProtocolAdaptor = 
ProtocolPluginFactory.getProtocolAdaptor(protocolType);
+        CloudEvent event = 
httpCommandProtocolAdaptor.toCloudEvent(asyncContext.getRequest());
 
         //validate event
-        if (StringUtils.isBlank(event.getId())
+        if (event != null
+                || StringUtils.isBlank(event.getId())
                 || event.getSource() != null
                 || event.getSpecVersion() != null
                 || StringUtils.isBlank(event.getType())
@@ -185,7 +187,7 @@ public class SendAsyncMessageProcessor implements 
HttpRequestProcessor {
         String ttl = String.valueOf(EventMeshConstants.DEFAULT_MSG_TTL_MILLS);
         if 
(StringUtils.isBlank(event.getExtension(SendMessageRequestBody.TTL).toString())
                 && 
!StringUtils.isNumeric(event.getExtension(SendMessageRequestBody.TTL).toString()))
 {
-            event = new 
CloudEventBuilder(event).withExtension(SendMessageRequestBody.TTL, ttl).build();
+            event = 
CloudEventBuilder.from(event).withExtension(SendMessageRequestBody.TTL, 
ttl).build();
         }
 
         try {
@@ -202,7 +204,7 @@ public class SendAsyncMessageProcessor implements 
HttpRequestProcessor {
 //            omsMsg.putUserProperties(Constants.PROPERTY_MESSAGE_TIMEOUT, 
ttl);
 //            // bizNo
 //            
omsMsg.putSystemProperties(Constants.PROPERTY_MESSAGE_SEARCH_KEYS, 
sendMessageRequestBody.getBizSeqNo());
-            event = new CloudEventBuilder(event)
+            event = CloudEventBuilder.from(event)
                     .withExtension("msgType", "persistent")
                     
.withExtension(EventMeshConstants.REQ_C2EVENTMESH_TIMESTAMP, 
String.valueOf(System.currentTimeMillis()))
                     
.withExtension(EventMeshConstants.REQ_EVENTMESH2MQ_TIMESTAMP, 
String.valueOf(System.currentTimeMillis()))
@@ -248,7 +250,7 @@ public class SendAsyncMessageProcessor implements 
HttpRequestProcessor {
 
 
         try {
-            event = new CloudEventBuilder(sendMessageContext.getEvent())
+            event = CloudEventBuilder.from(sendMessageContext.getEvent())
                     
.withExtension(EventMeshConstants.REQ_EVENTMESH2MQ_TIMESTAMP, 
String.valueOf(System.currentTimeMillis()))
                     .build();
             sendMessageContext.setEvent(event);
diff --git 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/SendSyncMessageProcessor.java
 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/SendSyncMessageProcessor.java
index 041429f..699325f 100644
--- 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/SendSyncMessageProcessor.java
+++ 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/SendSyncMessageProcessor.java
@@ -18,7 +18,7 @@
 package org.apache.eventmesh.runtime.core.protocol.http.processor;
 
 import io.cloudevents.CloudEvent;
-import io.cloudevents.core.v1.CloudEventBuilder;
+import io.cloudevents.core.builder.CloudEventBuilder;
 import org.apache.eventmesh.api.RRCallback;
 import org.apache.eventmesh.api.RequestReplyCallback;
 import org.apache.eventmesh.common.Constants;
@@ -84,8 +84,11 @@ public class SendSyncMessageProcessor implements 
HttpRequestProcessor {
             EventMeshConstants.PROTOCOL_HTTP,
             RemotingHelper.parseChannelRemoteAddr(ctx.channel()), 
IPUtil.getLocalAddress());
 
-        ProtocolAdaptor httpCommandProtocolAdaptor = 
ProtocolPluginFactory.getProtocolAdaptor("cloudevents");
-        CloudEvent event = 
httpCommandProtocolAdaptor.toCloudEventV1(asyncContext.getRequest());
+        SendMessageRequestHeader sendMessageRequestHeader = 
(SendMessageRequestHeader) asyncContext.getRequest().getHeader();
+
+        String protocolType = sendMessageRequestHeader.getProtocolType();
+        ProtocolAdaptor httpCommandProtocolAdaptor = 
ProtocolPluginFactory.getProtocolAdaptor(protocolType);
+        CloudEvent event = 
httpCommandProtocolAdaptor.toCloudEvent(asyncContext.getRequest());
 
         SendMessageResponseHeader sendMessageResponseHeader =
             SendMessageResponseHeader
@@ -96,7 +99,8 @@ public class SendSyncMessageProcessor implements 
HttpRequestProcessor {
                     
eventMeshHTTPServer.getEventMeshHttpConfiguration().eventMeshIDC);
 
         //validate event
-        if (StringUtils.isBlank(event.getId())
+        if (event != null
+                || StringUtils.isBlank(event.getId())
                 || event.getSource() != null
                 || event.getSpecVersion() != null
                 || StringUtils.isBlank(event.getType())
@@ -195,7 +199,7 @@ public class SendSyncMessageProcessor implements 
HttpRequestProcessor {
         }
 
         try {
-            event = new CloudEventBuilder(event)
+            event = CloudEventBuilder.from(event)
                     .withExtension("msgType", "persistent")
                     
.withExtension(EventMeshConstants.REQ_C2EVENTMESH_TIMESTAMP, 
String.valueOf(System.currentTimeMillis()))
                     
.withExtension(EventMeshConstants.REQ_EVENTMESH2MQ_TIMESTAMP, 
String.valueOf(System.currentTimeMillis()))
@@ -250,7 +254,7 @@ public class SendSyncMessageProcessor implements 
HttpRequestProcessor {
                         topic, bizNo, uniqueId);
 
                     try {
-                        event = new CloudEventBuilder(event)
+                        event = CloudEventBuilder.from(event)
                                 
.withExtension(EventMeshConstants.RSP_EVENTMESH2C_TIMESTAMP,
                                 String.valueOf(System.currentTimeMillis()))
                                 
.withExtension(EventMeshConstants.RSP_MQ2EVENTMESH_TIMESTAMP,
diff --git 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/SubscribeProcessor.java
 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/SubscribeProcessor.java
index 4be4fb2..d91b462 100644
--- 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/SubscribeProcessor.java
+++ 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/SubscribeProcessor.java
@@ -192,16 +192,25 @@ public class SubscribeProcessor implements 
HttpRequestProcessor {
                     for (String key : map.keySet()) {
                         if (StringUtils.equals(subTopic.getTopic(), key)) {
                             ConsumerGroupTopicConf latestTopicConf = new 
ConsumerGroupTopicConf();
-                            ConsumerGroupTopicConf currentTopicConf = 
map.get(key);
                             latestTopicConf.setConsumerGroup(consumerGroup);
                             latestTopicConf.setTopic(subTopic.getTopic());
                             latestTopicConf.setSubscriptionItem(subTopic);
                             latestTopicConf.setUrls(new 
HashSet<>(Arrays.asList(url)));
-                            
latestTopicConf.getUrls().addAll(currentTopicConf.getUrls());
 
+                            ConsumerGroupTopicConf currentTopicConf = 
map.get(key);
+                            
latestTopicConf.getUrls().addAll(currentTopicConf.getUrls());
                             latestTopicConf.setIdcUrls(idcUrls);
 
                             map.put(key, latestTopicConf);
+                        } else {
+                            //If there are multiple topics, append it
+                            ConsumerGroupTopicConf newTopicConf = new 
ConsumerGroupTopicConf();
+                            newTopicConf.setConsumerGroup(consumerGroup);
+                            newTopicConf.setTopic(subTopic.getTopic());
+                            newTopicConf.setSubscriptionItem(subTopic);
+                            newTopicConf.setUrls(new 
HashSet<>(Arrays.asList(url)));
+                            newTopicConf.setIdcUrls(idcUrls);
+                            map.put(subTopic.getTopic(), newTopicConf);
                         }
                     }
                 }
diff --git 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/push/AsyncHTTPPushRequest.java
 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/push/AsyncHTTPPushRequest.java
index 4259c94..1559e69 100644
--- 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/push/AsyncHTTPPushRequest.java
+++ 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/push/AsyncHTTPPushRequest.java
@@ -18,7 +18,7 @@
 package org.apache.eventmesh.runtime.core.protocol.http.push;
 
 import io.cloudevents.CloudEvent;
-import io.cloudevents.core.v1.CloudEventBuilder;
+import io.cloudevents.core.builder.CloudEventBuilder;
 import org.apache.eventmesh.common.Constants;
 import org.apache.eventmesh.common.IPUtil;
 import org.apache.eventmesh.common.RandomStringUtil;
@@ -107,7 +107,7 @@ public class AsyncHTTPPushRequest extends 
AbstractHTTPPushRequest {
         builder.addHeader(ProtocolKey.EventMeshInstanceKey.EVENTMESHIDC,
             
handleMsgContext.getEventMeshHTTPServer().getEventMeshHttpConfiguration().eventMeshIDC);
 
-        CloudEvent event = new CloudEventBuilder(handleMsgContext.getEvent())
+        CloudEvent event = CloudEventBuilder.from(handleMsgContext.getEvent())
                 .withExtension(EventMeshConstants.REQ_EVENTMESH2C_TIMESTAMP,
                         String.valueOf(System.currentTimeMillis()))
                 .build();
diff --git 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/group/ClientGroupWrapper.java
 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/group/ClientGroupWrapper.java
index 0768bb9..e789374 100644
--- 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/group/ClientGroupWrapper.java
+++ 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/group/ClientGroupWrapper.java
@@ -18,7 +18,7 @@
 package org.apache.eventmesh.runtime.core.protocol.tcp.client.group;
 
 import io.cloudevents.CloudEvent;
-import io.cloudevents.core.v1.CloudEventBuilder;
+import io.cloudevents.core.builder.CloudEventBuilder;
 import org.apache.eventmesh.api.*;
 import org.apache.eventmesh.api.EventListener;
 import org.apache.eventmesh.api.exception.OnExceptionContext;
@@ -445,7 +445,7 @@ public class ClientGroupWrapper {
                 public void consume(CloudEvent event, AsyncConsumeContext 
context) {
 
                     
eventMeshTcpMonitor.getMq2EventMeshMsgNum().incrementAndGet();
-                    event = new CloudEventBuilder(event)
+                    event = CloudEventBuilder.from(event)
                             
.withExtension(EventMeshConstants.REQ_MQ2EVENTMESH_TIMESTAMP,
                             String.valueOf(System.currentTimeMillis()))
                             
.withExtension(EventMeshConstants.REQ_RECEIVE_EVENTMESH_IP,
@@ -506,7 +506,7 @@ public class ClientGroupWrapper {
                 @Override
                 public void consume(CloudEvent event, AsyncConsumeContext 
context) {
                     
eventMeshTcpMonitor.getMq2EventMeshMsgNum().incrementAndGet();
-                    event = new CloudEventBuilder(event)
+                    event = CloudEventBuilder.from(event)
                             
.withExtension(EventMeshConstants.REQ_MQ2EVENTMESH_TIMESTAMP,
                                     String.valueOf(System.currentTimeMillis()))
                             
.withExtension(EventMeshConstants.REQ_RECEIVE_EVENTMESH_IP,
@@ -547,7 +547,7 @@ public class ClientGroupWrapper {
                                     consumerGroup, topic, bizSeqNo);
                             } else {
                                 sendBackTimes++;
-                                event = new CloudEventBuilder(event)
+                                event = CloudEventBuilder.from(event)
                                         
.withExtension(EventMeshConstants.EVENTMESH_SEND_BACK_TIMES,
                                                 sendBackTimes.toString())
                                         
.withExtension(EventMeshConstants.EVENTMESH_SEND_BACK_IP,
diff --git 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/session/push/SessionPusher.java
 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/session/push/SessionPusher.java
index 3e86f21..78de6fc 100644
--- 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/session/push/SessionPusher.java
+++ 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/session/push/SessionPusher.java
@@ -17,7 +17,7 @@
 
 package org.apache.eventmesh.runtime.core.protocol.tcp.client.session.push;
 
-import io.cloudevents.core.v1.CloudEventBuilder;
+import io.cloudevents.core.builder.CloudEventBuilder;
 import io.netty.channel.ChannelFuture;
 import io.netty.channel.ChannelFutureListener;
 import org.apache.commons.collections4.CollectionUtils;
@@ -80,7 +80,7 @@ public class SessionPusher {
 
         Package pkg = new Package();
 
-        downStreamMsgContext.event = new 
CloudEventBuilder(downStreamMsgContext.event)
+        downStreamMsgContext.event = 
CloudEventBuilder.from(downStreamMsgContext.event)
                 .withExtension(EventMeshConstants.REQ_EVENTMESH2C_TIMESTAMP, 
String.valueOf(System.currentTimeMillis()))
                 .build();
 //        
downStreamMsgContext.event.getSystemProperties().put(EventMeshConstants.REQ_EVENTMESH2C_TIMESTAMP,
 String.valueOf(System.currentTimeMillis()));
diff --git 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/session/send/SessionSender.java
 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/session/send/SessionSender.java
index c77ed6a..4fbea05 100644
--- 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/session/send/SessionSender.java
+++ 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/session/send/SessionSender.java
@@ -17,6 +17,7 @@
 
 package org.apache.eventmesh.runtime.core.protocol.tcp.client.session.send;
 
+import io.cloudevents.core.builder.CloudEventBuilder;
 import org.apache.commons.lang3.StringUtils;
 import org.apache.commons.lang3.time.DateFormatUtils;
 import org.apache.eventmesh.api.RRCallback;
@@ -39,7 +40,6 @@ import java.util.concurrent.TimeUnit;
 import java.util.concurrent.atomic.AtomicLong;
 
 import io.cloudevents.CloudEvent;
-import io.cloudevents.core.v1.CloudEventBuilder;
 
 public class SessionSender {
 
@@ -96,7 +96,7 @@ public class SessionSender {
                     if (!StringUtils.isEmpty(cluster)) {
                         String replyTopic = EventMeshConstants.RR_REPLY_TOPIC;
                         replyTopic = cluster + "-" + replyTopic;
-                        event = new 
CloudEventBuilder(event).withSubject(replyTopic).build();
+                        event = 
CloudEventBuilder.from(event).withSubject(replyTopic).build();
 //                        
msg.getSystemProperties().put(Constants.PROPERTY_MESSAGE_DESTINATION, 
replyTopic);
 //                        event(replyTopic);
                     }
@@ -142,7 +142,7 @@ public class SessionSender {
 //                    msg.putUserProperty(EventMeshConstants.STORE_TIMESTAMP, 
String.valueOf(((MessageExt) msg)
 //                            .getStoreTimestamp()));
 //                }
-                event = new CloudEventBuilder(event)
+                event = CloudEventBuilder.from(event)
                         
.withExtension(EventMeshConstants.RSP_MQ2EVENTMESH_TIMESTAMP, 
String.valueOf(System.currentTimeMillis()))
                         
.withExtension(EventMeshConstants.RSP_RECEIVE_EVENTMESH_IP, 
session.getEventMeshTCPConfiguration().eventMeshServerIp)
                         .build();
@@ -157,7 +157,7 @@ public class SessionSender {
                 }
                 Package pkg = new Package();
                 pkg.setHeader(new Header(cmd, OPStatus.SUCCESS.getCode(), 
null, seq));
-                event = new CloudEventBuilder(event)
+                event = CloudEventBuilder.from(event)
                         
.withExtension(EventMeshConstants.RSP_EVENTMESH2C_TIMESTAMP, 
String.valueOf(System.currentTimeMillis()))
                         .build();
 //                
msg.getSystemProperties().put(EventMeshConstants.RSP_EVENTMESH2C_TIMESTAMP, 
String.valueOf(System.currentTimeMillis()));
diff --git 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/task/MessageTransferTask.java
 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/task/MessageTransferTask.java
index 43f5b36..6a1bd59 100644
--- 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/task/MessageTransferTask.java
+++ 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/task/MessageTransferTask.java
@@ -18,7 +18,7 @@
 package org.apache.eventmesh.runtime.core.protocol.tcp.client.task;
 
 import io.cloudevents.CloudEvent;
-import io.cloudevents.core.v1.CloudEventBuilder;
+import io.cloudevents.core.builder.CloudEventBuilder;
 import io.netty.channel.ChannelFuture;
 import io.netty.channel.ChannelFutureListener;
 import io.netty.channel.ChannelHandlerContext;
@@ -74,7 +74,7 @@ public class MessageTransferTask extends AbstractTask {
         EventMeshTcpSendResult sendStatus;
         CloudEvent event = null;
         try {
-            event = protocolAdaptor.toCloudEventV1(pkg);
+            event = protocolAdaptor.toCloudEvent(pkg);
             if (event == null) {
                 throw new Exception("event is null");
             }
@@ -123,13 +123,13 @@ public class MessageTransferTask extends AbstractTask {
 
     private CloudEvent addTimestamp(CloudEvent event, Command cmd, long 
sendTime) {
         if (cmd.equals(RESPONSE_TO_SERVER)) {
-            event = new CloudEventBuilder(event)
+            event = CloudEventBuilder.from(event)
                     
.withExtension(EventMeshConstants.RSP_C2EVENTMESH_TIMESTAMP, 
String.valueOf(startTime))
                     
.withExtension(EventMeshConstants.RSP_EVENTMESH2MQ_TIMESTAMP, 
String.valueOf(sendTime))
                     .withExtension(EventMeshConstants.RSP_SEND_EVENTMESH_IP, 
eventMeshTCPServer.getEventMeshTCPConfiguration().eventMeshServerIp)
                     .build();
         } else {
-            event = new CloudEventBuilder(event)
+            event = CloudEventBuilder.from(event)
                     
.withExtension(EventMeshConstants.REQ_C2EVENTMESH_TIMESTAMP, 
String.valueOf(startTime))
                     
.withExtension(EventMeshConstants.REQ_EVENTMESH2MQ_TIMESTAMP, 
String.valueOf(sendTime))
                     .withExtension(EventMeshConstants.REQ_SEND_EVENTMESH_IP, 
eventMeshTCPServer.getEventMeshTCPConfiguration().eventMeshServerIp)
diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-rocketmq/gradle.properties 
b/eventmesh-security-plugin/eventmesh-security-acl/gradle.properties
similarity index 95%
copy from 
eventmesh-connector-plugin/eventmesh-connector-rocketmq/gradle.properties
copy to eventmesh-security-plugin/eventmesh-security-acl/gradle.properties
index 3d49f4c..719ba4f 100644
--- a/eventmesh-connector-plugin/eventmesh-connector-rocketmq/gradle.properties
+++ b/eventmesh-security-plugin/eventmesh-security-acl/gradle.properties
@@ -14,4 +14,5 @@
 # limitations under the License.
 #
 
-rocketmq_version=4.7.1
\ No newline at end of file
+pluginType=security
+pluginName=acl
\ No newline at end of file

---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to