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 @@

## 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.

-**EventMesh Ecosystem:**
-
-
-
**EventMesh Architecture:**
-
-
-**EventMesh Cloud Native:**
-
-
-
+
**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 |
| :---------------------------------------: |
:----------------------------------------------------: |
-|  |
 |
+|  |
 |
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]