This is an automated email from the ASF dual-hosted git repository.
mikexue pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/eventmesh.git
The following commit(s) were added to refs/heads/master by this push:
new cb96b920c [ISSUE #3339]Refactor EventMeshGrpcProducer
new dc6481d06 Merge pull request #3340 from mxsm/eventmesh-3338
cb96b920c is described below
commit cb96b920ccf8a4ab06f9a852ee6a8e66df024480
Author: mxsm <[email protected]>
AuthorDate: Sun Mar 5 09:39:04 2023 +0800
[ISSUE #3339]Refactor EventMeshGrpcProducer
---
.../apache/eventmesh/common/EventMeshMessage.java | 6 +-
.../common/enums/EventMeshProtocolType.java | 36 ++-
.../grpc/sub/CloudEventsAsyncSubscribe.java | 6 +-
.../grpc/sub/CloudEventsSubscribeReply.java | 6 +-
.../grpc/sub/EventMeshAsyncSubscribe.java | 7 +-
.../grpc/sub/EventMeshSubscribeBroadcast.java | 7 +-
.../grpc/sub/EventMeshSubscribeReply.java | 7 +-
.../grpc/sub/WorkflowExpressAsyncSubscribe.java | 7 +-
.../grpc/sub/WorkflowOrderAsyncSubscribe.java | 7 +-
.../grpc/sub/WorkflowPaymentAsyncSubscribe.java | 7 +-
.../grpc/consumer/EventMeshGrpcConsumer.java | 6 +-
.../client/grpc/consumer/ReceiveMsgHook.java | 4 +-
.../client/grpc/producer/CloudEventProducer.java | 19 +-
.../grpc/producer/EventMeshGrpcProducer.java | 100 ++-----
...Producer.java => EventMeshMessageProducer.java} | 77 ++---
.../GrpcProducer.java} | 21 +-
.../client/grpc/util/EventMeshClientUtil.java | 309 ++++++++++++---------
.../grpc/consumer/EventMeshGrpcConsumerTest.java | 5 +-
.../grpc/producer/EventMeshGrpcProducerTest.java | 35 ++-
.../producer/EventMeshMessageProducerTest.java | 120 ++++++++
.../client/grpc/util/EventMeshClientUtilTest.java | 26 +-
21 files changed, 465 insertions(+), 353 deletions(-)
diff --git
a/eventmesh-common/src/main/java/org/apache/eventmesh/common/EventMeshMessage.java
b/eventmesh-common/src/main/java/org/apache/eventmesh/common/EventMeshMessage.java
index 8fecc700f..9175ca45d 100644
---
a/eventmesh-common/src/main/java/org/apache/eventmesh/common/EventMeshMessage.java
+++
b/eventmesh-common/src/main/java/org/apache/eventmesh/common/EventMeshMessage.java
@@ -38,15 +38,13 @@ public class EventMeshMessage {
private String content;
- private Map<String, String> prop;
+ @Builder.Default
+ private Map<String, String> prop = new HashMap<>();
@Builder.Default
private final long createTime = System.currentTimeMillis();
public EventMeshMessage addProp(String key, String val) {
- if (prop == null) {
- prop = new HashMap<>();
- }
prop.put(key, val);
return this;
}
diff --git
a/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/grpc/consumer/ReceiveMsgHook.java
b/eventmesh-common/src/main/java/org/apache/eventmesh/common/enums/EventMeshProtocolType.java
similarity index 54%
copy from
eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/grpc/consumer/ReceiveMsgHook.java
copy to
eventmesh-common/src/main/java/org/apache/eventmesh/common/enums/EventMeshProtocolType.java
index 2e99ea2b8..6108f4ce5 100644
---
a/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/grpc/consumer/ReceiveMsgHook.java
+++
b/eventmesh-common/src/main/java/org/apache/eventmesh/common/enums/EventMeshProtocolType.java
@@ -15,23 +15,31 @@
* limitations under the License.
*/
-package org.apache.eventmesh.client.grpc.consumer;
+package org.apache.eventmesh.common.enums;
-import java.util.Optional;
+public enum EventMeshProtocolType {
-/**
- * @param <T>
- */
-public interface ReceiveMsgHook<T> {
+ CLOUD_EVENTS("cloudevents"),
+ EVENT_MESH_MESSAGE("eventmeshmessage"),
+ OPEN_MESSAGE("openmessage");
+
+ private String name;
+
+ EventMeshProtocolType(String name) {
+ this.name = name;
+ }
- /**
- * Handle the received message, return the response message.
- *
- * @param msg
- * @return
- */
- Optional<T> handle(T msg) throws Exception;
+ public String protocolTypeName() {
+ return this.name;
+ }
- String getProtocolType();
+ public static EventMeshProtocolType eventMeshProtocolType(String name) {
+ for (EventMeshProtocolType protocolType :
EventMeshProtocolType.values()) {
+ if (protocolType.protocolTypeName().equals(name)) {
+ return protocolType;
+ }
+ }
+ return null;
+ }
}
diff --git
a/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/CloudEventsAsyncSubscribe.java
b/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/CloudEventsAsyncSubscribe.java
index fba1a74e8..32f9baf8c 100644
---
a/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/CloudEventsAsyncSubscribe.java
+++
b/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/CloudEventsAsyncSubscribe.java
@@ -19,8 +19,8 @@ package org.apache.eventmesh.grpc.sub;
import org.apache.eventmesh.client.grpc.consumer.EventMeshGrpcConsumer;
import org.apache.eventmesh.client.grpc.consumer.ReceiveMsgHook;
-import org.apache.eventmesh.client.tcp.common.EventMeshCommon;
import org.apache.eventmesh.common.ExampleConstants;
+import org.apache.eventmesh.common.enums.EventMeshProtocolType;
import org.apache.eventmesh.common.protocol.SubscriptionItem;
import org.apache.eventmesh.common.protocol.SubscriptionMode;
import org.apache.eventmesh.common.protocol.SubscriptionType;
@@ -68,7 +68,7 @@ public class CloudEventsAsyncSubscribe extends
GrpcAbstractDemo implements Recei
}
@Override
- public String getProtocolType() {
- return EventMeshCommon.CLOUD_EVENTS_PROTOCOL_NAME;
+ public EventMeshProtocolType getProtocolType() {
+ return EventMeshProtocolType.CLOUD_EVENTS;
}
}
diff --git
a/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/CloudEventsSubscribeReply.java
b/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/CloudEventsSubscribeReply.java
index 3a288f029..53e95c68d 100644
---
a/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/CloudEventsSubscribeReply.java
+++
b/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/CloudEventsSubscribeReply.java
@@ -19,8 +19,8 @@ package org.apache.eventmesh.grpc.sub;
import org.apache.eventmesh.client.grpc.consumer.EventMeshGrpcConsumer;
import org.apache.eventmesh.client.grpc.consumer.ReceiveMsgHook;
-import org.apache.eventmesh.client.tcp.common.EventMeshCommon;
import org.apache.eventmesh.common.ExampleConstants;
+import org.apache.eventmesh.common.enums.EventMeshProtocolType;
import org.apache.eventmesh.common.protocol.SubscriptionItem;
import org.apache.eventmesh.common.protocol.SubscriptionMode;
import org.apache.eventmesh.common.protocol.SubscriptionType;
@@ -74,7 +74,7 @@ public class CloudEventsSubscribeReply extends
GrpcAbstractDemo implements Recei
}
@Override
- public String getProtocolType() {
- return EventMeshCommon.CLOUD_EVENTS_PROTOCOL_NAME;
+ public EventMeshProtocolType getProtocolType() {
+ return EventMeshProtocolType.CLOUD_EVENTS;
}
}
diff --git
a/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/EventMeshAsyncSubscribe.java
b/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/EventMeshAsyncSubscribe.java
index 7c6167fc9..b7b7206d3 100644
---
a/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/EventMeshAsyncSubscribe.java
+++
b/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/EventMeshAsyncSubscribe.java
@@ -19,9 +19,9 @@ package org.apache.eventmesh.grpc.sub;
import org.apache.eventmesh.client.grpc.consumer.EventMeshGrpcConsumer;
import org.apache.eventmesh.client.grpc.consumer.ReceiveMsgHook;
-import org.apache.eventmesh.client.tcp.common.EventMeshCommon;
import org.apache.eventmesh.common.EventMeshMessage;
import org.apache.eventmesh.common.ExampleConstants;
+import org.apache.eventmesh.common.enums.EventMeshProtocolType;
import org.apache.eventmesh.common.protocol.SubscriptionItem;
import org.apache.eventmesh.common.protocol.SubscriptionMode;
import org.apache.eventmesh.common.protocol.SubscriptionType;
@@ -33,7 +33,6 @@ import java.util.Collections;
import java.util.Optional;
import java.util.concurrent.TimeUnit;
-
import lombok.extern.slf4j.Slf4j;
@Slf4j
@@ -68,7 +67,7 @@ public class EventMeshAsyncSubscribe extends GrpcAbstractDemo
implements Receive
}
@Override
- public String getProtocolType() {
- return EventMeshCommon.EM_MESSAGE_PROTOCOL_NAME;
+ public EventMeshProtocolType getProtocolType() {
+ return EventMeshProtocolType.EVENT_MESH_MESSAGE;
}
}
diff --git
a/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/EventMeshSubscribeBroadcast.java
b/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/EventMeshSubscribeBroadcast.java
index 830a29def..7fcb04cf1 100644
---
a/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/EventMeshSubscribeBroadcast.java
+++
b/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/EventMeshSubscribeBroadcast.java
@@ -19,9 +19,9 @@ package org.apache.eventmesh.grpc.sub;
import org.apache.eventmesh.client.grpc.consumer.EventMeshGrpcConsumer;
import org.apache.eventmesh.client.grpc.consumer.ReceiveMsgHook;
-import org.apache.eventmesh.client.tcp.common.EventMeshCommon;
import org.apache.eventmesh.common.EventMeshMessage;
import org.apache.eventmesh.common.ExampleConstants;
+import org.apache.eventmesh.common.enums.EventMeshProtocolType;
import org.apache.eventmesh.common.protocol.SubscriptionItem;
import org.apache.eventmesh.common.protocol.SubscriptionMode;
import org.apache.eventmesh.common.protocol.SubscriptionType;
@@ -33,7 +33,6 @@ import java.util.Collections;
import java.util.Optional;
import java.util.concurrent.TimeUnit;
-
import lombok.extern.slf4j.Slf4j;
@Slf4j
@@ -69,7 +68,7 @@ public class EventMeshSubscribeBroadcast extends
GrpcAbstractDemo implements Rec
}
@Override
- public String getProtocolType() {
- return EventMeshCommon.EM_MESSAGE_PROTOCOL_NAME;
+ public EventMeshProtocolType getProtocolType() {
+ return EventMeshProtocolType.EVENT_MESH_MESSAGE;
}
}
diff --git
a/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/EventMeshSubscribeReply.java
b/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/EventMeshSubscribeReply.java
index e78bb411e..56cab1010 100644
---
a/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/EventMeshSubscribeReply.java
+++
b/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/EventMeshSubscribeReply.java
@@ -19,9 +19,9 @@ package org.apache.eventmesh.grpc.sub;
import org.apache.eventmesh.client.grpc.consumer.EventMeshGrpcConsumer;
import org.apache.eventmesh.client.grpc.consumer.ReceiveMsgHook;
-import org.apache.eventmesh.client.tcp.common.EventMeshCommon;
import org.apache.eventmesh.common.EventMeshMessage;
import org.apache.eventmesh.common.ExampleConstants;
+import org.apache.eventmesh.common.enums.EventMeshProtocolType;
import org.apache.eventmesh.common.protocol.SubscriptionItem;
import org.apache.eventmesh.common.protocol.SubscriptionMode;
import org.apache.eventmesh.common.protocol.SubscriptionType;
@@ -33,7 +33,6 @@ import java.util.Collections;
import java.util.Optional;
import java.util.concurrent.TimeUnit;
-
import lombok.extern.slf4j.Slf4j;
@Slf4j
@@ -72,7 +71,7 @@ public class EventMeshSubscribeReply extends GrpcAbstractDemo
implements Receive
}
@Override
- public String getProtocolType() {
- return EventMeshCommon.EM_MESSAGE_PROTOCOL_NAME;
+ public EventMeshProtocolType getProtocolType() {
+ return EventMeshProtocolType.EVENT_MESH_MESSAGE;
}
}
diff --git
a/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/WorkflowExpressAsyncSubscribe.java
b/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/WorkflowExpressAsyncSubscribe.java
index 4614ef9f0..4ba90adca 100644
---
a/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/WorkflowExpressAsyncSubscribe.java
+++
b/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/WorkflowExpressAsyncSubscribe.java
@@ -22,11 +22,11 @@ import
org.apache.eventmesh.client.catalog.config.EventMeshCatalogClientConfig;
import org.apache.eventmesh.client.grpc.consumer.EventMeshGrpcConsumer;
import org.apache.eventmesh.client.grpc.consumer.ReceiveMsgHook;
import org.apache.eventmesh.client.selector.SelectorFactory;
-import org.apache.eventmesh.client.tcp.common.EventMeshCommon;
import org.apache.eventmesh.client.workflow.EventMeshWorkflowClient;
import
org.apache.eventmesh.client.workflow.config.EventMeshWorkflowClientConfig;
import org.apache.eventmesh.common.EventMeshMessage;
import org.apache.eventmesh.common.ExampleConstants;
+import org.apache.eventmesh.common.enums.EventMeshProtocolType;
import org.apache.eventmesh.common.protocol.workflow.protos.ExecuteRequest;
import org.apache.eventmesh.common.protocol.workflow.protos.ExecuteResponse;
import org.apache.eventmesh.common.utils.ThreadUtils;
@@ -39,7 +39,6 @@ import java.util.Optional;
import java.util.Properties;
import java.util.concurrent.TimeUnit;
-
import lombok.extern.slf4j.Slf4j;
@Slf4j
@@ -104,7 +103,7 @@ public class WorkflowExpressAsyncSubscribe extends
GrpcAbstractDemo implements R
}
@Override
- public String getProtocolType() {
- return EventMeshCommon.EM_MESSAGE_PROTOCOL_NAME;
+ public EventMeshProtocolType getProtocolType() {
+ return EventMeshProtocolType.EVENT_MESH_MESSAGE;
}
}
diff --git
a/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/WorkflowOrderAsyncSubscribe.java
b/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/WorkflowOrderAsyncSubscribe.java
index 2f4634993..2db100392 100644
---
a/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/WorkflowOrderAsyncSubscribe.java
+++
b/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/WorkflowOrderAsyncSubscribe.java
@@ -22,11 +22,11 @@ import
org.apache.eventmesh.client.catalog.config.EventMeshCatalogClientConfig;
import org.apache.eventmesh.client.grpc.consumer.EventMeshGrpcConsumer;
import org.apache.eventmesh.client.grpc.consumer.ReceiveMsgHook;
import org.apache.eventmesh.client.selector.SelectorFactory;
-import org.apache.eventmesh.client.tcp.common.EventMeshCommon;
import org.apache.eventmesh.client.workflow.EventMeshWorkflowClient;
import
org.apache.eventmesh.client.workflow.config.EventMeshWorkflowClientConfig;
import org.apache.eventmesh.common.EventMeshMessage;
import org.apache.eventmesh.common.ExampleConstants;
+import org.apache.eventmesh.common.enums.EventMeshProtocolType;
import org.apache.eventmesh.common.protocol.workflow.protos.ExecuteRequest;
import org.apache.eventmesh.common.protocol.workflow.protos.ExecuteResponse;
import org.apache.eventmesh.common.utils.ThreadUtils;
@@ -39,7 +39,6 @@ import java.util.Optional;
import java.util.Properties;
import java.util.concurrent.TimeUnit;
-
import lombok.extern.slf4j.Slf4j;
@Slf4j
@@ -100,7 +99,7 @@ public class WorkflowOrderAsyncSubscribe extends
GrpcAbstractDemo implements Rec
}
@Override
- public String getProtocolType() {
- return EventMeshCommon.EM_MESSAGE_PROTOCOL_NAME;
+ public EventMeshProtocolType getProtocolType() {
+ return EventMeshProtocolType.EVENT_MESH_MESSAGE;
}
}
diff --git
a/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/WorkflowPaymentAsyncSubscribe.java
b/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/WorkflowPaymentAsyncSubscribe.java
index 396add5e4..7c7c9400f 100644
---
a/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/WorkflowPaymentAsyncSubscribe.java
+++
b/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/WorkflowPaymentAsyncSubscribe.java
@@ -22,11 +22,11 @@ import
org.apache.eventmesh.client.catalog.config.EventMeshCatalogClientConfig;
import org.apache.eventmesh.client.grpc.consumer.EventMeshGrpcConsumer;
import org.apache.eventmesh.client.grpc.consumer.ReceiveMsgHook;
import org.apache.eventmesh.client.selector.SelectorFactory;
-import org.apache.eventmesh.client.tcp.common.EventMeshCommon;
import org.apache.eventmesh.client.workflow.EventMeshWorkflowClient;
import
org.apache.eventmesh.client.workflow.config.EventMeshWorkflowClientConfig;
import org.apache.eventmesh.common.EventMeshMessage;
import org.apache.eventmesh.common.ExampleConstants;
+import org.apache.eventmesh.common.enums.EventMeshProtocolType;
import org.apache.eventmesh.common.protocol.workflow.protos.ExecuteRequest;
import org.apache.eventmesh.common.protocol.workflow.protos.ExecuteResponse;
import org.apache.eventmesh.common.utils.ThreadUtils;
@@ -39,7 +39,6 @@ import java.util.Optional;
import java.util.Properties;
import java.util.concurrent.TimeUnit;
-
import lombok.extern.slf4j.Slf4j;
@Slf4j
@@ -104,7 +103,7 @@ public class WorkflowPaymentAsyncSubscribe extends
GrpcAbstractDemo implements R
}
@Override
- public String getProtocolType() {
- return EventMeshCommon.EM_MESSAGE_PROTOCOL_NAME;
+ public EventMeshProtocolType getProtocolType() {
+ return EventMeshProtocolType.EVENT_MESH_MESSAGE;
}
}
diff --git
a/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/grpc/consumer/EventMeshGrpcConsumer.java
b/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/grpc/consumer/EventMeshGrpcConsumer.java
index f0e54ba1f..cf401c5b0 100644
---
a/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/grpc/consumer/EventMeshGrpcConsumer.java
+++
b/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/grpc/consumer/EventMeshGrpcConsumer.java
@@ -24,6 +24,7 @@ import
org.apache.eventmesh.client.grpc.config.EventMeshGrpcClientConfig;
import org.apache.eventmesh.client.grpc.util.EventMeshClientUtil;
import org.apache.eventmesh.client.tcp.common.EventMeshCommon;
import org.apache.eventmesh.common.EventMeshThreadFactory;
+import org.apache.eventmesh.common.enums.EventMeshProtocolType;
import org.apache.eventmesh.common.protocol.SubscriptionItem;
import org.apache.eventmesh.common.protocol.SubscriptionMode;
import org.apache.eventmesh.common.protocol.SubscriptionType;
@@ -204,7 +205,7 @@ public class EventMeshGrpcConsumer implements AutoCloseable
{
Objects.requireNonNull(subscriptionItems, "subscriptionItems can not
be null");
final Subscription.Builder builder = Subscription.newBuilder()
- .setHeader(EventMeshClientUtil.buildHeader(clientConfig,
EventMeshCommon.EM_MESSAGE_PROTOCOL_NAME))
+ .setHeader(EventMeshClientUtil.buildHeader(clientConfig,
EventMeshProtocolType.EVENT_MESH_MESSAGE))
.setConsumerGroup(clientConfig.getConsumerGroup());
if (StringUtils.isNotEmpty(url)) {
@@ -250,8 +251,7 @@ public class EventMeshGrpcConsumer implements AutoCloseable
{
}
private void heartBeat() {
- final RequestHeader header = EventMeshClientUtil.buildHeader(
- clientConfig, EventMeshCommon.EM_MESSAGE_PROTOCOL_NAME);
+ final RequestHeader header =
EventMeshClientUtil.buildHeader(clientConfig,
EventMeshProtocolType.EVENT_MESH_MESSAGE);
scheduler.scheduleAtFixedRate(() -> {
if (MapUtils.isEmpty(subscriptionMap)) {
diff --git
a/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/grpc/consumer/ReceiveMsgHook.java
b/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/grpc/consumer/ReceiveMsgHook.java
index 2e99ea2b8..24b8ca0ae 100644
---
a/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/grpc/consumer/ReceiveMsgHook.java
+++
b/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/grpc/consumer/ReceiveMsgHook.java
@@ -17,6 +17,8 @@
package org.apache.eventmesh.client.grpc.consumer;
+import org.apache.eventmesh.common.enums.EventMeshProtocolType;
+
import java.util.Optional;
/**
@@ -32,6 +34,6 @@ public interface ReceiveMsgHook<T> {
*/
Optional<T> handle(T msg) throws Exception;
- String getProtocolType();
+ EventMeshProtocolType getProtocolType();
}
diff --git
a/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/grpc/producer/CloudEventProducer.java
b/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/grpc/producer/CloudEventProducer.java
index 6d2793772..57b8c8805 100644
---
a/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/grpc/producer/CloudEventProducer.java
+++
b/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/grpc/producer/CloudEventProducer.java
@@ -19,8 +19,8 @@ package org.apache.eventmesh.client.grpc.producer;
import org.apache.eventmesh.client.grpc.config.EventMeshGrpcClientConfig;
import org.apache.eventmesh.client.grpc.util.EventMeshClientUtil;
-import org.apache.eventmesh.client.tcp.common.EventMeshCommon;
import org.apache.eventmesh.common.Constants;
+import org.apache.eventmesh.common.enums.EventMeshProtocolType;
import org.apache.eventmesh.common.protocol.grpc.common.ProtocolKey;
import org.apache.eventmesh.common.protocol.grpc.protos.BatchMessage;
import
org.apache.eventmesh.common.protocol.grpc.protos.PublisherServiceGrpc.PublisherServiceBlockingStub;
@@ -43,9 +43,9 @@ import io.cloudevents.core.builder.CloudEventBuilder;
import lombok.extern.slf4j.Slf4j;
@Slf4j
-public class CloudEventProducer {
+public class CloudEventProducer implements GrpcProducer<CloudEvent> {
- private static final String PROTOCOL_TYPE =
EventMeshCommon.CLOUD_EVENTS_PROTOCOL_NAME;
+ private static final EventMeshProtocolType PROTOCOL_TYPE =
EventMeshProtocolType.CLOUD_EVENTS;
private final transient EventMeshGrpcClientConfig clientConfig;
@@ -57,6 +57,7 @@ public class CloudEventProducer {
this.publisherClient = publisherClient;
}
+ @Override
public Response publish(final List<CloudEvent> events) {
if (log.isInfoEnabled()) {
log.info("BatchPublish message, batch size={}", events.size());
@@ -74,7 +75,7 @@ public class CloudEventProducer {
try {
final Response response =
publisherClient.batchPublish(enhancedMessage);
if (log.isInfoEnabled()) {
- log.info("Received response {}", response);
+ log.info("Received response:{}", response.toString());
}
return response;
} catch (Exception e) {
@@ -85,9 +86,10 @@ public class CloudEventProducer {
return null;
}
+ @Override
public Response publish(final CloudEvent cloudEvent) {
if (log.isInfoEnabled()) {
- log.info("Publish message {}", cloudEvent);
+ log.info("Publish message: {}", cloudEvent.toString());
}
final CloudEvent enhanceEvent = enhanceCloudEvent(cloudEvent, null);
@@ -96,7 +98,7 @@ public class CloudEventProducer {
try {
final Response response = publisherClient.publish(enhancedMessage);
if (log.isInfoEnabled()) {
- log.info("Received response {}", response);
+ log.info("Received response:{} ", response.toString());
}
return response;
} catch (Exception e) {
@@ -107,7 +109,8 @@ public class CloudEventProducer {
return null;
}
- public CloudEvent requestReply(final CloudEvent cloudEvent, final int
timeout) {
+ @Override
+ public CloudEvent requestReply(final CloudEvent cloudEvent, final long
timeout) {
if (log.isInfoEnabled()) {
log.info("RequestReply message {}", cloudEvent);
}
@@ -143,7 +146,7 @@ public class CloudEventProducer {
.withExtension(ProtocolKey.PID,
Long.toString(ThreadUtils.getPID()))
.withExtension(ProtocolKey.SYS, clientConfig.getSys())
.withExtension(ProtocolKey.LANGUAGE, Constants.LANGUAGE_JAVA)
- .withExtension(ProtocolKey.PROTOCOL_TYPE, PROTOCOL_TYPE)
+ .withExtension(ProtocolKey.PROTOCOL_TYPE,
PROTOCOL_TYPE.protocolTypeName())
.withExtension(ProtocolKey.PROTOCOL_DESC, Constants.PROTOCOL_GRPC)
.withExtension(ProtocolKey.PROTOCOL_VERSION,
cloudEvent.getSpecVersion().toString())
.withExtension(ProtocolKey.UNIQUE_ID,
RandomStringUtils.generateNum(30))
diff --git
a/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/grpc/producer/EventMeshGrpcProducer.java
b/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/grpc/producer/EventMeshGrpcProducer.java
index 297c5b2b3..b9f82af51 100644
---
a/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/grpc/producer/EventMeshGrpcProducer.java
+++
b/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/grpc/producer/EventMeshGrpcProducer.java
@@ -18,19 +18,15 @@
package org.apache.eventmesh.client.grpc.producer;
import org.apache.eventmesh.client.grpc.config.EventMeshGrpcClientConfig;
-import org.apache.eventmesh.client.grpc.util.EventMeshClientUtil;
import org.apache.eventmesh.client.tcp.common.EventMeshCommon;
import org.apache.eventmesh.common.EventMeshMessage;
-import org.apache.eventmesh.common.protocol.grpc.protos.BatchMessage;
import org.apache.eventmesh.common.protocol.grpc.protos.PublisherServiceGrpc;
import
org.apache.eventmesh.common.protocol.grpc.protos.PublisherServiceGrpc.PublisherServiceBlockingStub;
import org.apache.eventmesh.common.protocol.grpc.protos.Response;
-import org.apache.eventmesh.common.protocol.grpc.protos.SimpleMessage;
import org.apache.commons.collections4.CollectionUtils;
import java.util.List;
-import java.util.concurrent.TimeUnit;
import io.cloudevents.CloudEvent;
import io.grpc.ManagedChannel;
@@ -49,37 +45,31 @@ public class EventMeshGrpcProducer implements AutoCloseable
{
private final transient ManagedChannel channel;
- private transient PublisherServiceBlockingStub publisherClient;
+ private PublisherServiceBlockingStub publisherClient;
- private transient CloudEventProducer cloudEventProducer;
+ private CloudEventProducer cloudEventProducer;
+
+ private EventMeshMessageProducer eventMeshMessageProducer;
public EventMeshGrpcProducer(EventMeshGrpcClientConfig clientConfig) {
this.clientConfig = clientConfig;
- channel =
ManagedChannelBuilder.forAddress(clientConfig.getServerAddr(),
clientConfig.getServerPort())
- .usePlaintext().build();
- publisherClient = PublisherServiceGrpc.newBlockingStub(channel);
-
- cloudEventProducer = new CloudEventProducer(clientConfig,
publisherClient);
+ this.channel =
ManagedChannelBuilder.forAddress(clientConfig.getServerAddr(),
clientConfig.getServerPort()).usePlaintext().build();
+ this.publisherClient = PublisherServiceGrpc.newBlockingStub(channel);
+ this.cloudEventProducer = new CloudEventProducer(clientConfig,
publisherClient);
+ this.eventMeshMessageProducer = new
EventMeshMessageProducer(clientConfig, publisherClient);
}
- public Response publish(EventMeshMessage message) {
+ public <T> Response publish(T message) {
if (log.isInfoEnabled()) {
log.info("Publish message " + message.toString());
}
-
- SimpleMessage simpleMessage =
EventMeshClientUtil.buildSimpleMessage(message, clientConfig, PROTOCOL_TYPE);
- try {
- Response response = publisherClient.publish(simpleMessage);
- if (log.isInfoEnabled()) {
- log.info("Received response:{}", response);
- }
- return response;
- } catch (Exception e) {
- if (log.isErrorEnabled()) {
- log.error("Error in publishing message {}", message, e);
- }
+ if (message instanceof CloudEvent) {
+ return cloudEventProducer.publish((CloudEvent) message);
+ } else if (message instanceof EventMeshMessage) {
+ return eventMeshMessageProducer.publish((EventMeshMessage)
message);
+ } else {
+ throw new IllegalArgumentException("Not support message " +
message.getClass().getName());
}
- return null;
}
@SuppressWarnings("unchecked")
@@ -92,58 +82,26 @@ public class EventMeshGrpcProducer implements AutoCloseable
{
return null;
}
- if (messageList.get(0) instanceof CloudEvent) {
+ T target = messageList.get(0);
+ if (target instanceof CloudEvent) {
return cloudEventProducer.publish((List<CloudEvent>) messageList);
+ } else if (target instanceof EventMeshMessage) {
+ return eventMeshMessageProducer.publish((List<EventMeshMessage>)
messageList);
+ } else {
+ throw new IllegalArgumentException("Not support message " +
target.getClass().getName());
}
-
- BatchMessage batchMessage =
EventMeshClientUtil.buildBatchMessages(messageList, clientConfig,
PROTOCOL_TYPE);
- try {
- Response response = publisherClient.batchPublish(batchMessage);
- if (log.isInfoEnabled()) {
- log.info("Received response:{}", response);
- }
- return response;
- } catch (Exception e) {
- if (log.isErrorEnabled()) {
- log.error("Error in BatchPublish message {}", messageList, e);
- }
- }
- return null;
}
- public Response publish(final CloudEvent cloudEvent) {
- return cloudEventProducer.publish(cloudEvent);
- }
-
- public CloudEvent requestReply(final CloudEvent cloudEvent, final int
timeout) {
- return cloudEventProducer.requestReply(cloudEvent, timeout);
- }
-
- public EventMeshMessage requestReply(final EventMeshMessage message, final
int timeout) {
- if (log.isInfoEnabled()) {
- log.info("RequestReply message:{}", message);
- }
+ public <T> T requestReply(final T message, final long timeout) {
- SimpleMessage simpleMessage =
EventMeshClientUtil.buildSimpleMessage(message, clientConfig, PROTOCOL_TYPE);
- try {
- SimpleMessage reply = publisherClient.withDeadlineAfter(timeout,
TimeUnit.MILLISECONDS)
- .requestReply(simpleMessage);
- if (log.isInfoEnabled()) {
- log.info("Received reply message:{}", reply);
- }
-
- final Object msg = EventMeshClientUtil.buildMessage(reply,
PROTOCOL_TYPE);
- if (msg instanceof EventMeshMessage) {
- return (EventMeshMessage) msg;
- } else {
- return null;
- }
- } catch (Exception e) {
- if (log.isErrorEnabled()) {
- log.error("Error in RequestReply message {}", message, e);
- }
+ if (message instanceof CloudEvent) {
+ CloudEvent cloudEvent =
cloudEventProducer.requestReply((CloudEvent) message, timeout);
+ return (T) cloudEvent;
+ } else if (message instanceof EventMeshMessage) {
+ return (T)
eventMeshMessageProducer.requestReply((EventMeshMessage) message, timeout);
+ } else {
+ throw new IllegalArgumentException("Not support message " +
message.getClass().getName());
}
- return null;
}
@Override
diff --git
a/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/grpc/producer/EventMeshGrpcProducer.java
b/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/grpc/producer/EventMeshMessageProducer.java
similarity index 62%
copy from
eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/grpc/producer/EventMeshGrpcProducer.java
copy to
eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/grpc/producer/EventMeshMessageProducer.java
index 297c5b2b3..a298b0838 100644
---
a/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/grpc/producer/EventMeshGrpcProducer.java
+++
b/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/grpc/producer/EventMeshMessageProducer.java
@@ -19,10 +19,9 @@ package org.apache.eventmesh.client.grpc.producer;
import org.apache.eventmesh.client.grpc.config.EventMeshGrpcClientConfig;
import org.apache.eventmesh.client.grpc.util.EventMeshClientUtil;
-import org.apache.eventmesh.client.tcp.common.EventMeshCommon;
import org.apache.eventmesh.common.EventMeshMessage;
+import org.apache.eventmesh.common.enums.EventMeshProtocolType;
import org.apache.eventmesh.common.protocol.grpc.protos.BatchMessage;
-import org.apache.eventmesh.common.protocol.grpc.protos.PublisherServiceGrpc;
import
org.apache.eventmesh.common.protocol.grpc.protos.PublisherServiceGrpc.PublisherServiceBlockingStub;
import org.apache.eventmesh.common.protocol.grpc.protos.Response;
import org.apache.eventmesh.common.protocol.grpc.protos.SimpleMessage;
@@ -32,41 +31,32 @@ import org.apache.commons.collections4.CollectionUtils;
import java.util.List;
import java.util.concurrent.TimeUnit;
-import io.cloudevents.CloudEvent;
-import io.grpc.ManagedChannel;
-import io.grpc.ManagedChannelBuilder;
-
-import lombok.Data;
import lombok.extern.slf4j.Slf4j;
@Slf4j
-@Data
-public class EventMeshGrpcProducer implements AutoCloseable {
-
- private static final String PROTOCOL_TYPE =
EventMeshCommon.EM_MESSAGE_PROTOCOL_NAME;
+public class EventMeshMessageProducer implements
GrpcProducer<EventMeshMessage> {
- private final transient EventMeshGrpcClientConfig clientConfig;
+ private static final EventMeshProtocolType PROTOCOL_TYPE =
EventMeshProtocolType.EVENT_MESH_MESSAGE;
- private final transient ManagedChannel channel;
+ private final EventMeshGrpcClientConfig clientConfig;
- private transient PublisherServiceBlockingStub publisherClient;
+ private final PublisherServiceBlockingStub publisherClient;
- private transient CloudEventProducer cloudEventProducer;
-
- public EventMeshGrpcProducer(EventMeshGrpcClientConfig clientConfig) {
+ public EventMeshMessageProducer(EventMeshGrpcClientConfig clientConfig,
PublisherServiceBlockingStub publisherClient) {
this.clientConfig = clientConfig;
- channel =
ManagedChannelBuilder.forAddress(clientConfig.getServerAddr(),
clientConfig.getServerPort())
- .usePlaintext().build();
- publisherClient = PublisherServiceGrpc.newBlockingStub(channel);
-
- cloudEventProducer = new CloudEventProducer(clientConfig,
publisherClient);
+ this.publisherClient = publisherClient;
}
+ @Override
public Response publish(EventMeshMessage message) {
- if (log.isInfoEnabled()) {
- log.info("Publish message " + message.toString());
+
+ if (null == message) {
+ return null;
}
+ if (log.isDebugEnabled()) {
+ log.info("Publish message: {}", message.toString());
+ }
SimpleMessage simpleMessage =
EventMeshClientUtil.buildSimpleMessage(message, clientConfig, PROTOCOL_TYPE);
try {
Response response = publisherClient.publish(simpleMessage);
@@ -75,28 +65,18 @@ public class EventMeshGrpcProducer implements AutoCloseable
{
}
return response;
} catch (Exception e) {
- if (log.isErrorEnabled()) {
- log.error("Error in publishing message {}", message, e);
- }
+ log.error("Error in publishing message {}", message, e);
}
return null;
}
- @SuppressWarnings("unchecked")
- public <T> Response publish(List<T> messageList) {
- if (log.isInfoEnabled()) {
- log.info("BatchPublish message :{}", messageList);
- }
+ @Override
+ public Response publish(List<EventMeshMessage> messages) {
- if (CollectionUtils.isEmpty(messageList)) {
+ if (CollectionUtils.isEmpty(messages)) {
return null;
}
-
- if (messageList.get(0) instanceof CloudEvent) {
- return cloudEventProducer.publish((List<CloudEvent>) messageList);
- }
-
- BatchMessage batchMessage =
EventMeshClientUtil.buildBatchMessages(messageList, clientConfig,
PROTOCOL_TYPE);
+ BatchMessage batchMessage =
EventMeshClientUtil.buildBatchMessages(messages, clientConfig, PROTOCOL_TYPE);
try {
Response response = publisherClient.batchPublish(batchMessage);
if (log.isInfoEnabled()) {
@@ -105,21 +85,14 @@ public class EventMeshGrpcProducer implements
AutoCloseable {
return response;
} catch (Exception e) {
if (log.isErrorEnabled()) {
- log.error("Error in BatchPublish message {}", messageList, e);
+ log.error("Error in BatchPublish message {}", messages, e);
}
}
return null;
}
- public Response publish(final CloudEvent cloudEvent) {
- return cloudEventProducer.publish(cloudEvent);
- }
-
- public CloudEvent requestReply(final CloudEvent cloudEvent, final int
timeout) {
- return cloudEventProducer.requestReply(cloudEvent, timeout);
- }
-
- public EventMeshMessage requestReply(final EventMeshMessage message, final
int timeout) {
+ @Override
+ public EventMeshMessage requestReply(EventMeshMessage message, long
timeout) {
if (log.isInfoEnabled()) {
log.info("RequestReply message:{}", message);
}
@@ -139,15 +112,11 @@ public class EventMeshGrpcProducer implements
AutoCloseable {
return null;
}
} catch (Exception e) {
+ e.printStackTrace();
if (log.isErrorEnabled()) {
log.error("Error in RequestReply message {}", message, e);
}
}
return null;
}
-
- @Override
- public void close() {
- channel.shutdown();
- }
}
diff --git
a/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/grpc/consumer/ReceiveMsgHook.java
b/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/grpc/producer/GrpcProducer.java
similarity index 72%
copy from
eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/grpc/consumer/ReceiveMsgHook.java
copy to
eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/grpc/producer/GrpcProducer.java
index 2e99ea2b8..a209d1e85 100644
---
a/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/grpc/consumer/ReceiveMsgHook.java
+++
b/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/grpc/producer/GrpcProducer.java
@@ -15,23 +15,22 @@
* limitations under the License.
*/
-package org.apache.eventmesh.client.grpc.consumer;
+package org.apache.eventmesh.client.grpc.producer;
-import java.util.Optional;
+import org.apache.eventmesh.common.protocol.grpc.protos.Response;
+
+import java.util.List;
/**
+ *
* @param <T>
*/
-public interface ReceiveMsgHook<T> {
+public interface GrpcProducer<T> {
+
+ Response publish(T message);
- /**
- * Handle the received message, return the response message.
- *
- * @param msg
- * @return
- */
- Optional<T> handle(T msg) throws Exception;
+ Response publish(List<T> messages);
- String getProtocolType();
+ T requestReply(T message, long timeout);
}
diff --git
a/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/grpc/util/EventMeshClientUtil.java
b/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/grpc/util/EventMeshClientUtil.java
index 97ae6fb6a..ba46523b7 100644
---
a/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/grpc/util/EventMeshClientUtil.java
+++
b/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/grpc/util/EventMeshClientUtil.java
@@ -18,9 +18,9 @@
package org.apache.eventmesh.client.grpc.util;
import org.apache.eventmesh.client.grpc.config.EventMeshGrpcClientConfig;
-import org.apache.eventmesh.client.tcp.common.EventMeshCommon;
import org.apache.eventmesh.common.Constants;
import org.apache.eventmesh.common.EventMeshMessage;
+import org.apache.eventmesh.common.enums.EventMeshProtocolType;
import org.apache.eventmesh.common.protocol.grpc.common.ProtocolKey;
import org.apache.eventmesh.common.protocol.grpc.protos.BatchMessage;
import org.apache.eventmesh.common.protocol.grpc.protos.RequestHeader;
@@ -49,7 +49,7 @@ import com.fasterxml.jackson.core.type.TypeReference;
public class EventMeshClientUtil {
- public static RequestHeader buildHeader(EventMeshGrpcClientConfig
clientConfig, String protocolType) {
+ public static RequestHeader buildHeader(EventMeshGrpcClientConfig
clientConfig, EventMeshProtocolType protocolType) {
return RequestHeader.newBuilder()
.setEnv(clientConfig.getEnv())
.setIdc(clientConfig.getIdc())
@@ -59,7 +59,7 @@ public class EventMeshClientUtil {
.setLanguage(clientConfig.getLanguage())
.setUsername(clientConfig.getUserName())
.setPassword(clientConfig.getPassword())
- .setProtocolType(protocolType)
+ .setProtocolType(protocolType.protocolTypeName())
.setProtocolDesc(Constants.PROTOCOL_GRPC)
// default CloudEvents version is V1
.setProtocolVersion(SpecVersion.V1.toString())
@@ -67,7 +67,11 @@ public class EventMeshClientUtil {
}
@SuppressWarnings("unchecked")
- public static <T> T buildMessage(final SimpleMessage message, final String
protocolType) {
+ public static <T> T buildMessage(final SimpleMessage message, final
EventMeshProtocolType protocolType) {
+
+ if (null == message) {
+ return null;
+ }
final String seq = message.getSeqNum();
final String uniqueId = message.getUniqueId();
final String content = message.getContent();
@@ -78,151 +82,184 @@ public class EventMeshClientUtil {
new TypeReference<HashMap<String, String>>() {
});
}
+ if (null == protocolType) {
+ return null;
+ }
- if (EventMeshCommon.CLOUD_EVENTS_PROTOCOL_NAME.equals(protocolType)) {
- final String contentType =
message.getPropertiesOrDefault(ProtocolKey.CONTENT_TYPE,
JsonFormat.CONTENT_TYPE);
- final CloudEvent cloudEvent =
Objects.requireNonNull(EventFormatProvider.getInstance().resolveFormat(contentType))
- .deserialize(content.getBytes(StandardCharsets.UTF_8));
-
- final CloudEventBuilder cloudEventBuilder =
CloudEventBuilder.from(cloudEvent)
- .withSubject(message.getTopic())
- .withExtension(ProtocolKey.SEQ_NUM, message.getSeqNum())
- .withExtension(ProtocolKey.UNIQUE_ID, message.getUniqueId());
-
-
message.getPropertiesMap().forEach(cloudEventBuilder::withExtension);
-
- return (T) cloudEventBuilder.build();
- } else {
- return (T) EventMeshMessage.builder()
- .content(content)
- .topic(message.getTopic())
- .bizSeqNo(seq)
- .uniqueId(uniqueId)
- .prop(message.getPropertiesMap())
- .build();
+ switch (protocolType) {
+ case CLOUD_EVENTS:
+ return (T) switchSimpleMessage2CloudEvent(message, content);
+ case EVENT_MESH_MESSAGE:
+ return (T) switchSimpleMessage2EventMeshMessage(message, seq,
uniqueId, content);
+ case OPEN_MESSAGE:
+ default:
+ return null;
}
}
- public static <T> SimpleMessage buildSimpleMessage(final T message, final
EventMeshGrpcClientConfig clientConfig,
- final String protocolType) {
- if (EventMeshCommon.CLOUD_EVENTS_PROTOCOL_NAME.equals(protocolType)) {
- final CloudEvent cloudEvent = (CloudEvent) message;
- final String contentType =
StringUtils.isEmpty(cloudEvent.getDataContentType())
- ? Constants.CONTENT_TYPE_CLOUDEVENTS_JSON
- : cloudEvent.getDataContentType();
- final byte[] bodyByte =
Objects.requireNonNull(EventFormatProvider.getInstance().resolveFormat(contentType))
- .serialize(cloudEvent);
- final String content = new String(bodyByte,
StandardCharsets.UTF_8);
- final String ttl =
cloudEvent.getExtension(Constants.EVENTMESH_MESSAGE_CONST_TTL) == null
- ? Constants.DEFAULT_EVENTMESH_MESSAGE_TTL
- :
Objects.requireNonNull(cloudEvent.getExtension(Constants.EVENTMESH_MESSAGE_CONST_TTL)).toString();
-
- final String seqNum = cloudEvent.getExtension(ProtocolKey.SEQ_NUM)
== null
- ? RandomStringUtils.generateNum(30)
- :
Objects.requireNonNull(cloudEvent.getExtension(ProtocolKey.SEQ_NUM)).toString();
+ private static EventMeshMessage
switchSimpleMessage2EventMeshMessage(SimpleMessage message, String seq, String
uniqueId, String content) {
+ return EventMeshMessage.builder()
+ .content(content)
+ .topic(message.getTopic())
+ .bizSeqNo(seq)
+ .uniqueId(uniqueId)
+ .prop(message.getPropertiesMap())
+ .build();
+ }
- final String uniqueId =
cloudEvent.getExtension(ProtocolKey.UNIQUE_ID) == null
- ? RandomStringUtils.generateNum(30)
- :
Objects.requireNonNull(cloudEvent.getExtension(ProtocolKey.UNIQUE_ID)).toString();
+ private static CloudEvent switchSimpleMessage2CloudEvent(SimpleMessage
message, String content) {
+ final String contentType =
message.getPropertiesOrDefault(ProtocolKey.CONTENT_TYPE,
JsonFormat.CONTENT_TYPE);
+ final CloudEvent cloudEvent =
Objects.requireNonNull(EventFormatProvider.getInstance().resolveFormat(contentType))
+ .deserialize(content.getBytes(StandardCharsets.UTF_8));
- final SimpleMessage.Builder builder = SimpleMessage.newBuilder()
- .setHeader(EventMeshClientUtil.buildHeader(clientConfig,
protocolType))
- .setProducerGroup(clientConfig.getProducerGroup())
- .setTopic(cloudEvent.getSubject())
- .setTtl(ttl)
- .setSeqNum(seqNum)
- .setUniqueId(uniqueId)
- .setContent(content)
- .putProperties(ProtocolKey.CONTENT_TYPE, contentType);
+ final CloudEventBuilder cloudEventBuilder =
CloudEventBuilder.from(cloudEvent)
+ .withSubject(message.getTopic())
+ .withExtension(ProtocolKey.SEQ_NUM, message.getSeqNum())
+ .withExtension(ProtocolKey.UNIQUE_ID, message.getUniqueId());
- cloudEvent.getExtensionNames().forEach(extName -> {
- builder.putProperties(extName,
Objects.requireNonNull(cloudEvent.getExtension(extName)).toString());
- });
- return builder.build();
+ message.getPropertiesMap().forEach(cloudEventBuilder::withExtension);
- } else {
- final EventMeshMessage eventMeshMessage = (EventMeshMessage)
message;
+ return cloudEventBuilder.build();
+ }
- final String ttl =
eventMeshMessage.getProp(Constants.EVENTMESH_MESSAGE_CONST_TTL) == null
- ? Constants.DEFAULT_EVENTMESH_MESSAGE_TTL
- :
eventMeshMessage.getProp(Constants.EVENTMESH_MESSAGE_CONST_TTL);
- final Map<String, String> props = eventMeshMessage.getProp() ==
null
- ? new HashMap<>() : eventMeshMessage.getProp();
-
- final String seqNum = eventMeshMessage.getBizSeqNo() == null ?
RandomStringUtils.generateNum(30)
- : eventMeshMessage.getBizSeqNo();
-
- final String uniqueId = eventMeshMessage.getUniqueId() == null ?
RandomStringUtils.generateNum(30)
- : eventMeshMessage.getUniqueId();
-
- return SimpleMessage.newBuilder()
- .setHeader(EventMeshClientUtil.buildHeader(clientConfig,
protocolType))
- .setProducerGroup(clientConfig.getProducerGroup())
- .setTopic(eventMeshMessage.getTopic())
- .setContent(eventMeshMessage.getContent())
- .setSeqNum(seqNum)
- .setUniqueId(uniqueId)
- .setTtl(ttl)
- .putAllProperties(props)
- .build();
+ public static <T> SimpleMessage buildSimpleMessage(final T message, final
EventMeshGrpcClientConfig clientConfig,
+ final EventMeshProtocolType protocolType) {
+
+ switch (protocolType) {
+ case CLOUD_EVENTS:
+ return switchCloudEvents2SimpleMessage((CloudEvent) message,
clientConfig, protocolType);
+ case EVENT_MESH_MESSAGE:
+ return switchEventMessage2SimpleMessage((EventMeshMessage)
message, clientConfig, protocolType);
+ case OPEN_MESSAGE:
+ default:
+ return null;
}
}
+ private static SimpleMessage
switchEventMessage2SimpleMessage(EventMeshMessage message,
EventMeshGrpcClientConfig clientConfig,
+ EventMeshProtocolType protocolType) {
+ final EventMeshMessage eventMeshMessage = message;
+
+ final String ttl =
eventMeshMessage.getProp(Constants.EVENTMESH_MESSAGE_CONST_TTL) == null
+ ? Constants.DEFAULT_EVENTMESH_MESSAGE_TTL :
eventMeshMessage.getProp(Constants.EVENTMESH_MESSAGE_CONST_TTL);
+ final Map<String, String> props = eventMeshMessage.getProp() == null ?
new HashMap<>() : eventMeshMessage.getProp();
+ final String seqNum = eventMeshMessage.getBizSeqNo() == null ?
RandomStringUtils.generateNum(30) : eventMeshMessage.getBizSeqNo();
+ final String uniqueId = eventMeshMessage.getUniqueId() == null ?
RandomStringUtils.generateNum(30) : eventMeshMessage.getUniqueId();
+
+ return SimpleMessage.newBuilder()
+ .setHeader(EventMeshClientUtil.buildHeader(clientConfig,
protocolType))
+ .setProducerGroup(clientConfig.getProducerGroup())
+ .setTopic(eventMeshMessage.getTopic())
+ .setContent(eventMeshMessage.getContent())
+ .setSeqNum(seqNum)
+ .setUniqueId(uniqueId)
+ .setTtl(ttl)
+ .putAllProperties(props)
+ .build();
+ }
+
+ private static SimpleMessage switchCloudEvents2SimpleMessage(CloudEvent
message, EventMeshGrpcClientConfig clientConfig,
+ EventMeshProtocolType protocolType) {
+ final CloudEvent cloudEvent = message;
+ final String contentType =
StringUtils.isEmpty(cloudEvent.getDataContentType())
+ ? Constants.CONTENT_TYPE_CLOUDEVENTS_JSON
+ : cloudEvent.getDataContentType();
+ final byte[] bodyByte =
Objects.requireNonNull(EventFormatProvider.getInstance().resolveFormat(contentType))
+ .serialize(cloudEvent);
+ final String content = new String(bodyByte, StandardCharsets.UTF_8);
+ final String ttl =
cloudEvent.getExtension(Constants.EVENTMESH_MESSAGE_CONST_TTL) == null
+ ? Constants.DEFAULT_EVENTMESH_MESSAGE_TTL
+ :
Objects.requireNonNull(cloudEvent.getExtension(Constants.EVENTMESH_MESSAGE_CONST_TTL)).toString();
+
+ final String seqNum = cloudEvent.getExtension(ProtocolKey.SEQ_NUM) ==
null
+ ? RandomStringUtils.generateNum(30)
+ :
Objects.requireNonNull(cloudEvent.getExtension(ProtocolKey.SEQ_NUM)).toString();
+
+ final String uniqueId = cloudEvent.getExtension(ProtocolKey.UNIQUE_ID)
== null
+ ? RandomStringUtils.generateNum(30)
+ :
Objects.requireNonNull(cloudEvent.getExtension(ProtocolKey.UNIQUE_ID)).toString();
+
+ final SimpleMessage.Builder builder = SimpleMessage.newBuilder()
+ .setHeader(EventMeshClientUtil.buildHeader(clientConfig,
protocolType))
+ .setProducerGroup(clientConfig.getProducerGroup())
+ .setTopic(cloudEvent.getSubject())
+ .setTtl(ttl)
+ .setSeqNum(seqNum)
+ .setUniqueId(uniqueId)
+ .setContent(content)
+ .putProperties(ProtocolKey.CONTENT_TYPE, contentType);
+
+ cloudEvent.getExtensionNames().forEach(extName -> {
+ builder.putProperties(extName,
Objects.requireNonNull(cloudEvent.getExtension(extName)).toString());
+ });
+ return builder.build();
+ }
+
@SuppressWarnings("unchecked")
- public static <T> BatchMessage buildBatchMessages(final List<T>
messageList,
- final EventMeshGrpcClientConfig clientConfig,
- final String protocolType) {
- if (EventMeshCommon.CLOUD_EVENTS_PROTOCOL_NAME.equals(protocolType)) {
- List<CloudEvent> events = (List<CloudEvent>) messageList;
- BatchMessage.Builder messageBuilder = BatchMessage.newBuilder()
- .setHeader(EventMeshClientUtil.buildHeader(clientConfig,
protocolType))
- .setProducerGroup(clientConfig.getProducerGroup())
- .setTopic(events.get(0).getSubject());
-
- events.forEach(event -> {
- final String contentType =
StringUtils.isEmpty(event.getDataContentType())
- ? Constants.CONTENT_TYPE_CLOUDEVENTS_JSON
- : event.getDataContentType();
-
- final byte[] bodyByte =
Objects.requireNonNull(EventFormatProvider.getInstance().resolveFormat(contentType))
- .serialize(event);
- final String content = new String(bodyByte,
StandardCharsets.UTF_8);
-
- final String ttl =
event.getExtension(Constants.EVENTMESH_MESSAGE_CONST_TTL) == null
- ? Constants.DEFAULT_EVENTMESH_MESSAGE_TTL
- :
Objects.requireNonNull(event.getExtension(Constants.EVENTMESH_MESSAGE_CONST_TTL)).toString();
-
- BatchMessage.MessageItem messageItem =
BatchMessage.MessageItem.newBuilder()
- .setContent(content)
- .setTtl(ttl)
-
.setSeqNum(Objects.requireNonNull(event.getExtension(ProtocolKey.SEQ_NUM)).toString())
-
.setUniqueId(Objects.requireNonNull(event.getExtension(ProtocolKey.UNIQUE_ID)).toString())
- .putProperties(ProtocolKey.CONTENT_TYPE, contentType)
- .build();
-
- messageBuilder.addMessageItem(messageItem);
- });
- return messageBuilder.build();
- } else {
- final List<EventMeshMessage> eventMeshMessages =
(List<EventMeshMessage>) messageList;
- final BatchMessage.Builder messageBuilder =
BatchMessage.newBuilder()
- .setHeader(EventMeshClientUtil.buildHeader(clientConfig,
protocolType))
- .setProducerGroup(clientConfig.getProducerGroup())
- .setTopic(eventMeshMessages.get(0).getTopic());
-
- eventMeshMessages.forEach(message -> {
- BatchMessage.MessageItem item =
BatchMessage.MessageItem.newBuilder()
- .setContent(message.getContent())
- .setUniqueId(message.getUniqueId())
- .setSeqNum(message.getBizSeqNo())
-
.setTtl(Optional.ofNullable(message.getProp(Constants.EVENTMESH_MESSAGE_CONST_TTL))
- .orElseGet(() ->
Constants.DEFAULT_EVENTMESH_MESSAGE_TTL))
- .putAllProperties(message.getProp())
- .build();
- messageBuilder.addMessageItem(item);
- });
-
- return messageBuilder.build();
+ public static <T> BatchMessage buildBatchMessages(final List<T>
messageList, final EventMeshGrpcClientConfig clientConfig,
+ final EventMeshProtocolType protocolType) {
+
+ switch (protocolType) {
+ case CLOUD_EVENTS:
+ return switchCloudEvents2BatchMessage((List<CloudEvent>)
messageList, clientConfig, protocolType);
+ case EVENT_MESH_MESSAGE:
+ return
switchEventMessages2BatchMessage((List<EventMeshMessage>) messageList,
clientConfig, protocolType);
+ case OPEN_MESSAGE:
+ default:
+ return null;
}
}
+
+ private static BatchMessage
switchEventMessages2BatchMessage(List<EventMeshMessage> messageList,
EventMeshGrpcClientConfig clientConfig,
+ EventMeshProtocolType protocolType) {
+
+ final List<EventMeshMessage> eventMeshMessages = messageList;
+ final BatchMessage.Builder messageBuilder = BatchMessage.newBuilder()
+ .setHeader(EventMeshClientUtil.buildHeader(clientConfig,
protocolType))
+ .setProducerGroup(clientConfig.getProducerGroup())
+ .setTopic(eventMeshMessages.get(0).getTopic());
+
+ eventMeshMessages.forEach(message -> {
+ BatchMessage.MessageItem item =
BatchMessage.MessageItem.newBuilder()
+ .setContent(message.getContent())
+ .setUniqueId(message.getUniqueId())
+ .setSeqNum(message.getBizSeqNo())
+
.setTtl(Optional.ofNullable(message.getProp(Constants.EVENTMESH_MESSAGE_CONST_TTL)).orElse(Constants.DEFAULT_EVENTMESH_MESSAGE_TTL))
+ .putAllProperties(message.getProp())
+ .build();
+ messageBuilder.addMessageItem(item);
+ });
+
+ return messageBuilder.build();
+ }
+
+ private static BatchMessage
switchCloudEvents2BatchMessage(List<CloudEvent> messageList,
EventMeshGrpcClientConfig clientConfig,
+ EventMeshProtocolType protocolType) {
+ List<CloudEvent> events = messageList;
+ BatchMessage.Builder messageBuilder = BatchMessage.newBuilder()
+ .setHeader(EventMeshClientUtil.buildHeader(clientConfig,
protocolType))
+ .setProducerGroup(clientConfig.getProducerGroup())
+ .setTopic(events.get(0).getSubject());
+
+ events.forEach(event -> {
+ final String contentType =
StringUtils.isEmpty(event.getDataContentType())
+ ? Constants.CONTENT_TYPE_CLOUDEVENTS_JSON :
event.getDataContentType();
+
+ final byte[] bodyByte =
Objects.requireNonNull(EventFormatProvider.getInstance().resolveFormat(contentType)).serialize(event);
+ final String content = new String(bodyByte,
Constants.DEFAULT_CHARSET);
+ final String ttl =
event.getExtension(Constants.EVENTMESH_MESSAGE_CONST_TTL) == null
+ ? Constants.DEFAULT_EVENTMESH_MESSAGE_TTL
+ :
Objects.requireNonNull(event.getExtension(Constants.EVENTMESH_MESSAGE_CONST_TTL)).toString();
+
+ BatchMessage.MessageItem messageItem =
BatchMessage.MessageItem.newBuilder()
+ .setContent(content)
+ .setTtl(ttl)
+
.setSeqNum(Objects.requireNonNull(event.getExtension(ProtocolKey.SEQ_NUM)).toString())
+
.setUniqueId(Objects.requireNonNull(event.getExtension(ProtocolKey.UNIQUE_ID)).toString())
+ .putProperties(ProtocolKey.CONTENT_TYPE, contentType)
+ .build();
+ messageBuilder.addMessageItem(messageItem);
+ });
+ return messageBuilder.build();
+ }
}
diff --git
a/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/grpc/consumer/EventMeshGrpcConsumerTest.java
b/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/grpc/consumer/EventMeshGrpcConsumerTest.java
index 49d7eedd5..fa60200e8 100644
---
a/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/grpc/consumer/EventMeshGrpcConsumerTest.java
+++
b/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/grpc/consumer/EventMeshGrpcConsumerTest.java
@@ -26,6 +26,7 @@ import static org.mockito.Mockito.when;
import org.apache.eventmesh.client.grpc.config.EventMeshGrpcClientConfig;
import org.apache.eventmesh.common.EventMeshMessage;
+import org.apache.eventmesh.common.enums.EventMeshProtocolType;
import org.apache.eventmesh.common.protocol.SubscriptionItem;
import org.apache.eventmesh.common.protocol.SubscriptionMode;
import org.apache.eventmesh.common.protocol.SubscriptionType;
@@ -127,8 +128,8 @@ public class EventMeshGrpcConsumerTest {
}
@Override
- public String getProtocolType() {
- return null;
+ public EventMeshProtocolType getProtocolType() {
+ return EventMeshProtocolType.EVENT_MESH_MESSAGE;
}
});
eventMeshGrpcConsumer.subscribe(Collections.singletonList(buildMockSubscriptionItem()));
diff --git
a/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/grpc/producer/EventMeshGrpcProducerTest.java
b/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/grpc/producer/EventMeshGrpcProducerTest.java
index f2b22e634..e3050a1c4 100644
---
a/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/grpc/producer/EventMeshGrpcProducerTest.java
+++
b/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/grpc/producer/EventMeshGrpcProducerTest.java
@@ -43,10 +43,13 @@ import org.junit.Before;
import org.junit.Test;
import org.junit.runner.RunWith;
import org.mockito.Mock;
+import org.mockito.Mockito;
import org.powermock.core.classloader.annotations.PowerMockIgnore;
import org.powermock.core.classloader.annotations.PrepareForTest;
import org.powermock.modules.junit4.PowerMockRunner;
+import io.cloudevents.CloudEvent;
+
@RunWith(PowerMockRunner.class)
@PrepareForTest(PublisherServiceBlockingStub.class)
@PowerMockIgnore({"javax.management.*", "com.sun.org.apache.xerces.*",
"javax.xml.*", "org.xml.*", "org.w3c.*"})
@@ -58,11 +61,16 @@ public class EventMeshGrpcProducerTest {
@Mock
private PublisherServiceBlockingStub stub;
+ @Mock
+ private EventMeshMessageProducer eventMeshMessageProducer;
+
@Before
public void setUp() throws Exception {
producer = new
EventMeshGrpcProducer(EventMeshGrpcClientConfig.builder().build());
producer.setCloudEventProducer(cloudEventProducer);
+ producer.setEventMeshMessageProducer(eventMeshMessageProducer);
producer.setPublisherClient(stub);
+
doThrow(RuntimeException.class).when(stub).publish(
argThat(argument -> argument != null &&
StringUtils.equals(argument.getContent(),
"mockExceptionContent")));
@@ -70,20 +78,35 @@ public class EventMeshGrpcProducerTest {
argThat(argument -> argument != null &&
StringUtils.equals(argument.getContent(), "mockContent")));
doReturn(Response.getDefaultInstance()).when(stub).batchPublish(
argThat(argument -> argument != null &&
StringUtils.equals(argument.getTopic(), "mockTopic")));
+ doReturn(stub).when(stub).withDeadlineAfter(1000L,
TimeUnit.MILLISECONDS);
doAnswer(invocation -> {
SimpleMessage simpleMessage = invocation.getArgument(0);
if (StringUtils.isEmpty(simpleMessage.getContent())) {
return SimpleMessage.getDefaultInstance();
}
return SimpleMessage.newBuilder(simpleMessage).build();
- }).when(stub).requestReply(any());
- doReturn(stub).when(stub).withDeadlineAfter(1000,
TimeUnit.MILLISECONDS);
+ }).when(stub).requestReply(Mockito.isA(SimpleMessage.class));
when(cloudEventProducer.publish(anyList())).thenReturn(Response.getDefaultInstance());
+
when(cloudEventProducer.publish(Mockito.isA(CloudEvent.class))).thenReturn(Response.getDefaultInstance());
+
when(eventMeshMessageProducer.publish(anyList())).thenReturn(Response.getDefaultInstance());
+
when(eventMeshMessageProducer.publish(Mockito.isA(EventMeshMessage.class))).thenReturn(Response.getDefaultInstance());
+ doAnswer(invocation -> {
+ EventMeshMessage eventMeshMessage = invocation.getArgument(0);
+ if (StringUtils.isEmpty(eventMeshMessage.getContent())) {
+ return null;
+ }
+ return eventMeshMessage;
+ }).when(eventMeshMessageProducer).requestReply(any(),
Mockito.anyLong());
+
}
@Test
public void testPublishWithException() {
-
assertThat(producer.publish(defaultEventMeshMessageBuilder().content("mockExceptionContent").build())).isNull();
+ try {
+
producer.publish(defaultEventMeshMessageBuilder().content("mockExceptionContent").build());
+ } catch (Exception e) {
+ assertThat(e).isNotNull();
+ }
}
@Test
@@ -109,9 +132,9 @@ public class EventMeshGrpcProducerTest {
@Test
public void testRequestReply() {
assertThat(producer.requestReply(defaultEventMeshMessageBuilder().content(StringUtils.EMPTY).build(),
- 1000)).isNull();
+ 1000L)).isNull();
EventMeshMessage eventMeshMessage =
defaultEventMeshMessageBuilder().build();
- assertThat(producer.requestReply(eventMeshMessage,
1000)).hasFieldOrPropertyWithValue("content",
+ assertThat(producer.requestReply(eventMeshMessage,
1000L)).hasFieldOrPropertyWithValue("content",
eventMeshMessage.getContent()).hasFieldOrPropertyWithValue("topic",
eventMeshMessage.getTopic());
}
@@ -120,4 +143,4 @@ public class EventMeshGrpcProducerTest {
.createTime(System.currentTimeMillis()).uniqueId("mockUniqueId").topic("mockTopic")
.prop(Collections.emptyMap());
}
-}
\ No newline at end of file
+}
diff --git
a/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/grpc/producer/EventMeshMessageProducerTest.java
b/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/grpc/producer/EventMeshMessageProducerTest.java
new file mode 100644
index 000000000..fcdbca10e
--- /dev/null
+++
b/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/grpc/producer/EventMeshMessageProducerTest.java
@@ -0,0 +1,120 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * 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.
+ */
+
+package org.apache.eventmesh.client.grpc.producer;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.Mockito.doAnswer;
+import static org.mockito.Mockito.doReturn;
+import static org.mockito.Mockito.when;
+
+import org.apache.eventmesh.client.grpc.config.EventMeshGrpcClientConfig;
+import org.apache.eventmesh.common.EventMeshMessage;
+import org.apache.eventmesh.common.protocol.grpc.protos.BatchMessage;
+import
org.apache.eventmesh.common.protocol.grpc.protos.PublisherServiceGrpc.PublisherServiceBlockingStub;
+import org.apache.eventmesh.common.protocol.grpc.protos.Response;
+import org.apache.eventmesh.common.protocol.grpc.protos.SimpleMessage;
+
+import org.apache.commons.lang3.StringUtils;
+
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.List;
+import java.util.concurrent.TimeUnit;
+
+import org.junit.Before;
+import org.junit.Test;
+import org.junit.runner.RunWith;
+import org.mockito.Mock;
+import org.mockito.Mockito;
+import org.powermock.core.classloader.annotations.PowerMockIgnore;
+import org.powermock.core.classloader.annotations.PrepareForTest;
+import org.powermock.modules.junit4.PowerMockRunner;
+
+@RunWith(PowerMockRunner.class)
+@PrepareForTest({PublisherServiceBlockingStub.class, Response.class})
+@PowerMockIgnore({"javax.management.*", "com.sun.org.apache.xerces.*",
"javax.xml.*", "org.xml.*", "org.w3c.*"})
+public class EventMeshMessageProducerTest {
+
+ private EventMeshMessageProducer eventMeshMessageProducer;
+
+ @Mock
+ private PublisherServiceBlockingStub blockingStub;
+
+ @Mock
+ private Response mockResponse;
+
+ @Mock
+ private SimpleMessage mockSimpleMessage;
+
+ @Before
+ public void setUp() throws Exception {
+ eventMeshMessageProducer = new
EventMeshMessageProducer(EventMeshGrpcClientConfig.builder().build(),
blockingStub);
+
when(blockingStub.batchPublish(Mockito.isA(BatchMessage.class))).thenReturn(mockResponse);
+
when(blockingStub.publish(Mockito.isA(SimpleMessage.class))).thenReturn(mockResponse);
+
when(blockingStub.requestReply(Mockito.isA(SimpleMessage.class))).thenReturn(mockSimpleMessage);
+ doAnswer(invocation -> {
+ SimpleMessage simpleMessage = invocation.getArgument(0);
+ if (StringUtils.isEmpty(simpleMessage.getContent())) {
+ return SimpleMessage.getDefaultInstance();
+ }
+ return SimpleMessage.newBuilder(simpleMessage).build();
+ }).when(blockingStub).requestReply(any());
+ doReturn(blockingStub).when(blockingStub).withDeadlineAfter(1000,
TimeUnit.MILLISECONDS);
+ }
+
+ @Test
+ public void testPublishSingle() {
+
+ EventMeshMessage eventMeshMessage = EventMeshMessage.builder()
+ .topic("mxsm")
+ .content("mxsm")
+ .bizSeqNo("mxsm")
+ .uniqueId("mxsm")
+ .build();
+
assertThat(eventMeshMessageProducer.publish(eventMeshMessage)).isEqualTo(mockResponse);
+ }
+
+ @Test
+ public void testPublishMulti() {
+
+ List<EventMeshMessage> messageArrayList = new ArrayList<>();
+ for (int i = 0; i < 10; ++i) {
+ EventMeshMessage eventMeshMessage = EventMeshMessage.builder()
+ .topic("mxsm")
+ .content("mxsm")
+ .bizSeqNo("mxsm" + i)
+ .uniqueId("mxsm" + i)
+ .build();
+ messageArrayList.add(eventMeshMessage);
+ }
+
assertThat(eventMeshMessageProducer.publish(messageArrayList)).isEqualTo(mockResponse);
+
assertThat(eventMeshMessageProducer.publish(Collections.emptyList())).isNull();
+ }
+
+ @Test
+ public void requestReply() {
+ EventMeshMessage eventMeshMessage = EventMeshMessage.builder()
+ .topic("mxsm")
+ .content("mxsm")
+ .bizSeqNo("mxsm")
+ .uniqueId("mxsm")
+ .build();
+ assertThat(eventMeshMessageProducer.requestReply(eventMeshMessage,
1000L)).isNotNull();
+ }
+}
diff --git
a/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/grpc/util/EventMeshClientUtilTest.java
b/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/grpc/util/EventMeshClientUtilTest.java
index 6dc714a71..b7e621886 100644
---
a/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/grpc/util/EventMeshClientUtilTest.java
+++
b/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/grpc/util/EventMeshClientUtilTest.java
@@ -20,9 +20,9 @@ package org.apache.eventmesh.client.grpc.util;
import static org.assertj.core.api.Assertions.assertThat;
import org.apache.eventmesh.client.grpc.config.EventMeshGrpcClientConfig;
-import org.apache.eventmesh.client.tcp.common.EventMeshCommon;
import org.apache.eventmesh.common.Constants;
import org.apache.eventmesh.common.EventMeshMessage;
+import org.apache.eventmesh.common.enums.EventMeshProtocolType;
import org.apache.eventmesh.common.protocol.grpc.common.ProtocolKey;
import org.apache.eventmesh.common.protocol.grpc.protos.BatchMessage;
import org.apache.eventmesh.common.protocol.grpc.protos.SimpleMessage;
@@ -50,7 +50,7 @@ public class EventMeshClientUtilTest {
@Test
public void testBuildHeader() {
EventMeshGrpcClientConfig clientConfig =
EventMeshGrpcClientConfig.builder().build();
- assertThat(EventMeshClientUtil.buildHeader(clientConfig,
"protocolType")).hasFieldOrPropertyWithValue("env",
+ assertThat(EventMeshClientUtil.buildHeader(clientConfig,
EventMeshProtocolType.CLOUD_EVENTS)).hasFieldOrPropertyWithValue("env",
clientConfig.getEnv()).hasFieldOrPropertyWithValue("idc",
clientConfig.getIdc())
.hasFieldOrPropertyWithValue("ip", IPUtils.getLocalAddress())
.hasFieldOrPropertyWithValue("pid",
Long.toString(ThreadUtils.getPID()))
@@ -58,7 +58,7 @@ public class EventMeshClientUtilTest {
.hasFieldOrPropertyWithValue("language",
clientConfig.getLanguage())
.hasFieldOrPropertyWithValue("username",
clientConfig.getUserName())
.hasFieldOrPropertyWithValue("password",
clientConfig.getPassword())
- .hasFieldOrPropertyWithValue("protocolType", "protocolType")
+ .hasFieldOrPropertyWithValue("protocolType",
EventMeshProtocolType.CLOUD_EVENTS.protocolTypeName())
.hasFieldOrPropertyWithValue("protocolDesc",
Constants.PROTOCOL_GRPC)
.hasFieldOrPropertyWithValue("protocolVersion",
SpecVersion.V1.toString());
}
@@ -76,7 +76,7 @@ public class EventMeshClientUtilTest {
SimpleMessage message =
SimpleMessage.newBuilder().setSeqNum("1").setUniqueId(RandomStringUtils.generateNum(5))
.setTopic("mockTopic")
.setContent("{\"specversion\":\"1.0\",\"id\":\"id\",\"source\":\"source\",\"type\":\"type\"}").build();
- Object buildMessage = EventMeshClientUtil.buildMessage(message,
EventMeshCommon.CLOUD_EVENTS_PROTOCOL_NAME);
+ Object buildMessage = EventMeshClientUtil.buildMessage(message,
EventMeshProtocolType.CLOUD_EVENTS);
assertThat(buildMessage).isInstanceOf(CloudEvent.class);
CloudEvent cloudEvent = (CloudEvent) buildMessage;
assertThat(cloudEvent).isNotNull().hasFieldOrPropertyWithValue("subject",
message.getTopic());
@@ -88,7 +88,7 @@ public class EventMeshClientUtilTest {
public void testBuildMessageWithDefaultProto() {
SimpleMessage message =
SimpleMessage.newBuilder().setSeqNum("1").setUniqueId(RandomStringUtils.generateNum(5))
.setTopic("mockTopic").setContent("mockContent").build();
- Object buildMessage = EventMeshClientUtil.buildMessage(message, null);
+ Object buildMessage = EventMeshClientUtil.buildMessage(message,
EventMeshProtocolType.EVENT_MESH_MESSAGE);
assertThat(buildMessage).isInstanceOf(EventMeshMessage.class)
.hasFieldOrPropertyWithValue("content", message.getContent())
.hasFieldOrPropertyWithValue("topic", message.getTopic())
@@ -103,8 +103,8 @@ public class EventMeshClientUtilTest {
.withExtension(ProtocolKey.UNIQUE_ID, "uniqueId").build();
EventMeshGrpcClientConfig clientConfig =
EventMeshGrpcClientConfig.builder().build();
assertThat(EventMeshClientUtil.buildSimpleMessage(cloudEvent,
clientConfig,
-
EventMeshCommon.CLOUD_EVENTS_PROTOCOL_NAME)).hasFieldOrPropertyWithValue("header",
- EventMeshClientUtil.buildHeader(clientConfig,
EventMeshCommon.CLOUD_EVENTS_PROTOCOL_NAME))
+
EventMeshProtocolType.CLOUD_EVENTS)).hasFieldOrPropertyWithValue("header",
+ EventMeshClientUtil.buildHeader(clientConfig,
EventMeshProtocolType.CLOUD_EVENTS))
.hasFieldOrPropertyWithValue("producerGroup",
clientConfig.getProducerGroup())
.hasFieldOrPropertyWithValue("topic", cloudEvent.getSubject())
.hasFieldOrPropertyWithValue("ttl", "4000")
@@ -121,8 +121,8 @@ public class EventMeshClientUtilTest {
.uniqueId("mockUniqueId").bizSeqNo("mockBizSeqNo").build();
EventMeshGrpcClientConfig clientConfig =
EventMeshGrpcClientConfig.builder().build();
assertThat(
- EventMeshClientUtil.buildSimpleMessage(eventMeshMessage,
clientConfig, "")).hasFieldOrPropertyWithValue(
- "header", EventMeshClientUtil.buildHeader(clientConfig, ""))
+ EventMeshClientUtil.buildSimpleMessage(eventMeshMessage,
clientConfig, EventMeshProtocolType.EVENT_MESH_MESSAGE))
+ .hasFieldOrPropertyWithValue("header",
EventMeshClientUtil.buildHeader(clientConfig,
EventMeshProtocolType.EVENT_MESH_MESSAGE))
.hasFieldOrPropertyWithValue("producerGroup",
clientConfig.getProducerGroup())
.hasFieldOrPropertyWithValue("topic", eventMeshMessage.getTopic())
.hasFieldOrPropertyWithValue("ttl", "4000")
@@ -139,9 +139,9 @@ public class EventMeshClientUtilTest {
.withExtension(ProtocolKey.UNIQUE_ID, "uniqueId").build());
EventMeshGrpcClientConfig clientConfig =
EventMeshGrpcClientConfig.builder().build();
BatchMessage batchMessage =
EventMeshClientUtil.buildBatchMessages(cloudEvents, clientConfig,
- EventMeshCommon.CLOUD_EVENTS_PROTOCOL_NAME);
+ EventMeshProtocolType.CLOUD_EVENTS);
assertThat(batchMessage).hasFieldOrPropertyWithValue("header",
- EventMeshClientUtil.buildHeader(clientConfig,
EventMeshCommon.CLOUD_EVENTS_PROTOCOL_NAME))
+ EventMeshClientUtil.buildHeader(clientConfig,
EventMeshProtocolType.CLOUD_EVENTS))
.hasFieldOrPropertyWithValue("topic",
cloudEvents.get(0).getSubject())
.hasFieldOrPropertyWithValue("producerGroup",
clientConfig.getProducerGroup());
assertThat(batchMessage.getMessageItemList()).hasSize(1).first().hasFieldOrPropertyWithValue("content",
@@ -161,9 +161,9 @@ public class EventMeshClientUtilTest {
.bizSeqNo("mockBizSeqNo")
.prop(Collections.singletonMap(Constants.EVENTMESH_MESSAGE_CONST_TTL,
"4000")).build());
EventMeshGrpcClientConfig clientConfig =
EventMeshGrpcClientConfig.builder().build();
- BatchMessage batchMessage =
EventMeshClientUtil.buildBatchMessages(eventMeshMessages, clientConfig, "");
+ BatchMessage batchMessage =
EventMeshClientUtil.buildBatchMessages(eventMeshMessages, clientConfig,
EventMeshProtocolType.EVENT_MESH_MESSAGE);
assertThat(batchMessage).hasFieldOrPropertyWithValue("header",
- EventMeshClientUtil.buildHeader(clientConfig, ""))
+ EventMeshClientUtil.buildHeader(clientConfig,
EventMeshProtocolType.EVENT_MESH_MESSAGE))
.hasFieldOrPropertyWithValue("topic",
eventMeshMessages.get(0).getTopic())
.hasFieldOrPropertyWithValue("producerGroup",
clientConfig.getProducerGroup());
EventMeshMessage firstMeshMessage = eventMeshMessages.get(0);
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]