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 fcd1218 update Eventmeshmessage plugin (#599)
fcd1218 is described below
commit fcd1218fc231d76e375e407c0d2a34f8c115e46b
Author: mike_xwm <[email protected]>
AuthorDate: Fri Nov 19 17:29:27 2021 +0800
update Eventmeshmessage plugin (#599)
* [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
* [Feature #562] Implement EventMeshMessage protocol adaptor
* supplement apache header
* 1.update package name
2.support build.gradle and gradle.properties
3.support ProtocolTransportObject
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]>
---
.../build.gradle} | 10 ++++----
.../gradle.properties} | 4 +++-
.../meshmessage/MeshMessageProtocolAdaptor.java} | 27 ++++++++++++----------
.../meshmessage/MeshMessageProtocolConstant.java} | 4 ++--
.../http/SendMessageBatchProtocolResolver.java | 2 +-
.../http/SendMessageBatchV2ProtocolResolver.java | 6 +----
.../http/SendMessageRequestProtocolResolver.java | 5 +---
.../resolver/tcp/TcpMessageProtocolResolver.java | 6 ++---
...g.apache.eventmesh.protocol.api.ProtocolAdaptor | 2 +-
settings.gradle | 2 +-
10 files changed, 34 insertions(+), 34 deletions(-)
diff --git
a/eventmesh-protocol-plugin/eventmesh-protocol-eventmeshmessage/src/main/java/org/apache/eventmesh/protocol/eventmeshmessage/EventMeshMessageProtocolConstant.java
b/eventmesh-protocol-plugin/eventmesh-protocol-meshmessage/build.gradle
similarity index 71%
copy from
eventmesh-protocol-plugin/eventmesh-protocol-eventmeshmessage/src/main/java/org/apache/eventmesh/protocol/eventmeshmessage/EventMeshMessageProtocolConstant.java
copy to eventmesh-protocol-plugin/eventmesh-protocol-meshmessage/build.gradle
index f1c744f..4a82353 100644
---
a/eventmesh-protocol-plugin/eventmesh-protocol-eventmeshmessage/src/main/java/org/apache/eventmesh/protocol/eventmeshmessage/EventMeshMessageProtocolConstant.java
+++ b/eventmesh-protocol-plugin/eventmesh-protocol-meshmessage/build.gradle
@@ -15,9 +15,11 @@
* limitations under the License.
*/
-package org.apache.eventmesh.protocol.eventmeshmessage;
+dependencies {
+ compileOnly project(":eventmesh-protocol-plugin:eventmesh-protocol-api")
+ implementation "io.cloudevents:cloudevents-core"
-public enum EventMeshMessageProtocolConstant {
- ;
- public static final String PROTOCOL_NAME = "eventmeshmessage";
+ testImplementation
project(":eventmesh-protocol-plugin:eventmesh-protocol-api")
+ testImplementation "io.cloudevents:cloudevents-core"
+ testImplementation "junit:junit"
}
diff --git
a/eventmesh-protocol-plugin/eventmesh-protocol-eventmeshmessage/src/main/resources/META-INF.eventmesh/org.apache.eventmesh.protocol.api.ProtocolAdaptor
b/eventmesh-protocol-plugin/eventmesh-protocol-meshmessage/gradle.properties
similarity index 89%
copy from
eventmesh-protocol-plugin/eventmesh-protocol-eventmeshmessage/src/main/resources/META-INF.eventmesh/org.apache.eventmesh.protocol.api.ProtocolAdaptor
copy to
eventmesh-protocol-plugin/eventmesh-protocol-meshmessage/gradle.properties
index 9be39ed..07476be 100644
---
a/eventmesh-protocol-plugin/eventmesh-protocol-eventmeshmessage/src/main/resources/META-INF.eventmesh/org.apache.eventmesh.protocol.api.ProtocolAdaptor
+++ b/eventmesh-protocol-plugin/eventmesh-protocol-meshmessage/gradle.properties
@@ -12,5 +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.
+#
-eventmeshmessage=org.apache.eventmesh.protocol.eventmeshmessage.EventMeshMessageProtocolAdaptor
\ No newline at end of file
+pluginType=protocol
+pluginName=eventmeshmessage
\ No newline at end of file
diff --git
a/eventmesh-protocol-plugin/eventmesh-protocol-eventmeshmessage/src/main/java/org/apache/eventmesh/protocol/eventmeshmessage/EventMeshMessageProtocolAdaptor.java
b/eventmesh-protocol-plugin/eventmesh-protocol-meshmessage/src/main/java/org/apache/eventmesh/protocol/meshmessage/MeshMessageProtocolAdaptor.java
similarity index 76%
rename from
eventmesh-protocol-plugin/eventmesh-protocol-eventmeshmessage/src/main/java/org/apache/eventmesh/protocol/eventmeshmessage/EventMeshMessageProtocolAdaptor.java
rename to
eventmesh-protocol-plugin/eventmesh-protocol-meshmessage/src/main/java/org/apache/eventmesh/protocol/meshmessage/MeshMessageProtocolAdaptor.java
index 4b8eeef..3781147 100644
---
a/eventmesh-protocol-plugin/eventmesh-protocol-eventmeshmessage/src/main/java/org/apache/eventmesh/protocol/eventmeshmessage/EventMeshMessageProtocolAdaptor.java
+++
b/eventmesh-protocol-plugin/eventmesh-protocol-meshmessage/src/main/java/org/apache/eventmesh/protocol/meshmessage/MeshMessageProtocolAdaptor.java
@@ -15,31 +15,32 @@
* limitations under the License.
*/
-package org.apache.eventmesh.protocol.eventmeshmessage;
+package org.apache.eventmesh.protocol.meshmessage;
import io.cloudevents.CloudEvent;
import org.apache.commons.lang3.StringUtils;
import org.apache.eventmesh.common.Constants;
+import org.apache.eventmesh.common.ProtocolTransportObject;
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.EventMeshMessage;
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.eventmeshmessage.resolver.http.SendMessageBatchProtocolResolver;
-import
org.apache.eventmesh.protocol.eventmeshmessage.resolver.http.SendMessageBatchV2ProtocolResolver;
-import
org.apache.eventmesh.protocol.eventmeshmessage.resolver.http.SendMessageRequestProtocolResolver;
-import
org.apache.eventmesh.protocol.eventmeshmessage.resolver.tcp.TcpMessageProtocolResolver;
+import
org.apache.eventmesh.protocol.meshmessage.resolver.http.SendMessageBatchProtocolResolver;
+import
org.apache.eventmesh.protocol.meshmessage.resolver.http.SendMessageBatchV2ProtocolResolver;
+import
org.apache.eventmesh.protocol.meshmessage.resolver.http.SendMessageRequestProtocolResolver;
+import
org.apache.eventmesh.protocol.meshmessage.resolver.tcp.TcpMessageProtocolResolver;
import java.nio.charset.StandardCharsets;
import java.util.List;
-public class EventMeshMessageProtocolAdaptor<T> implements ProtocolAdaptor<T> {
+public class MeshMessageProtocolAdaptor<T extends ProtocolTransportObject>
+ implements ProtocolAdaptor<ProtocolTransportObject> {
@Override
- public CloudEvent toCloudEvent(T protocol) throws ProtocolHandleException {
+ public CloudEvent toCloudEvent(ProtocolTransportObject protocol) throws
ProtocolHandleException {
if (protocol instanceof Package) {
Header header = ((Package) protocol).getHeader();
Object body = ((Package) protocol).getBody();
@@ -78,16 +79,18 @@ public class EventMeshMessageProtocolAdaptor<T> implements
ProtocolAdaptor<T> {
}
@Override
- public List<CloudEvent> toBatchCloudEvent(T protocol) throws
ProtocolHandleException {
+ public List<CloudEvent> toBatchCloudEvent(ProtocolTransportObject
protocol) throws ProtocolHandleException {
return null;
}
@Override
- public Object fromCloudEvent(CloudEvent cloudEvent) throws
ProtocolHandleException {
+ public ProtocolTransportObject 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);
+ // todo: return command, set cloudEvent.getData() to content?
+ return null;
+// return new String(cloudEvent.getData().toBytes(),
StandardCharsets.UTF_8);
} else if (StringUtils.equals("tcp", protocolDesc)) {
return
TcpMessageProtocolResolver.buildEventMeshMessage(cloudEvent);
} else {
@@ -97,6 +100,6 @@ public class EventMeshMessageProtocolAdaptor<T> implements
ProtocolAdaptor<T> {
@Override
public String getProtocolType() {
- return EventMeshMessageProtocolConstant.PROTOCOL_NAME;
+ return MeshMessageProtocolConstant.PROTOCOL_NAME;
}
}
diff --git
a/eventmesh-protocol-plugin/eventmesh-protocol-eventmeshmessage/src/main/java/org/apache/eventmesh/protocol/eventmeshmessage/EventMeshMessageProtocolConstant.java
b/eventmesh-protocol-plugin/eventmesh-protocol-meshmessage/src/main/java/org/apache/eventmesh/protocol/meshmessage/MeshMessageProtocolConstant.java
similarity index 89%
rename from
eventmesh-protocol-plugin/eventmesh-protocol-eventmeshmessage/src/main/java/org/apache/eventmesh/protocol/eventmeshmessage/EventMeshMessageProtocolConstant.java
rename to
eventmesh-protocol-plugin/eventmesh-protocol-meshmessage/src/main/java/org/apache/eventmesh/protocol/meshmessage/MeshMessageProtocolConstant.java
index f1c744f..0a667f6 100644
---
a/eventmesh-protocol-plugin/eventmesh-protocol-eventmeshmessage/src/main/java/org/apache/eventmesh/protocol/eventmeshmessage/EventMeshMessageProtocolConstant.java
+++
b/eventmesh-protocol-plugin/eventmesh-protocol-meshmessage/src/main/java/org/apache/eventmesh/protocol/meshmessage/MeshMessageProtocolConstant.java
@@ -15,9 +15,9 @@
* limitations under the License.
*/
-package org.apache.eventmesh.protocol.eventmeshmessage;
+package org.apache.eventmesh.protocol.meshmessage;
-public enum EventMeshMessageProtocolConstant {
+public enum MeshMessageProtocolConstant {
;
public static final String PROTOCOL_NAME = "eventmeshmessage";
}
diff --git
a/eventmesh-protocol-plugin/eventmesh-protocol-eventmeshmessage/src/main/java/org/apache/eventmesh/protocol/eventmeshmessage/resolver/http/SendMessageBatchProtocolResolver.java
b/eventmesh-protocol-plugin/eventmesh-protocol-meshmessage/src/main/java/org/apache/eventmesh/protocol/meshmessage/resolver/http/SendMessageBatchProtocolResolver.java
similarity index 94%
rename from
eventmesh-protocol-plugin/eventmesh-protocol-eventmeshmessage/src/main/java/org/apache/eventmesh/protocol/eventmeshmessage/resolver/http/SendMessageBatchProtocolResolver.java
rename to
eventmesh-protocol-plugin/eventmesh-protocol-meshmessage/src/main/java/org/apache/eventmesh/protocol/meshmessage/resolver/http/SendMessageBatchProtocolResolver.java
index 9c82038..cd8f817 100644
---
a/eventmesh-protocol-plugin/eventmesh-protocol-eventmeshmessage/src/main/java/org/apache/eventmesh/protocol/eventmeshmessage/resolver/http/SendMessageBatchProtocolResolver.java
+++
b/eventmesh-protocol-plugin/eventmesh-protocol-meshmessage/src/main/java/org/apache/eventmesh/protocol/meshmessage/resolver/http/SendMessageBatchProtocolResolver.java
@@ -15,7 +15,7 @@
* limitations under the License.
*/
-package org.apache.eventmesh.protocol.eventmeshmessage.resolver.http;
+package org.apache.eventmesh.protocol.meshmessage.resolver.http;
import io.cloudevents.CloudEvent;
import org.apache.eventmesh.common.protocol.http.body.Body;
diff --git
a/eventmesh-protocol-plugin/eventmesh-protocol-eventmeshmessage/src/main/java/org/apache/eventmesh/protocol/eventmeshmessage/resolver/http/SendMessageBatchV2ProtocolResolver.java
b/eventmesh-protocol-plugin/eventmesh-protocol-meshmessage/src/main/java/org/apache/eventmesh/protocol/meshmessage/resolver/http/SendMessageBatchV2ProtocolResolver.java
similarity index 96%
rename from
eventmesh-protocol-plugin/eventmesh-protocol-eventmeshmessage/src/main/java/org/apache/eventmesh/protocol/eventmeshmessage/resolver/http/SendMessageBatchV2ProtocolResolver.java
rename to
eventmesh-protocol-plugin/eventmesh-protocol-meshmessage/src/main/java/org/apache/eventmesh/protocol/meshmessage/resolver/http/SendMessageBatchV2ProtocolResolver.java
index adc1e7e..6e2ad28 100644
---
a/eventmesh-protocol-plugin/eventmesh-protocol-eventmeshmessage/src/main/java/org/apache/eventmesh/protocol/eventmeshmessage/resolver/http/SendMessageBatchV2ProtocolResolver.java
+++
b/eventmesh-protocol-plugin/eventmesh-protocol-meshmessage/src/main/java/org/apache/eventmesh/protocol/meshmessage/resolver/http/SendMessageBatchV2ProtocolResolver.java
@@ -15,13 +15,11 @@
* limitations under the License.
*/
-package org.apache.eventmesh.protocol.eventmeshmessage.resolver.http;
+package org.apache.eventmesh.protocol.meshmessage.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;
@@ -29,8 +27,6 @@ import
org.apache.eventmesh.common.protocol.http.common.ProtocolKey;
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.protocol.tcp.EventMeshMessage;
-import org.apache.eventmesh.common.utils.JsonUtils;
import org.apache.eventmesh.protocol.api.exception.ProtocolHandleException;
import java.nio.charset.StandardCharsets;
diff --git
a/eventmesh-protocol-plugin/eventmesh-protocol-eventmeshmessage/src/main/java/org/apache/eventmesh/protocol/eventmeshmessage/resolver/http/SendMessageRequestProtocolResolver.java
b/eventmesh-protocol-plugin/eventmesh-protocol-meshmessage/src/main/java/org/apache/eventmesh/protocol/meshmessage/resolver/http/SendMessageRequestProtocolResolver.java
similarity index 97%
rename from
eventmesh-protocol-plugin/eventmesh-protocol-eventmeshmessage/src/main/java/org/apache/eventmesh/protocol/eventmeshmessage/resolver/http/SendMessageRequestProtocolResolver.java
rename to
eventmesh-protocol-plugin/eventmesh-protocol-meshmessage/src/main/java/org/apache/eventmesh/protocol/meshmessage/resolver/http/SendMessageRequestProtocolResolver.java
index 40e2e9f..9abbdaf 100644
---
a/eventmesh-protocol-plugin/eventmesh-protocol-eventmeshmessage/src/main/java/org/apache/eventmesh/protocol/eventmeshmessage/resolver/http/SendMessageRequestProtocolResolver.java
+++
b/eventmesh-protocol-plugin/eventmesh-protocol-meshmessage/src/main/java/org/apache/eventmesh/protocol/meshmessage/resolver/http/SendMessageRequestProtocolResolver.java
@@ -15,13 +15,11 @@
* limitations under the License.
*/
-package org.apache.eventmesh.protocol.eventmeshmessage.resolver.http;
+package org.apache.eventmesh.protocol.meshmessage.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;
@@ -30,7 +28,6 @@ import
org.apache.eventmesh.common.protocol.http.common.ProtocolKey;
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;
import java.nio.charset.StandardCharsets;
diff --git
a/eventmesh-protocol-plugin/eventmesh-protocol-eventmeshmessage/src/main/java/org/apache/eventmesh/protocol/eventmeshmessage/resolver/tcp/TcpMessageProtocolResolver.java
b/eventmesh-protocol-plugin/eventmesh-protocol-meshmessage/src/main/java/org/apache/eventmesh/protocol/meshmessage/resolver/tcp/TcpMessageProtocolResolver.java
similarity index 94%
rename from
eventmesh-protocol-plugin/eventmesh-protocol-eventmeshmessage/src/main/java/org/apache/eventmesh/protocol/eventmeshmessage/resolver/tcp/TcpMessageProtocolResolver.java
rename to
eventmesh-protocol-plugin/eventmesh-protocol-meshmessage/src/main/java/org/apache/eventmesh/protocol/meshmessage/resolver/tcp/TcpMessageProtocolResolver.java
index 38faf38..6235759 100644
---
a/eventmesh-protocol-plugin/eventmesh-protocol-eventmeshmessage/src/main/java/org/apache/eventmesh/protocol/eventmeshmessage/resolver/tcp/TcpMessageProtocolResolver.java
+++
b/eventmesh-protocol-plugin/eventmesh-protocol-meshmessage/src/main/java/org/apache/eventmesh/protocol/meshmessage/resolver/tcp/TcpMessageProtocolResolver.java
@@ -15,7 +15,7 @@
* limitations under the License.
*/
-package org.apache.eventmesh.protocol.eventmeshmessage.resolver.tcp;
+package org.apache.eventmesh.protocol.meshmessage.resolver.tcp;
import io.cloudevents.CloudEvent;
import io.cloudevents.SpecVersion;
@@ -26,7 +26,7 @@ import
org.apache.eventmesh.common.protocol.tcp.EventMeshMessage;
import org.apache.eventmesh.common.protocol.tcp.Header;
import org.apache.eventmesh.common.protocol.tcp.Package;
import org.apache.eventmesh.protocol.api.exception.ProtocolHandleException;
-import
org.apache.eventmesh.protocol.eventmeshmessage.EventMeshMessageProtocolConstant;
+import org.apache.eventmesh.protocol.meshmessage.MeshMessageProtocolConstant;
import java.nio.charset.StandardCharsets;
import java.util.HashMap;
@@ -50,7 +50,7 @@ public class TcpMessageProtocolResolver {
protocolType, protocolVersion, protocolDesc));
}
- if
(!StringUtils.equals(EventMeshMessageProtocolConstant.PROTOCOL_NAME,
protocolType)) {
+ if (!StringUtils.equals(MeshMessageProtocolConstant.PROTOCOL_NAME,
protocolType)) {
throw new ProtocolHandleException(String.format("Unsupported
protocolType: %s", protocolType));
}
diff --git
a/eventmesh-protocol-plugin/eventmesh-protocol-eventmeshmessage/src/main/resources/META-INF.eventmesh/org.apache.eventmesh.protocol.api.ProtocolAdaptor
b/eventmesh-protocol-plugin/eventmesh-protocol-meshmessage/src/main/resources/META-INF.eventmesh/org.apache.eventmesh.protocol.api.ProtocolAdaptor
similarity index 89%
rename from
eventmesh-protocol-plugin/eventmesh-protocol-eventmeshmessage/src/main/resources/META-INF.eventmesh/org.apache.eventmesh.protocol.api.ProtocolAdaptor
rename to
eventmesh-protocol-plugin/eventmesh-protocol-meshmessage/src/main/resources/META-INF.eventmesh/org.apache.eventmesh.protocol.api.ProtocolAdaptor
index 9be39ed..45ffc15 100644
---
a/eventmesh-protocol-plugin/eventmesh-protocol-eventmeshmessage/src/main/resources/META-INF.eventmesh/org.apache.eventmesh.protocol.api.ProtocolAdaptor
+++
b/eventmesh-protocol-plugin/eventmesh-protocol-meshmessage/src/main/resources/META-INF.eventmesh/org.apache.eventmesh.protocol.api.ProtocolAdaptor
@@ -13,4 +13,4 @@
# See the License for the specific language governing permissions and
# limitations under the License.
-eventmeshmessage=org.apache.eventmesh.protocol.eventmeshmessage.EventMeshMessageProtocolAdaptor
\ No newline at end of file
+eventmeshmessage=org.apache.eventmesh.protocol.meshmessage.MeshMessageProtocolAdaptor
\ No newline at end of file
diff --git a/settings.gradle b/settings.gradle
index 787313b..0606b2e 100644
--- a/settings.gradle
+++ b/settings.gradle
@@ -37,5 +37,5 @@ include 'eventmesh-protocol-plugin'
include 'eventmesh-protocol-plugin:eventmesh-protocol-api'
include 'eventmesh-protocol-plugin:eventmesh-protocol-openmessage'
include 'eventmesh-protocol-plugin:eventmesh-protocol-cloudevents'
-include 'eventmesh-protocol-plugin:eventmesh-protocol-eventmeshmessage'
+include 'eventmesh-protocol-plugin:eventmesh-protocol-meshmessage'
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]