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]

Reply via email to