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

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


The following commit(s) were added to refs/heads/master by this push:
     new e52e41557 [ISSUE #2687]change go server project structure (#2688)
e52e41557 is described below

commit e52e41557b34d314f3f621d31e9a7c26da7106e0
Author: walleliu <[email protected]>
AuthorDate: Wed Dec 28 09:05:58 2022 +0800

    [ISSUE #2687]change go server project structure (#2688)
    
    * change api to interface, used to generate mock code
    
    * add mock code
    
    * add config test
    
    * add config test
    
    * change the project structure in different directory
    
    * add mock consumer service
    
    * add consumer test case
---
 eventmesh-server-go/config/config_test.go          |  96 ++++++
 .../config/testdata/test_config.yaml               |  74 +++++
 eventmesh-server-go/configs/eventmesh-server.yaml  |  19 ++
 .../grpc/producer_group.go => consts/grpc.go}      |  11 +-
 .../protocol/grpc/{ => consumer}/consumer_group.go |  12 +-
 .../grpc/{ => consumer}/consumer_group_client.go   |   8 +-
 .../grpc/{ => consumer}/consumer_group_option.go   |  27 +-
 .../grpc/consumer/consumer_group_option_test.go    | 116 ++++++++
 .../grpc/{ => consumer}/consumer_manager.go        |  51 ++--
 .../grpc/consumer/consumer_manager_test.go         | 196 +++++++++++++
 .../protocol/grpc/{ => consumer}/consumer_mesh.go  |  48 +--
 .../protocol/grpc/consumer/consumer_processor.go   | 321 +++++++++++++++++++++
 .../grpc/{ => consumer}/consumer_service.go        |  42 +--
 .../grpc/consumer/consumer_service_test.go         | 123 ++++++++
 .../{emitter.go => consumer/message_context.go}    |  21 +-
 .../grpc/{ => consumer}/message_handler.go         |  23 +-
 .../{request.go => consumer/message_request.go}    |  39 ++-
 .../grpc/consumer/mocks/consumer_group_option.go   | 193 +++++++++++++
 .../grpc/consumer/mocks/consumer_manager.go        | 132 +++++++++
 .../core/protocol/grpc/consumer_service_test.go    | 154 ----------
 .../runtime/core/protocol/grpc/context.go          | 101 -------
 .../core/protocol/grpc/{ => emitter}/emitter.go    |  17 +-
 .../core/protocol/grpc/emitter/mocks/emitter.go    |  50 ++++
 .../runtime/core/protocol/grpc/fake_client.go      | 113 --------
 .../runtime/core/protocol/grpc/generate_mocks.sh   |  12 +
 .../protocol/grpc/heartbeat/heartbeat_processor.go |  61 ++++
 .../grpc/{ => heartbeat}/heartbeat_service.go      |  15 +-
 .../core/protocol/grpc/mocks/consumer_service.go   |  92 ++++++
 .../message_context.go}                            |  17 +-
 .../protocol/grpc/{ => producer}/producer_group.go |   2 +-
 .../grpc/{ => producer}/producer_manager.go        |  31 +-
 .../protocol/grpc/{ => producer}/producer_mesh.go  |  32 +-
 .../producer_processor.go}                         | 252 +++-------------
 .../grpc/{ => producer}/producer_service.go        |  18 +-
 .../grpc/{ => producer}/producer_service_test.go   |   2 +-
 .../core/protocol/grpc/{ => retry}/retry.go        |   8 +-
 .../protocol/grpc/{ => validator}/validator.go     |   7 +-
 eventmesh-server-go/runtime/emserver/grpc.go       |  22 +-
 eventmesh-server-go/runtime/emserver/grpc_test.go  |  57 +---
 .../runtime/emserver/mocks/mock_graceful.go        |  62 ++++
 40 files changed, 1872 insertions(+), 805 deletions(-)

diff --git a/eventmesh-server-go/config/config_test.go 
b/eventmesh-server-go/config/config_test.go
new file mode 100644
index 000000000..96e8a95bc
--- /dev/null
+++ b/eventmesh-server-go/config/config_test.go
@@ -0,0 +1,96 @@
+// 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 config
+
+import (
+       "github.com/apache/incubator-eventmesh/eventmesh-server-go/plugin"
+       testifyassert "github.com/stretchr/testify/assert"
+       "gopkg.in/yaml.v3"
+       "os"
+       "testing"
+       "time"
+)
+
+func TestConfig_Load(t *testing.T) {
+       assert := testifyassert.New(t)
+
+       config := &Config{}
+       config.Server.GRPCOption = &GRPCOption{
+               Port: "10010",
+               TLSOption: &TLSOption{
+                       EnableInsecure: false,
+                       CA:             "",
+                       Certfile:       "",
+                       Keyfile:        "",
+               },
+               PProfOption: &PProfOption{
+                       Port: "10011",
+               },
+               SendPoolSize:          10,
+               SubscribePoolSize:     10,
+               RetryPoolSize:         10,
+               PushMessagePoolSize:   10,
+               ReplyPoolSize:         10,
+               MsgReqNumPerSecond:    5,
+               RegistryName:          "test",
+               Cluster:               "test",
+               Env:                   "env",
+               IDC:                   "idc1",
+               SessionExpiredInMills: 5 * time.Second,
+               SendMessageTimeout:    5 * time.Second,
+       }
+       config.Server.HTTPOption = &HTTPOption{
+               Port: "10010",
+               TLSOption: &TLSOption{
+                       EnableInsecure: false,
+                       CA:             "",
+                       Certfile:       "",
+                       Keyfile:        "",
+               },
+               PProfOption: &PProfOption{
+                       Port: "10011",
+               },
+       }
+       config.Server.TCPOption = &TCPOption{
+               Port: "10010",
+               TLSOption: &TLSOption{
+                       EnableInsecure: false,
+                       CA:             "",
+                       Certfile:       "",
+                       Keyfile:        "",
+               },
+               PProfOption: &PProfOption{
+                       Port: "10011",
+               },
+               Multicore: false,
+       }
+       config.ActivePlugins = map[string]string{
+               "registry":  "nacos",
+               "connector": "rocketmq",
+               "log":       "default",
+       }
+       config.Plugins = plugin.Config{}
+
+       configYAML := &Config{}
+       contentYAML, err := os.ReadFile("./testdata/test_config.yaml")
+       assert.NoError(err)
+       if err := yaml.Unmarshal(contentYAML, &configYAML); err != nil {
+               t.Fatal(err)
+       }
+       configYAML.Plugins = plugin.Config{}
+
+       assert.EqualValues(configYAML, config)
+}
diff --git a/eventmesh-server-go/config/testdata/test_config.yaml 
b/eventmesh-server-go/config/testdata/test_config.yaml
new file mode 100644
index 000000000..e4b52e588
--- /dev/null
+++ b/eventmesh-server-go/config/testdata/test_config.yaml
@@ -0,0 +1,74 @@
+server:
+  grpc:
+    port: 10010
+    tls:
+      enable-secure: false
+      ca: ""
+      certfile: ""
+      keyfile: ""
+    pprof:
+      port: 10011
+    send-pool-size: 10
+    subscribe-pool-size: 10
+    retry-pool-size: 10
+    push-message-pool-size: 10
+    reply-pool-size: 10
+    msg-req-num-per-second: 5
+    cluster: "test"
+    registry-name: "test"
+    env: "env"
+    idc: "idc1"
+    session-expired-in-mills: 5s
+    send-message-timeout: 5s
+  http:
+    port: 10010
+    tls:
+      enable-secure: false
+      ca: ""
+      certfile: ""
+      keyfile: ""
+    pprof:
+      port: 10011
+  tcp:
+    port: 10010
+    tls:
+      enable-secure: false
+      ca: ""
+      certfile: ""
+      keyfile: ""
+    pprof:
+      port: 10011
+active-plugins:
+  registry: nacos
+  connector: rocketmq
+  log: default
+plugins:
+  registry:
+    nacos:
+      address_list: your-nacos-addr
+      cache-dir: your-nacos-cache-dir
+  connector:
+    standalone:
+    rocketmq:
+      access_points: your-nameserver-addr
+      namespace: grpcnamespace
+      instance_name: grpcinstance
+      group_name: grpcgroupname
+      send_msg_timeout: 5000
+      producer_retry_times: 3
+      compress_msg_body_threshold: 1024
+      consumer_group: grpcconsumergroup
+      max_reconsume_times: 3
+      message_model: CLUSTERING
+  log:
+    default:
+      - writer: console
+        level: debug
+      - writer: file
+        level: info
+        writer_config:
+          filename: ./eventmesh.log
+          max_size: 10
+          max_backups: 10
+          max_age: 7
+          compress: false
\ No newline at end of file
diff --git a/eventmesh-server-go/configs/eventmesh-server.yaml 
b/eventmesh-server-go/configs/eventmesh-server.yaml
index 924a01cfd..2a112bd36 100644
--- a/eventmesh-server-go/configs/eventmesh-server.yaml
+++ b/eventmesh-server-go/configs/eventmesh-server.yaml
@@ -31,9 +31,28 @@ server:
     reply-pool-size: 10
     msg-req-num-per-second: 5
     cluster: "test"
+    env: "env"
     idc: "idc1"
     session-expired-in-mills: 5s
     send-message-timeout: 5s
+  http:
+    port: 10010
+    tls:
+      enable-secure: false
+      ca: ""
+      certfile: ""
+      keyfile: ""
+    pprof:
+      port: 10011
+  tcp:
+    port: 10010
+    tls:
+      enable-secure: false
+      ca: ""
+      certfile: ""
+      keyfile: ""
+    pprof:
+      port: 10011
 active-plugins:
   registry: nacos
   connector: rocketmq
diff --git a/eventmesh-server-go/runtime/core/protocol/grpc/producer_group.go 
b/eventmesh-server-go/runtime/consts/grpc.go
similarity index 88%
copy from eventmesh-server-go/runtime/core/protocol/grpc/producer_group.go
copy to eventmesh-server-go/runtime/consts/grpc.go
index 519840507..9d8fb031f 100644
--- a/eventmesh-server-go/runtime/core/protocol/grpc/producer_group.go
+++ b/eventmesh-server-go/runtime/consts/grpc.go
@@ -13,8 +13,11 @@
 // See the License for the specific language governing permissions and
 // limitations under the License.
 
-package grpc
+package consts
 
-type ProducerGroupConfig struct {
-       GroupName string `json:"groupName"`
-}
+type GRPCType string
+
+const (
+       WEBHOOK GRPCType = "WEBHOOK"
+       STREAM  GRPCType = "STREAM"
+)
diff --git a/eventmesh-server-go/runtime/core/protocol/grpc/consumer_group.go 
b/eventmesh-server-go/runtime/core/protocol/grpc/consumer/consumer_group.go
similarity index 94%
rename from eventmesh-server-go/runtime/core/protocol/grpc/consumer_group.go
rename to 
eventmesh-server-go/runtime/core/protocol/grpc/consumer/consumer_group.go
index fb8f891e8..cec44678e 100644
--- a/eventmesh-server-go/runtime/core/protocol/grpc/consumer_group.go
+++ b/eventmesh-server-go/runtime/core/protocol/grpc/consumer/consumer_group.go
@@ -13,9 +13,10 @@
 // See the License for the specific language governing permissions and
 // limitations under the License.
 
-package grpc
+package consumer
 
 import (
+       
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/consts"
        
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/proto/pb"
        "github.com/liyue201/gostl/ds/set"
        "sync"
@@ -28,13 +29,6 @@ type ConsumerGroupConfig struct {
        ConsumerGroupTopicConfigs *sync.Map
 }
 
-type GRPCType string
-
-const (
-       WEBHOOK GRPCType = "WEBHOOK"
-       STREAM  GRPCType = "STREAM"
-)
-
 type StateAction string
 
 const (
@@ -47,7 +41,7 @@ type ConsumerGroupTopicConfig struct {
        ConsumerGroup    string
        Topic            string
        SubscriptionMode pb.Subscription_SubscriptionItem_SubscriptionMode
-       GRPCType         GRPCType
+       GRPCType         consts.GRPCType
        // IDCWebhookURLs webhook urls seperated by IDC
        // key is IDC, value is vector.Vector
        IDCWebhookURLs *sync.Map
diff --git 
a/eventmesh-server-go/runtime/core/protocol/grpc/consumer_group_client.go 
b/eventmesh-server-go/runtime/core/protocol/grpc/consumer/consumer_group_client.go
similarity index 83%
rename from 
eventmesh-server-go/runtime/core/protocol/grpc/consumer_group_client.go
rename to 
eventmesh-server-go/runtime/core/protocol/grpc/consumer/consumer_group_client.go
index 4b03b54fe..48eb084bd 100644
--- a/eventmesh-server-go/runtime/core/protocol/grpc/consumer_group_client.go
+++ 
b/eventmesh-server-go/runtime/core/protocol/grpc/consumer/consumer_group_client.go
@@ -13,9 +13,11 @@
 // See the License for the specific language governing permissions and
 // limitations under the License.
 
-package grpc
+package consumer
 
 import (
+       
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/consts"
+       
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/core/protocol/grpc/emitter"
        
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/proto/pb"
        "time"
 )
@@ -26,7 +28,7 @@ type GroupClient struct {
        IDC              string
        ConsumerGroup    string
        Topic            string
-       GRPCType         GRPCType
+       GRPCType         consts.GRPCType
        URL              string
        SubscriptionMode pb.Subscription_SubscriptionItem_SubscriptionMode
        SYS              string
@@ -35,5 +37,5 @@ type GroupClient struct {
        Hostname         string
        APIVersion       string
        LastUPTime       time.Time
-       Emiter           *EventEmitter
+       Emiter           emitter.EventEmitter
 }
diff --git 
a/eventmesh-server-go/runtime/core/protocol/grpc/consumer_group_option.go 
b/eventmesh-server-go/runtime/core/protocol/grpc/consumer/consumer_group_option.go
similarity index 88%
rename from 
eventmesh-server-go/runtime/core/protocol/grpc/consumer_group_option.go
rename to 
eventmesh-server-go/runtime/core/protocol/grpc/consumer/consumer_group_option.go
index 339126e66..77d9d8ab4 100644
--- a/eventmesh-server-go/runtime/core/protocol/grpc/consumer_group_option.go
+++ 
b/eventmesh-server-go/runtime/core/protocol/grpc/consumer/consumer_group_option.go
@@ -13,11 +13,13 @@
 // See the License for the specific language governing permissions and
 // limitations under the License.
 
-package grpc
+package consumer
 
 import (
        "fmt"
        "github.com/apache/incubator-eventmesh/eventmesh-server-go/log"
+       
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/consts"
+       
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/core/protocol/grpc/emitter"
        
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/proto/pb"
        "github.com/liyue201/gostl/ds/set"
        "sync"
@@ -26,11 +28,13 @@ import (
 type RegisterClient func(*GroupClient)
 type DeregisterClient func(*GroupClient)
 
+//go:generate mockgen -destination ./mocks/consumer_group_option.go -package 
mocks 
github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/core/protocol/grpc/consumer
 ConsumerGroupTopicOption
+
 type ConsumerGroupTopicOption interface {
        ConsumerGroup() string
        Topic() string
        SubscriptionMode() pb.Subscription_SubscriptionItem_SubscriptionMode
-       GRPCType() GRPCType
+       GRPCType() consts.GRPCType
        RegisterClient() RegisterClient
        DeregisterClient() DeregisterClient
        IDCURLs() *sync.Map
@@ -45,13 +49,16 @@ type BaseConsumerGroupTopicOption struct {
        consumerGroup    string
        topic            string
        subscriptionMode pb.Subscription_SubscriptionItem_SubscriptionMode
-       gRPCType         GRPCType
+       gRPCType         consts.GRPCType
        registerClient   RegisterClient
        deregisterClient DeregisterClient
 }
 
-func NewConsumerGroupTopicOption(cg string, topic string, mode 
pb.Subscription_SubscriptionItem_SubscriptionMode, gtype GRPCType) 
ConsumerGroupTopicOption {
-       if gtype == WEBHOOK {
+func NewConsumerGroupTopicOption(
+       cg string, topic string,
+       mode pb.Subscription_SubscriptionItem_SubscriptionMode,
+       gtype consts.GRPCType) ConsumerGroupTopicOption {
+       if gtype == consts.WEBHOOK {
                return NewWebhookGroupTopicOption(cg, topic, mode, gtype)
        }
        return NewWStreamGroupTopicOption(cg, topic, mode, gtype)
@@ -78,7 +85,7 @@ func (b *BaseConsumerGroupTopicOption) Topic() string {
 func (b *BaseConsumerGroupTopicOption) SubscriptionMode() 
pb.Subscription_SubscriptionItem_SubscriptionMode {
        return b.subscriptionMode
 }
-func (b *BaseConsumerGroupTopicOption) GRPCType() GRPCType {
+func (b *BaseConsumerGroupTopicOption) GRPCType() consts.GRPCType {
        return b.gRPCType
 }
 func (b *BaseConsumerGroupTopicOption) RegisterClient() RegisterClient {
@@ -111,7 +118,7 @@ func (b *WebhookGroupTopicOption) Size() int {
 func NewWebhookGroupTopicOption(cg string,
        topic string,
        mode pb.Subscription_SubscriptionItem_SubscriptionMode,
-       gtype GRPCType) ConsumerGroupTopicOption {
+       gtype consts.GRPCType) ConsumerGroupTopicOption {
        opt := &WebhookGroupTopicOption{
                BaseConsumerGroupTopicOption: &BaseConsumerGroupTopicOption{
                        consumerGroup:    cg,
@@ -123,7 +130,7 @@ func NewWebhookGroupTopicOption(cg string,
                allURLs:        set.New(set.WithGoroutineSafe()),
        }
        opt.BaseConsumerGroupTopicOption.registerClient = func(cli 
*GroupClient) {
-               if cli.GRPCType != WEBHOOK {
+               if cli.GRPCType != consts.WEBHOOK {
                        log.Warnf("invalid grpc type:%v, with provide WEBHOOK", 
cli.GRPCType)
                        return
                }
@@ -195,7 +202,7 @@ func (b *StreamGroupTopicOption) buildIdcEmitter() {
                e1 := value.(*sync.Map)
                elist := set.New(set.WithGoroutineSafe())
                e1.Range(func(k1, v1 interface{}) bool {
-                       elist.Insert(v1.(*EventEmitter))
+                       elist.Insert(v1.(emitter.EventEmitter))
                        return true
                })
                newIDCEmiters.Store(key, elist)
@@ -223,7 +230,7 @@ func uniqClient(ip, pid string) string {
 func NewWStreamGroupTopicOption(cg string,
        topic string,
        mode pb.Subscription_SubscriptionItem_SubscriptionMode,
-       gtype GRPCType) ConsumerGroupTopicOption {
+       gtype consts.GRPCType) ConsumerGroupTopicOption {
        opt := &StreamGroupTopicOption{
                BaseConsumerGroupTopicOption: &BaseConsumerGroupTopicOption{
                        consumerGroup:    cg,
diff --git 
a/eventmesh-server-go/runtime/core/protocol/grpc/consumer/consumer_group_option_test.go
 
b/eventmesh-server-go/runtime/core/protocol/grpc/consumer/consumer_group_option_test.go
new file mode 100644
index 000000000..96e1a147f
--- /dev/null
+++ 
b/eventmesh-server-go/runtime/core/protocol/grpc/consumer/consumer_group_option_test.go
@@ -0,0 +1,116 @@
+// 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 consumer
+
+import (
+       
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/consts"
+       
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/proto/pb"
+       "github.com/stretchr/testify/assert"
+       "testing"
+)
+
+func Test_ConsumerGroupTopicOption(t *testing.T) {
+       tests := []struct {
+               name          string
+               consumerGroup string
+               topic         string
+               mode          pb.Subscription_SubscriptionItem_SubscriptionMode
+               grpcType      consts.GRPCType
+               expect        func(t *testing.T, option 
ConsumerGroupTopicOption)
+       }{
+               {
+                       name:          "create grpc stream broadcasting option",
+                       consumerGroup: "consumergroup",
+                       topic:         "topic",
+                       mode:          
pb.Subscription_SubscriptionItem_BROADCASTING,
+                       grpcType:      consts.STREAM,
+                       expect: func(t *testing.T, option 
ConsumerGroupTopicOption) {
+                               assert.NotNil(t, option)
+                               assert.Equal(t, option.Topic(), "topic")
+                               assert.Equal(t, option.GRPCType(), 
consts.STREAM)
+                               assert.Equal(t, option.ConsumerGroup(), 
"consumergroup")
+                               assert.Equal(t, option.SubscriptionMode(), 
pb.Subscription_SubscriptionItem_BROADCASTING)
+                               assert.NotNil(t, option.RegisterClient())
+                               assert.NotNil(t, option.DeregisterClient())
+                               assert.NotNil(t, option.IDCEmiters())
+                               assert.NotNil(t, option.AllEmiters())
+                       },
+               },
+               {
+                       name:          "create grpc stream clustering option",
+                       consumerGroup: "consumergroup",
+                       topic:         "topic",
+                       mode:          
pb.Subscription_SubscriptionItem_CLUSTERING,
+                       grpcType:      consts.STREAM,
+                       expect: func(t *testing.T, option 
ConsumerGroupTopicOption) {
+                               assert.NotNil(t, option)
+                               assert.Equal(t, option.Topic(), "topic")
+                               assert.Equal(t, option.GRPCType(), 
consts.STREAM)
+                               assert.Equal(t, option.ConsumerGroup(), 
"consumergroup")
+                               assert.Equal(t, option.SubscriptionMode(), 
pb.Subscription_SubscriptionItem_CLUSTERING)
+                               assert.NotNil(t, option.RegisterClient())
+                               assert.NotNil(t, option.DeregisterClient())
+                               assert.NotNil(t, option.IDCEmiters())
+                               assert.NotNil(t, option.AllEmiters())
+                       },
+               },
+               {
+                       name:          "create webhook clustering option",
+                       consumerGroup: "consumergroup",
+                       topic:         "topic",
+                       mode:          
pb.Subscription_SubscriptionItem_CLUSTERING,
+                       grpcType:      consts.WEBHOOK,
+                       expect: func(t *testing.T, option 
ConsumerGroupTopicOption) {
+                               assert.NotNil(t, option)
+                               assert.Equal(t, option.Topic(), "topic")
+                               assert.Equal(t, option.GRPCType(), 
consts.WEBHOOK)
+                               assert.Equal(t, option.ConsumerGroup(), 
"consumergroup")
+                               assert.Equal(t, option.SubscriptionMode(), 
pb.Subscription_SubscriptionItem_CLUSTERING)
+                               assert.NotNil(t, option.RegisterClient())
+                               assert.NotNil(t, option.DeregisterClient())
+                               assert.NotNil(t, option.IDCURLs())
+                               assert.NotNil(t, option.AllURLs())
+                               assert.Equal(t, option.Size(), 0)
+                       },
+               },
+               {
+                       name:          "create webhook broadcasting option",
+                       consumerGroup: "consumergroup",
+                       topic:         "topic",
+                       mode:          
pb.Subscription_SubscriptionItem_BROADCASTING,
+                       grpcType:      consts.WEBHOOK,
+                       expect: func(t *testing.T, option 
ConsumerGroupTopicOption) {
+                               assert.NotNil(t, option)
+                               assert.Equal(t, option.Topic(), "topic")
+                               assert.Equal(t, option.GRPCType(), 
consts.WEBHOOK)
+                               assert.Equal(t, option.ConsumerGroup(), 
"consumergroup")
+                               assert.Equal(t, option.SubscriptionMode(), 
pb.Subscription_SubscriptionItem_BROADCASTING)
+                               assert.NotNil(t, option.RegisterClient())
+                               assert.NotNil(t, option.DeregisterClient())
+                               assert.NotNil(t, option.IDCURLs())
+                               assert.NotNil(t, option.AllURLs())
+                               assert.Equal(t, option.Size(), 0)
+                       },
+               },
+       }
+
+       for _, tc := range tests {
+               t.Run(tc.name, func(t *testing.T) {
+                       option := NewConsumerGroupTopicOption(tc.consumerGroup, 
tc.topic, tc.mode, tc.grpcType)
+                       tc.expect(t, option)
+               })
+       }
+}
diff --git a/eventmesh-server-go/runtime/core/protocol/grpc/consumer_manager.go 
b/eventmesh-server-go/runtime/core/protocol/grpc/consumer/consumer_manager.go
similarity index 78%
rename from eventmesh-server-go/runtime/core/protocol/grpc/consumer_manager.go
rename to 
eventmesh-server-go/runtime/core/protocol/grpc/consumer/consumer_manager.go
index 80fc5666c..c06f1383c 100644
--- a/eventmesh-server-go/runtime/core/protocol/grpc/consumer_manager.go
+++ 
b/eventmesh-server-go/runtime/core/protocol/grpc/consumer/consumer_manager.go
@@ -13,7 +13,7 @@
 // See the License for the specific language governing permissions and
 // limitations under the License.
 
-package grpc
+package consumer
 
 import (
        config2 
"github.com/apache/incubator-eventmesh/eventmesh-server-go/config"
@@ -29,7 +29,18 @@ var (
        ErrNoConsumerClient = errors.New("no consumer group client")
 )
 
-type ConsumerManager struct {
+//go:generate mockgen -destination ./mocks/consumer_manager.go -package mocks 
github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/core/protocol/grpc/consumer
 ConsumerManager
+type ConsumerManager interface {
+       GetConsumer(consumerGroup string) (EventMeshConsumer, error)
+       RegisterClient(cli *GroupClient) error
+       DeRegisterClient(cli *GroupClient) error
+       UpdateClientTime(cli *GroupClient)
+       RestartConsumer(consumerGroup string) error
+       Start() error
+       Stop() error
+}
+
+type consumerManager struct {
        // consumerClients store all consumer clients
        // key is consumer group, value is set *GroupClient
        consumerGroupClients *sync.Map
@@ -40,17 +51,17 @@ type ConsumerManager struct {
 }
 
 // NewConsumerManager create new consumer manager
-func NewConsumerManager() (*ConsumerManager, error) {
-       return &ConsumerManager{
+func NewConsumerManager() (ConsumerManager, error) {
+       return &consumerManager{
                consumers:            new(sync.Map),
                consumerGroupClients: new(sync.Map),
        }, nil
 }
 
-func (c *ConsumerManager) GetConsumer(consumerGroup string) 
(*EventMeshConsumer, error) {
+func (c *consumerManager) GetConsumer(consumerGroup string) 
(EventMeshConsumer, error) {
        val, ok := c.consumers.Load(consumerGroup)
        if ok {
-               return val.(*EventMeshConsumer), nil
+               return val.(EventMeshConsumer), nil
        }
        cu, err := NewEventMeshConsumer(consumerGroup)
        if err != nil {
@@ -60,7 +71,7 @@ func (c *ConsumerManager) GetConsumer(consumerGroup string) 
(*EventMeshConsumer,
        return cu, nil
 }
 
-func (c *ConsumerManager) RegisterClient(cli *GroupClient) error {
+func (c *consumerManager) RegisterClient(cli *GroupClient) error {
        val, ok := c.consumerGroupClients.Load(cli.ConsumerGroup)
        if !ok {
                cliset := set.New(set.WithGoroutineSafe())
@@ -72,13 +83,13 @@ func (c *ConsumerManager) RegisterClient(cli *GroupClient) 
error {
        found := false
        for iter := localClients.Begin(); iter.IsValid(); iter.Next() {
                lc := iter.Value().(*GroupClient)
-               if lc.GRPCType == WEBHOOK {
+               if lc.GRPCType == consts.WEBHOOK {
                        lc.URL = cli.URL
                        lc.LastUPTime = cli.LastUPTime
                        found = true
                        break
                }
-               if lc.GRPCType == STREAM {
+               if lc.GRPCType == consts.STREAM {
                        lc.Emiter = cli.Emiter
                        lc.LastUPTime = cli.LastUPTime
                        found = true
@@ -91,7 +102,7 @@ func (c *ConsumerManager) RegisterClient(cli *GroupClient) 
error {
        return nil
 }
 
-func (c *ConsumerManager) DeRegisterClient(cli *GroupClient) error {
+func (c *consumerManager) DeRegisterClient(cli *GroupClient) error {
        val, ok := c.consumerGroupClients.Load(cli.ConsumerGroup)
        if !ok {
                log.Debugf("no consumer group client found, name:%v", 
cli.ConsumerGroup)
@@ -101,7 +112,7 @@ func (c *ConsumerManager) DeRegisterClient(cli 
*GroupClient) error {
        for iter := localClients.Begin(); iter.IsValid(); iter.Next() {
                lc := iter.Value().(*GroupClient)
                if lc.Topic == cli.Topic {
-                       if lc.GRPCType == STREAM {
+                       if lc.GRPCType == consts.STREAM {
                                // TODO
                                // close the GRPC client stream before removing 
it
                        }
@@ -114,13 +125,13 @@ func (c *ConsumerManager) DeRegisterClient(cli 
*GroupClient) error {
        return nil
 }
 
-func (c *ConsumerManager) restartConsumer(consumerGroup string) error {
+func (c *consumerManager) RestartConsumer(consumerGroup string) error {
        val, ok := c.consumers.Load(consumerGroup)
        if !ok {
                return nil
        }
-       emconsumer := val.(*EventMeshConsumer)
-       if emconsumer.ServiceState == consts.RUNNING {
+       emconsumer := val.(EventMeshConsumer)
+       if emconsumer.ServiceState() == consts.RUNNING {
                if err := emconsumer.Shutdown(); err != nil {
                        return err
                }
@@ -131,14 +142,14 @@ func (c *ConsumerManager) restartConsumer(consumerGroup 
string) error {
        if err := emconsumer.Start(); err != nil {
                return err
        }
-       if emconsumer.ServiceState != consts.RUNNING {
+       if emconsumer.ServiceState() != consts.RUNNING {
                log.Warnf("restart eventmesh consumer failed, status:%v", 
emconsumer.ServiceState)
                c.consumers.Delete(consumerGroup)
        }
        return nil
 }
 
-func (c *ConsumerManager) UpdateClientTime(cli *GroupClient) {
+func (c *consumerManager) UpdateClientTime(cli *GroupClient) {
        val, ok := c.consumerGroupClients.Load(cli.ConsumerGroup)
        if !ok {
                log.Debugf("no consumer group client found, name:%v", 
cli.ConsumerGroup)
@@ -150,7 +161,7 @@ func (c *ConsumerManager) UpdateClientTime(cli 
*GroupClient) {
        }
 }
 
-func (c *ConsumerManager) clientCheck() {
+func (c *consumerManager) clientCheck() {
        sessionExpiredInMills := 
config2.GlobalConfig().Server.GRPCOption.SessionExpiredInMills
        tk := time.NewTicker(sessionExpiredInMills)
        go func() {
@@ -180,7 +191,7 @@ func (c *ConsumerManager) clientCheck() {
                                        }
                                }
                                for _, rs := range consumerGroupRestart {
-                                       if err := c.restartConsumer(rs); err != 
nil {
+                                       if err := c.RestartConsumer(rs); err != 
nil {
                                                log.Warnf("deregistry 
consumer:%v  err:%v", rs, err)
                                                return true
                                        }
@@ -191,11 +202,11 @@ func (c *ConsumerManager) clientCheck() {
        }()
 }
 
-func (c *ConsumerManager) Start() error {
+func (c *consumerManager) Start() error {
        log.Infof("start consumer manager")
        return nil
 }
 
-func (c *ConsumerManager) Stop() error {
+func (c *consumerManager) Stop() error {
        return nil
 }
diff --git 
a/eventmesh-server-go/runtime/core/protocol/grpc/consumer/consumer_manager_test.go
 
b/eventmesh-server-go/runtime/core/protocol/grpc/consumer/consumer_manager_test.go
new file mode 100644
index 000000000..9aafa54d1
--- /dev/null
+++ 
b/eventmesh-server-go/runtime/core/protocol/grpc/consumer/consumer_manager_test.go
@@ -0,0 +1,196 @@
+// 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 consumer
+
+import (
+       "fmt"
+       
"github.com/apache/incubator-eventmesh/eventmesh-server-go/pkg/common/protocol/grpc"
+       "github.com/golang/mock/gomock"
+       "testing"
+       "time"
+
+       "github.com/liyue201/gostl/ds/set"
+       "github.com/stretchr/testify/assert"
+
+       "github.com/apache/incubator-eventmesh/eventmesh-server-go/pkg/util"
+       "github.com/apache/incubator-eventmesh/eventmesh-server-go/plugin"
+       _ 
"github.com/apache/incubator-eventmesh/eventmesh-server-go/plugin/connector/standalone"
+       
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/consts"
+       
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/core/protocol/grpc/emitter"
+       emitermock 
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/core/protocol/grpc/emitter/mocks"
+       
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/proto/pb"
+)
+
+func Test_NewConsumerManager(t *testing.T) {
+       mgr, err := NewConsumerManager()
+       assert.NoError(t, err)
+       assert.NotNil(t, mgr)
+}
+
+func Test_GetConsumer(t *testing.T) {
+       tests := []struct {
+               name   string
+               expect func(t *testing.T, mgr ConsumerManager)
+       }{
+               {
+                       name: "add one consumer",
+                       expect: func(t *testing.T, mgr ConsumerManager) {
+                               mesh, err := mgr.GetConsumer("consumergroup")
+                               assert.NoError(t, err)
+                               assert.NotNil(t, mesh)
+                       },
+               },
+               {
+                       name: "return exist one",
+                       expect: func(t *testing.T, mgr ConsumerManager) {
+                               mesh, err := mgr.GetConsumer("consumergroup")
+                               assert.NoError(t, err)
+                               assert.NotNil(t, mesh)
+
+                               mesh2, err := mgr.GetConsumer("consumergroup")
+                               assert.NoError(t, err)
+                               assert.NotNil(t, mesh2)
+                               assert.Equal(t, mesh, mesh2)
+                       },
+               },
+       }
+
+       plugin.SetActivePlugin(map[string]string{
+               "connector": "standalone",
+       })
+       for _, tc := range tests {
+               t.Run(tc.name, func(t *testing.T) {
+                       mgr, err := NewConsumerManager()
+                       assert.NoError(t, err)
+                       assert.NotNil(t, mgr)
+                       tc.expect(t, mgr)
+               })
+       }
+}
+
+func Test_RegisterClient(t *testing.T) {
+       getClientInConsumer := func(cg string, mgr ConsumerManager) 
(*GroupClient, bool) {
+               cm := mgr.(*consumerManager)
+               v, ok := cm.consumerGroupClients.Load(cg)
+               if !ok {
+                       return nil, false
+               }
+               return v.(*set.Set).First().Value().(*GroupClient), true
+       }
+       tests := []struct {
+               name   string
+               expect func(t *testing.T, mgr ConsumerManager)
+       }{
+               {
+                       name: "register new",
+                       expect: func(t *testing.T, mgr ConsumerManager) {
+                               err := mgr.RegisterClient(&GroupClient{
+                                       ENV:              "env",
+                                       IDC:              "IDC",
+                                       ConsumerGroup:    "ConsumerGroup",
+                                       Topic:            "Topic",
+                                       GRPCType:         consts.WEBHOOK,
+                                       URL:              "http://test.com";,
+                                       SubscriptionMode: 
pb.Subscription_SubscriptionItem_CLUSTERING,
+                                       SYS:              "SYS",
+                                       IP:               "IP",
+                                       PID:              util.PID(),
+                                       Hostname:         "test",
+                                       APIVersion:       "v1",
+                                       LastUPTime:       time.Now(),
+                                       Emiter:           
emitter.NewEventEmitter(nil),
+                               })
+                               assert.NoError(t, err)
+                       },
+               },
+               {
+                       name: "webhook register exist and update time",
+                       expect: func(t *testing.T, mgr ConsumerManager) {
+                               firstTime := time.Now()
+                               oldUrl := "http://old.test.com";
+                               cli := &GroupClient{
+                                       ENV:              "env",
+                                       IDC:              "IDC",
+                                       ConsumerGroup:    "ConsumerGroup",
+                                       Topic:            "Topic",
+                                       GRPCType:         consts.WEBHOOK,
+                                       URL:              oldUrl,
+                                       SubscriptionMode: 
pb.Subscription_SubscriptionItem_CLUSTERING,
+                                       SYS:              "SYS",
+                                       IP:               "IP",
+                                       PID:              util.PID(),
+                                       Hostname:         "test",
+                                       APIVersion:       "v1",
+                                       LastUPTime:       firstTime,
+                                       Emiter:           
emitter.NewEventEmitter(nil),
+                               }
+                               assert.NoError(t, mgr.RegisterClient(cli))
+                               time.Sleep(time.Second)
+                               cli.URL = "http://new.test.com";
+                               cli.LastUPTime = time.Now()
+                               assert.NoError(t, mgr.RegisterClient(cli))
+                               cliInMgr, ok := 
getClientInConsumer(cli.ConsumerGroup, mgr)
+                               assert.True(t, ok)
+                               assert.NotEqual(t, 
cliInMgr.LastUPTime.Second(), firstTime.Second())
+                               assert.Equal(t, cliInMgr.URL, 
"http://new.test.com";)
+                       },
+               },
+               {
+                       name: "stream register exist and update time",
+                       expect: func(t *testing.T, mgr ConsumerManager) {
+                               mockctl := gomock.NewController(t)
+                               oldEmiter := 
emitermock.NewMockEventEmitter(mockctl)
+                               
oldEmiter.EXPECT().SendStreamResp(&pb.RequestHeader{}, 
&grpc.StatusCode{}).Return(nil).Times(1)
+                               firstTime := time.Now()
+                               cli := &GroupClient{
+                                       ENV:              "env",
+                                       IDC:              "IDC",
+                                       ConsumerGroup:    "ConsumerGroup",
+                                       Topic:            "Topic",
+                                       GRPCType:         consts.WEBHOOK,
+                                       SubscriptionMode: 
pb.Subscription_SubscriptionItem_CLUSTERING,
+                                       SYS:              "SYS",
+                                       IP:               "IP",
+                                       PID:              util.PID(),
+                                       Hostname:         "test",
+                                       APIVersion:       "v1",
+                                       LastUPTime:       firstTime,
+                                       Emiter:           oldEmiter,
+                               }
+                               assert.NoError(t, mgr.RegisterClient(cli))
+                               time.Sleep(time.Second)
+                               newEmiter := 
emitermock.NewMockEventEmitter(mockctl)
+                               
newEmiter.EXPECT().SendStreamResp(&pb.RequestHeader{}, 
&grpc.StatusCode{}).Return(fmt.Errorf("error")).Times(1)
+                               cli.Emiter = newEmiter
+                               cli.LastUPTime = time.Now()
+                               assert.NoError(t, mgr.RegisterClient(cli))
+                               cliInMgr, ok := 
getClientInConsumer(cli.ConsumerGroup, mgr)
+                               assert.True(t, ok)
+                               assert.NotEqual(t, 
cliInMgr.LastUPTime.Second(), firstTime.Second())
+                               assert.Nil(t, 
oldEmiter.SendStreamResp(&pb.RequestHeader{}, &grpc.StatusCode{}))
+                               assert.Error(t, 
cliInMgr.Emiter.SendStreamResp(&pb.RequestHeader{}, &grpc.StatusCode{}))
+                       },
+               },
+       }
+       for _, tc := range tests {
+               t.Run(tc.name, func(t *testing.T) {
+                       mgr, err := NewConsumerManager()
+                       assert.NoError(t, err)
+                       assert.NotNil(t, mgr)
+                       tc.expect(t, mgr)
+               })
+       }
+}
diff --git a/eventmesh-server-go/runtime/core/protocol/grpc/consumer_mesh.go 
b/eventmesh-server-go/runtime/core/protocol/grpc/consumer/consumer_mesh.go
similarity index 88%
rename from eventmesh-server-go/runtime/core/protocol/grpc/consumer_mesh.go
rename to 
eventmesh-server-go/runtime/core/protocol/grpc/consumer/consumer_mesh.go
index ca2da90e3..bece674f0 100644
--- a/eventmesh-server-go/runtime/core/protocol/grpc/consumer_mesh.go
+++ b/eventmesh-server-go/runtime/core/protocol/grpc/consumer/consumer_mesh.go
@@ -13,12 +13,15 @@
 // See the License for the specific language governing permissions and
 // limitations under the License.
 
-package grpc
+package consumer
 
 import (
        "sync"
        "time"
 
+       cloudv2 "github.com/cloudevents/sdk-go/v2"
+       "github.com/pkg/errors"
+
        "github.com/apache/incubator-eventmesh/eventmesh-server-go/config"
        "github.com/apache/incubator-eventmesh/eventmesh-server-go/log"
        "github.com/apache/incubator-eventmesh/eventmesh-server-go/pkg/util"
@@ -26,8 +29,6 @@ import (
        
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/consts"
        
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/core/wrapper"
        
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/proto/pb"
-       cloudv2 "github.com/cloudevents/sdk-go/v2"
-       "github.com/pkg/errors"
 )
 
 var (
@@ -36,18 +37,27 @@ var (
        ErrNewProducerConnector = errors.New("create producer connector err")
 )
 
-type EventMeshConsumer struct {
+type EventMeshConsumer interface {
+       Init() error
+       Start() error
+       ServiceState() consts.ServiceState
+       RegisterClient(cli *GroupClient) bool
+       DeRegisterClient(cli *GroupClient) bool
+       Shutdown() error
+}
+
+type eventMeshConsumer struct {
        ConsumerGroup      string
        persistentConsumer *wrapper.Consumer
        broadcastConsumer  *wrapper.Consumer
-       messageHandler     *MessageHandler
-       ServiceState       consts.ServiceState
+       messageHandler     MessageHandler
+       serviceState       consts.ServiceState
        // consumerGroupTopicConfig key is topic
        // value is ConsumerGroupTopicOption
        consumerGroupTopicConfig *sync.Map
 }
 
-func NewEventMeshConsumer(consumerGroup string) (*EventMeshConsumer, error) {
+func NewEventMeshConsumer(consumerGroup string) (EventMeshConsumer, error) {
        pushHandler, err := NewMessageHandler(consumerGroup)
        if err != nil {
                return nil, err
@@ -60,7 +70,7 @@ func NewEventMeshConsumer(consumerGroup string) 
(*EventMeshConsumer, error) {
        if err != nil {
                return nil, ErrNewConsumerConnector
        }
-       return &EventMeshConsumer{
+       return &eventMeshConsumer{
                ConsumerGroup:            consumerGroup,
                messageHandler:           pushHandler,
                persistentConsumer:       cons,
@@ -69,7 +79,11 @@ func NewEventMeshConsumer(consumerGroup string) 
(*EventMeshConsumer, error) {
        }, nil
 }
 
-func (e *EventMeshConsumer) Init() error {
+func (e *eventMeshConsumer) ServiceState() consts.ServiceState {
+       return e.serviceState
+}
+
+func (e *eventMeshConsumer) Init() error {
        // no topics, don't init the consumer
        if e.ConsumerGroupSize() == 0 {
                return nil
@@ -98,13 +112,13 @@ func (e *EventMeshConsumer) Init() error {
        }
        broadcastEventListener := 
e.createEventListener(pb.Subscription_SubscriptionItem_BROADCASTING)
        e.broadcastConsumer.RegisterListener(broadcastEventListener)
-       e.ServiceState = consts.INITED
+       e.serviceState = consts.INITED
 
        log.Infof("init the eventmesh consumer success, group:%v", 
e.ConsumerGroup)
        return nil
 }
 
-func (e *EventMeshConsumer) Start() error {
+func (e *eventMeshConsumer) Start() error {
        // no topics, don't start the consumer
        if e.ConsumerGroupSize() == 0 {
                return nil
@@ -131,13 +145,13 @@ func (e *EventMeshConsumer) Start() error {
                return err
        }
 
-       e.ServiceState = consts.RUNNING
+       e.serviceState = consts.RUNNING
        return nil
 }
 
 // RegisterClient Register client's topic information
 // return true if this EventMeshConsumer required restart because of the topic 
changes
-func (e *EventMeshConsumer) RegisterClient(cli *GroupClient) bool {
+func (e *eventMeshConsumer) RegisterClient(cli *GroupClient) bool {
        var (
                consumerTopicOption ConsumerGroupTopicOption
                restart             = false
@@ -157,7 +171,7 @@ func (e *EventMeshConsumer) RegisterClient(cli 
*GroupClient) bool {
 // DeRegisterClient deregister client's topic information and return true if 
this EventMeshConsumer
 // required restart because of the topic changes
 // return true if the underlining EventMeshConsumer needs to restart later; 
false otherwise
-func (e *EventMeshConsumer) DeRegisterClient(cli *GroupClient) bool {
+func (e *eventMeshConsumer) DeRegisterClient(cli *GroupClient) bool {
        var (
                consumerTopicOption ConsumerGroupTopicOption
        )
@@ -171,7 +185,7 @@ func (e *EventMeshConsumer) DeRegisterClient(cli 
*GroupClient) bool {
        return true
 }
 
-func (e *EventMeshConsumer) Shutdown() error {
+func (e *eventMeshConsumer) Shutdown() error {
        if err := e.persistentConsumer.Shutdown(); err != nil {
                return err
        }
@@ -181,7 +195,7 @@ func (e *EventMeshConsumer) Shutdown() error {
        return nil
 }
 
-func (e *EventMeshConsumer) ConsumerGroupSize() int {
+func (e *eventMeshConsumer) ConsumerGroupSize() int {
        count := 0
        e.consumerGroupTopicConfig.Range(func(key, value any) bool {
                count++
@@ -190,7 +204,7 @@ func (e *EventMeshConsumer) ConsumerGroupSize() int {
        return count
 }
 
-func (e *EventMeshConsumer) createEventListener(mode 
pb.Subscription_SubscriptionItem_SubscriptionMode) *connector.EventListener {
+func (e *eventMeshConsumer) createEventListener(mode 
pb.Subscription_SubscriptionItem_SubscriptionMode) *connector.EventListener {
        return &connector.EventListener{
                Consume: func(event *cloudv2.Event, commitFunc 
connector.CommitFunc) error {
                        var commitAction connector.EventMeshAction
diff --git 
a/eventmesh-server-go/runtime/core/protocol/grpc/consumer/consumer_processor.go 
b/eventmesh-server-go/runtime/core/protocol/grpc/consumer/consumer_processor.go
new file mode 100644
index 000000000..53c0b2616
--- /dev/null
+++ 
b/eventmesh-server-go/runtime/core/protocol/grpc/consumer/consumer_processor.go
@@ -0,0 +1,321 @@
+// 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 consumer
+
+import (
+       "context"
+       "fmt"
+       "github.com/apache/incubator-eventmesh/eventmesh-server-go/log"
+       
"github.com/apache/incubator-eventmesh/eventmesh-server-go/pkg/common/protocol/grpc"
+       "github.com/apache/incubator-eventmesh/eventmesh-server-go/plugin"
+       
"github.com/apache/incubator-eventmesh/eventmesh-server-go/plugin/connector"
+       
"github.com/apache/incubator-eventmesh/eventmesh-server-go/plugin/protocol"
+       
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/consts"
+       
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/core/protocol/grpc/emitter"
+       
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/core/protocol/grpc/producer"
+       
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/core/protocol/grpc/validator"
+       
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/proto/pb"
+       "time"
+)
+
+var (
+       ErrProtocolPluginNotFound = fmt.Errorf("protocol plugin not found")
+)
+
+type Processor interface {
+       Subscribe(consumerMgr ConsumerManager, msg *pb.Subscription) 
(*pb.Response, error)
+       UnSubscribe(consumerMgr ConsumerManager, msg *pb.Subscription) 
(*pb.Response, error)
+       SubscribeStream(consumerMgr ConsumerManager, emiter 
emitter.EventEmitter, msg *pb.Subscription) error
+       Heartbeat(consumerMgr ConsumerManager, msg *pb.Heartbeat) 
(*pb.Response, error)
+       ReplyMessage(ctx context.Context, producerMgr producer.ProducerManager, 
emiter emitter.EventEmitter, msg *pb.SimpleMessage) error
+}
+
+type processor struct {
+}
+
+func NewProcessor() Processor {
+       return &processor{}
+}
+
+func (p *processor) Subscribe(consumerMgr ConsumerManager, msg 
*pb.Subscription) (*pb.Response, error) {
+       hdr := msg.Header
+       if err := validator.ValidateHeader(hdr); err != nil {
+               log.Warnf("invalid header:%v", err)
+               return buildPBResponse(grpc.EVENTMESH_PROTOCOL_HEADER_ERR), err
+       }
+       if err := validator.ValidateSubscription(consts.WEBHOOK, msg); err != 
nil {
+               log.Warnf("invalid body:%v", err)
+               return buildPBResponse(grpc.EVENTMESH_PROTOCOL_BODY_ERR), err
+       }
+       consumerGroup := msg.ConsumerGroup
+       url := msg.Url
+       items := msg.SubscriptionItems
+       var newClients []*GroupClient
+       for _, item := range items {
+               newClients = append(newClients, &GroupClient{
+                       ENV:              hdr.Env,
+                       IDC:              hdr.Idc,
+                       SYS:              hdr.Sys,
+                       IP:               hdr.Ip,
+                       PID:              hdr.Pid,
+                       ConsumerGroup:    consumerGroup,
+                       Topic:            item.Topic,
+                       SubscriptionMode: item.Mode,
+                       GRPCType:         consts.WEBHOOK,
+                       URL:              url,
+                       LastUPTime:       time.Now(),
+               })
+       }
+       for _, cli := range newClients {
+               if err := consumerMgr.RegisterClient(cli); err != nil {
+                       return 
buildPBResponse(grpc.EVENTMESH_Subscribe_Register_ERR), err
+               }
+       }
+       meshConsumer, err := consumerMgr.GetConsumer(consumerGroup)
+       if err != nil {
+               return buildPBResponse(grpc.EVENTMESH_Consumer_NotFound_ERR), 
err
+       }
+       requireRestart := false
+       for _, cli := range newClients {
+               if meshConsumer.RegisterClient(cli) {
+                       requireRestart = true
+               }
+       }
+       if requireRestart {
+               log.Infof("ConsumerGroup %v topic info changed, restart 
EventMesh Consumer", consumerGroup)
+               if err := consumerMgr.RestartConsumer(consumerGroup); err != 
nil {
+                       return 
buildPBResponse(grpc.EVENTMESH_Consumer_NotFound_ERR), err
+               }
+       } else {
+               log.Warnf("EventMesh consumer [%v] didn't restart.", 
consumerGroup)
+       }
+       return buildPBResponse(grpc.SUCCESS), nil
+}
+
+func (p *processor) UnSubscribe(consumerMgr ConsumerManager, msg 
*pb.Subscription) (*pb.Response, error) {
+       hdr := msg.Header
+       if err := validator.ValidateHeader(hdr); err != nil {
+               log.Warnf("invalid header:%v", err)
+               return buildPBResponse(grpc.EVENTMESH_PROTOCOL_HEADER_ERR), err
+       }
+       if err := validator.ValidateSubscription(consts.WEBHOOK, msg); err != 
nil {
+               log.Warnf("invalid body:%v", err)
+               return buildPBResponse(grpc.EVENTMESH_PROTOCOL_BODY_ERR), err
+       }
+       consumerGroup := msg.ConsumerGroup
+       url := msg.Url
+       items := msg.SubscriptionItems
+       var removeClients []*GroupClient
+       for _, item := range items {
+               removeClients = append(removeClients, &GroupClient{
+                       ENV:              hdr.Env,
+                       IDC:              hdr.Idc,
+                       SYS:              hdr.Sys,
+                       IP:               hdr.Ip,
+                       PID:              hdr.Pid,
+                       ConsumerGroup:    consumerGroup,
+                       Topic:            item.Topic,
+                       SubscriptionMode: item.Mode,
+                       GRPCType:         consts.WEBHOOK,
+                       URL:              url,
+                       LastUPTime:       time.Now(),
+               })
+       }
+       for _, cli := range removeClients {
+               if err := consumerMgr.DeRegisterClient(cli); err != nil {
+                       return 
buildPBResponse(grpc.EVENTMESH_Subscribe_Register_ERR), err
+               }
+       }
+       meshConsumer, err := consumerMgr.GetConsumer(consumerGroup)
+       if err != nil {
+               return buildPBResponse(grpc.EVENTMESH_Consumer_NotFound_ERR), 
err
+       }
+       requireRestart := false
+       for _, cli := range removeClients {
+               if meshConsumer.DeRegisterClient(cli) {
+                       requireRestart = true
+               }
+       }
+       if requireRestart {
+               log.Infof("ConsumerGroup %v topic info changed, restart 
EventMesh Consumer", consumerGroup)
+               if err := consumerMgr.RestartConsumer(consumerGroup); err != 
nil {
+                       return 
buildPBResponse(grpc.EVENTMESH_Consumer_NotFound_ERR), err
+               }
+       } else {
+               log.Warnf("EventMesh consumer [%v] didn't restart.", 
consumerGroup)
+       }
+       return buildPBResponse(grpc.SUCCESS), nil
+}
+
+func (p *processor) SubscribeStream(consumerMgr ConsumerManager, emiter 
emitter.EventEmitter, msg *pb.Subscription) error {
+       hdr := msg.Header
+       if err := validator.ValidateHeader(hdr); err != nil {
+               log.Warnf("invalid header:%v", err)
+               emiter.SendStreamResp(hdr, grpc.EVENTMESH_PROTOCOL_HEADER_ERR)
+               return err
+       }
+       if err := validator.ValidateSubscription(consts.STREAM, msg); err != 
nil {
+               log.Warnf("invalid body:%v", err)
+               emiter.SendStreamResp(hdr, grpc.EVENTMESH_PROTOCOL_BODY_ERR)
+               return err
+       }
+       consumerGroup := msg.ConsumerGroup
+       var clients []*GroupClient
+       for _, item := range msg.SubscriptionItems {
+               clients = append(clients, &GroupClient{
+                       ENV:              hdr.Env,
+                       IDC:              hdr.Idc,
+                       SYS:              hdr.Sys,
+                       IP:               hdr.Ip,
+                       PID:              hdr.Pid,
+                       ConsumerGroup:    consumerGroup,
+                       Topic:            item.Topic,
+                       SubscriptionMode: item.Mode,
+                       GRPCType:         consts.STREAM,
+                       LastUPTime:       time.Now(),
+                       Emiter:           emiter,
+               })
+       }
+       for _, cli := range clients {
+               if err := consumerMgr.RegisterClient(cli); err != nil {
+                       return err
+               }
+       }
+       meshConsumer, err := consumerMgr.GetConsumer(consumerGroup)
+       if err != nil {
+               return err
+       }
+       requireRestart := false
+       for _, cli := range clients {
+               if meshConsumer.RegisterClient(cli) {
+                       requireRestart = true
+               }
+       }
+       if requireRestart {
+               log.Infof("ConsumerGroup %v topic info changed, restart 
EventMesh Consumer", consumerGroup)
+               return consumerMgr.RestartConsumer(consumerGroup)
+       } else {
+               log.Warnf("EventMesh consumer [%v] didn't restart.", 
consumerGroup)
+       }
+
+       return nil
+}
+
+func (p *processor) Heartbeat(consumerMgr ConsumerManager, msg *pb.Heartbeat) 
(*pb.Response, error) {
+       hdr := msg.Header
+       if err := validator.ValidateHeader(hdr); err != nil {
+               log.Warnf("invalid header:%v", err)
+               return buildPBResponse(grpc.EVENTMESH_PROTOCOL_HEADER_ERR), err
+       }
+       if err := validator.ValidateHeartBeat(msg); err != nil {
+               log.Warnf("invalid body:%v", err)
+               return buildPBResponse(grpc.EVENTMESH_PROTOCOL_BODY_ERR), err
+       }
+       if msg.ClientType != pb.Heartbeat_SUB {
+               log.Warnf("client type err, not sub")
+               return buildPBResponse(grpc.EVENTMESH_Heartbeat_Protocol_ERR), 
fmt.Errorf("protocol not sub")
+       }
+       consumerGroup := msg.ConsumerGroup
+       for _, item := range msg.HeartbeatItems {
+               cli := &GroupClient{
+                       ENV:           hdr.Env,
+                       IDC:           hdr.Idc,
+                       SYS:           hdr.Sys,
+                       IP:            hdr.Ip,
+                       PID:           hdr.Pid,
+                       ConsumerGroup: consumerGroup,
+                       Topic:         item.Topic,
+                       LastUPTime:    time.Now(),
+               }
+               consumerMgr.UpdateClientTime(cli)
+       }
+       return buildPBResponse(grpc.SUCCESS), nil
+}
+
+func (p *processor) ReplyMessage(ctx context.Context, producerMgr 
producer.ProducerManager, emiter emitter.EventEmitter, msg *pb.SimpleMessage) 
error {
+       hdr := msg.Header
+       if err := validator.ValidateHeader(hdr); err != nil {
+               log.Warnf("invalid header:%v", err)
+               emiter.SendStreamResp(hdr, grpc.EVENTMESH_PROTOCOL_HEADER_ERR)
+               return err
+       }
+       if err := validator.ValidateMessage(msg); err != nil {
+               log.Warnf("invalid body:%v", err)
+               emiter.SendStreamResp(hdr, grpc.EVENTMESH_PROTOCOL_BODY_ERR)
+               return err
+       }
+       seqNum := msg.SeqNum
+       uniqID := msg.UniqueId
+       producerGroup := msg.ProducerGroup
+       mqCluster := 
defaultIfEmpty(msg.Properties[consts.PROPERTY_MESSAGE_CLUSTER], 
"defaultCluster")
+       replyTopic := mqCluster + "_" + consts.RR_REPLY_TOPIC
+       msg.Topic = replyTopic
+       protocolType := hdr.ProtocolType
+       adp := plugin.Get(plugin.Protocol, protocolType).(protocol.Adapter)
+       if adp == nil {
+               log.Warnf("protocol plugin not found:%v", protocolType)
+               emiter.SendStreamResp(hdr, grpc.EVENTMESH_Plugin_NotFound_ERR)
+               return ErrProtocolPluginNotFound
+       }
+       cevt, err := adp.ToCloudEvent(&grpc.SimpleMessageWrapper{SimpleMessage: 
msg})
+       if err != nil {
+               log.Warnf("transfer to cloud event msg err:%v", err)
+               emiter.SendStreamResp(hdr, grpc.EVENTMESH_Transfer_Protocol_ERR)
+               return err
+       }
+       emProducer, err := producerMgr.GetProducer(producerGroup)
+       if err != nil {
+               log.Warnf("no eventmesh producer found, err:%v, group:%v", err, 
producerGroup)
+               emiter.SendStreamResp(hdr, 
grpc.EVENTMESH_Producer_Group_NotFound_ERR)
+               return err
+       }
+       start := time.Now()
+       return emProducer.Reply(
+               producer.SendMessageContext{
+                       Ctx:         ctx,
+                       Event:       cevt,
+                       BizSeqNO:    seqNum,
+                       ProducerAPI: emProducer,
+                       CreateTime:  time.Now(),
+               },
+               &connector.SendCallback{
+                       OnSuccess: func(result *connector.SendResult) {
+                               
log.Infof("message|mq2eventmesh|REPLY|ReplyToServer|send2MQCost=%vms|topic=%v|bizSeqNo=%v|uniqueId=%v",
+                                       time.Now().Sub(start).Milliseconds(), 
replyTopic, seqNum, uniqID)
+                       },
+                       OnError: func(result *connector.ErrorResult) {
+                               emiter.SendStreamResp(hdr, 
grpc.EVENTMESH_REPLY_MSG_ERR)
+                               
log.Errorf("message|mq2eventmesh|REPLY|ReplyToServer|send2MQCost=%vms|topic=%v|bizSeqNo=%v|uniqueId=%v",
+                                       time.Now().Sub(start).Milliseconds(), 
replyTopic, seqNum, uniqID, result.Err)
+                       },
+               },
+       )
+}
+
+func defaultIfEmpty(in interface{}, def string) string {
+       if in == nil {
+               return def
+       }
+       return in.(string)
+}
+
+func buildPBResponse(code *grpc.StatusCode) *pb.Response {
+       return &pb.Response{
+               RespCode: code.RetCode,
+               RespMsg:  code.ErrMsg,
+               RespTime: fmt.Sprintf("%v", time.Now().UnixMilli()),
+       }
+}
diff --git a/eventmesh-server-go/runtime/core/protocol/grpc/consumer_service.go 
b/eventmesh-server-go/runtime/core/protocol/grpc/consumer/consumer_service.go
similarity index 79%
rename from eventmesh-server-go/runtime/core/protocol/grpc/consumer_service.go
rename to 
eventmesh-server-go/runtime/core/protocol/grpc/consumer/consumer_service.go
index 2af3f6613..7199f160e 100644
--- a/eventmesh-server-go/runtime/core/protocol/grpc/consumer_service.go
+++ 
b/eventmesh-server-go/runtime/core/protocol/grpc/consumer/consumer_service.go
@@ -13,28 +13,34 @@
 // See the License for the specific language governing permissions and
 // limitations under the License.
 
-package grpc
+package consumer
 
 import (
        "context"
        "github.com/apache/incubator-eventmesh/eventmesh-server-go/config"
        "github.com/apache/incubator-eventmesh/eventmesh-server-go/log"
+       
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/core/protocol/grpc/emitter"
+       
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/core/protocol/grpc/producer"
        
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/proto/pb"
        "github.com/panjf2000/ants/v2"
        "io"
+       "time"
 )
 
+var defaultAsyncTimeout = time.Second * 5
+
 // ConsumerService grpc service
 type ConsumerService struct {
        pb.UnimplementedConsumerServiceServer
-       gctx          *GRPCContext
-       subscribePool *ants.Pool
-       replyPool     *ants.Pool
-       msgToClient   chan *pb.SimpleMessage
-       subFromClient chan *pb.Subscription
+       consumerManager ConsumerManager
+       producerManager producer.ProducerManager
+       subscribePool   *ants.Pool
+       replyPool       *ants.Pool
+       msgToClient     chan *pb.SimpleMessage
+       subFromClient   chan *pb.Subscription
 }
 
-func NewConsumerServiceServer(gctx *GRPCContext) (*ConsumerService, error) {
+func NewConsumerServiceServer(consumerManager ConsumerManager, producerManager 
producer.ProducerManager) (*ConsumerService, error) {
        ss := config.GlobalConfig().Server.GRPCOption.SubscribePoolSize
        subPool, err := ants.NewPool(ss)
        if err != nil {
@@ -46,9 +52,10 @@ func NewConsumerServiceServer(gctx *GRPCContext) 
(*ConsumerService, error) {
                return nil, err
        }
        return &ConsumerService{
-               gctx:          gctx,
-               subscribePool: subPool,
-               replyPool:     replyPool,
+               consumerManager: consumerManager,
+               producerManager: producerManager,
+               subscribePool:   subPool,
+               replyPool:       replyPool,
        }, nil
 }
 
@@ -62,7 +69,7 @@ func (c *ConsumerService) Subscribe(ctx context.Context, sub 
*pb.Subscription) (
                err     error
        )
        c.subscribePool.Submit(func() {
-               resp, err = ProcessSubscribe(c.gctx, sub)
+               resp, err = NewProcessor().Subscribe(c.consumerManager, sub)
                errChan <- err
        })
        select {
@@ -82,7 +89,6 @@ func (c *ConsumerService) Subscribe(ctx context.Context, sub 
*pb.Subscription) (
 // for Recv() goroutine if got err==io.EOF as the client close the 
stream(即客户端关闭stream)
 // for Send() refers to https://github.com/grpc/grpc-go/issues/444
 func (c *ConsumerService) SubscribeStream(stream 
pb.ConsumerService_SubscribeStreamServer) error {
-       //go func() {
        for {
                req, err := stream.Recv()
                if err == io.EOF {
@@ -101,8 +107,6 @@ func (c *ConsumerService) SubscribeStream(stream 
pb.ConsumerService_SubscribeStr
                        c.handleSubscriptionStream(req, stream)
                }
        }
-       //}()
-
        return nil
 }
 
@@ -116,7 +120,7 @@ func (c *ConsumerService) Unsubscribe(ctx context.Context, 
sub *pb.Subscription)
                err     error
        )
        c.subscribePool.Submit(func() {
-               resp, err = ProcessUnSubscribe(c.gctx, sub)
+               resp, err = NewProcessor().UnSubscribe(c.consumerManager, sub)
                errChan <- err
        })
        select {
@@ -133,17 +137,17 @@ func (c *ConsumerService) Unsubscribe(ctx 
context.Context, sub *pb.Subscription)
 
 func (c *ConsumerService) handleSubscriptionStream(sub *pb.Subscription, 
stream pb.ConsumerService_SubscribeStreamServer) error {
        c.subscribePool.Submit(func() {
-               emiter := &EventEmitter{emitter: stream}
-               ProcessSubscribeStream(context.TODO(), c.gctx, emiter, sub)
+               emiter := emitter.NewEventEmitter(stream)
+               NewProcessor().SubscribeStream(c.consumerManager, emiter, sub)
        })
        return nil
 }
 
 func (c *ConsumerService) handleSubscribeReply(sub *pb.Subscription, stream 
pb.ConsumerService_SubscribeStreamServer) error {
        c.replyPool.Submit(func() {
-               emiter := &EventEmitter{emitter: stream}
+               emiter := emitter.NewEventEmitter(stream)
                reply := sub.Reply
-               ProcessReplyMessage(context.TODO(), c.gctx, emiter, 
&pb.SimpleMessage{
+               NewProcessor().ReplyMessage(context.TODO(), c.producerManager, 
emiter, &pb.SimpleMessage{
                        Header:        sub.Header,
                        ProducerGroup: reply.ProducerGroup,
                        Content:       reply.Content,
diff --git 
a/eventmesh-server-go/runtime/core/protocol/grpc/consumer/consumer_service_test.go
 
b/eventmesh-server-go/runtime/core/protocol/grpc/consumer/consumer_service_test.go
new file mode 100644
index 000000000..179cd8416
--- /dev/null
+++ 
b/eventmesh-server-go/runtime/core/protocol/grpc/consumer/consumer_service_test.go
@@ -0,0 +1,123 @@
+// 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 consumer
+
+import (
+       "context"
+       "github.com/apache/incubator-eventmesh/eventmesh-server-go/pkg/util"
+       
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/core/protocol/grpc/mocks"
+       
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/proto/pb"
+       "github.com/golang/mock/gomock"
+       "testing"
+)
+
+func Test_Subscribe(t *testing.T) {
+       mockctl := gomock.NewController(t)
+       mockConsumer := mocks.NewMockConsumerServiceServer(mockctl)
+       mockConsumer.EXPECT().Subscribe(context.TODO(), &pb.Subscription{
+               Header: &pb.RequestHeader{
+                       Env:             "grpc-env",
+                       Region:          "sh",
+                       Idc:             "idc-sh",
+                       Ip:              util.GetIP(),
+                       Pid:             util.PID(),
+                       Sys:             "grpc-sys",
+                       Username:        "grpc-username",
+                       Password:        "grpc-passwd",
+                       Language:        "Go",
+                       ProtocolType:    "cloudevents",
+                       ProtocolVersion: "1.0",
+                       ProtocolDesc:    "cloudevents",
+               },
+               ConsumerGroup: "grpc-stream-consumergroup",
+               SubscriptionItems: []*pb.Subscription_SubscriptionItem{
+                       {
+                               Topic: "test_topic",
+                               Mode:  
pb.Subscription_SubscriptionItem_CLUSTERING,
+                               Type:  pb.Subscription_SubscriptionItem_SYNC,
+                       },
+               },
+               Url: "http://127.0.0.1:18080/onmessage";,
+       }).Return()
+}
+
+func Test_unsubscribe(t *testing.T) {
+       //cli := grpc.newTestClient(t)
+       //assert.NotNil(t, cli)
+       //resp, err := cli.consumerClient.Unsubscribe(context.TODO(), 
&pb.Subscription{
+       //      Header: &pb.RequestHeader{
+       //              Env:             "grpc-env",
+       //              Region:          "sh",
+       //              Idc:             "idc-sh",
+       //              Ip:              util.GetIP(),
+       //              Pid:             util.PID(),
+       //              Sys:             "grpc-sys",
+       //              Username:        "grpc-username",
+       //              Password:        "grpc-passwd",
+       //              Language:        "Go",
+       //              ProtocolType:    "cloudevents",
+       //              ProtocolVersion: "1.0",
+       //              ProtocolDesc:    "cloudevents",
+       //      },
+       //      ConsumerGroup: "grpc-stream-consumergroup",
+       //      SubscriptionItems: []*pb.Subscription_SubscriptionItem{
+       //              {
+       //                      Topic: grpc._testWebhookTopic,
+       //                      Mode:  
pb.Subscription_SubscriptionItem_CLUSTERING,
+       //                      Type:  pb.Subscription_SubscriptionItem_SYNC,
+       //              },
+       //      },
+       //      Url: "http://127.0.0.1:18080/onmessage";,
+       //})
+       //assert.NoError(t, err)
+       //assert.NotNil(t, cli)
+       //t.Log(resp.String())
+       //assert.Equal(t, resp.RespCode, 0)
+}
+
+func Test_subscribeStream(t *testing.T) {
+       //cli := grpc.newTestClient(t)
+       //assert.NotNil(t, cli)
+       //
+       //stream, err := cli.consumerClient.SubscribeStream(context.TODO())
+       //assert.NoError(t, err)
+       //err = stream.Send(&pb.Subscription{
+       //      Header: &pb.RequestHeader{
+       //              Env:             "grpc-env",
+       //              Region:          "sh",
+       //              Idc:             "idc-sh",
+       //              Ip:              util.GetIP(),
+       //              Pid:             util.PID(),
+       //              Sys:             "grpc-sys",
+       //              Username:        "grpc-username",
+       //              Password:        "grpc-passwd",
+       //              Language:        "Go",
+       //              ProtocolType:    "cloudevents",
+       //              ProtocolVersion: "1.0",
+       //              ProtocolDesc:    "cloudevents",
+       //      },
+       //      ConsumerGroup: "grpc-stream-consumergroup",
+       //      SubscriptionItems: []*pb.Subscription_SubscriptionItem{
+       //              {
+       //                      Topic: grpc._testStreamTopic,
+       //                      Mode:  
pb.Subscription_SubscriptionItem_CLUSTERING,
+       //                      Type:  pb.Subscription_SubscriptionItem_ASYNC,
+       //              },
+       //      },
+       //})
+       //assert.NoError(t, err)
+       //time.Sleep(time.Hour)
+}
diff --git a/eventmesh-server-go/runtime/core/protocol/grpc/emitter.go 
b/eventmesh-server-go/runtime/core/protocol/grpc/consumer/message_context.go
similarity index 69%
copy from eventmesh-server-go/runtime/core/protocol/grpc/emitter.go
copy to 
eventmesh-server-go/runtime/core/protocol/grpc/consumer/message_context.go
index 7ccc3a065..8e13a0361 100644
--- a/eventmesh-server-go/runtime/core/protocol/grpc/emitter.go
+++ b/eventmesh-server-go/runtime/core/protocol/grpc/consumer/message_context.go
@@ -13,20 +13,19 @@
 // See the License for the specific language governing permissions and
 // limitations under the License.
 
-package grpc
+package consumer
 
 import (
-       
"github.com/apache/incubator-eventmesh/eventmesh-server-go/pkg/common/protocol/grpc"
+       
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/consts"
        
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/proto/pb"
+       cloudv2 "github.com/cloudevents/sdk-go/v2"
 )
 
-type EventEmitter struct {
-       emitter pb.ConsumerService_SubscribeStreamServer
-}
-
-func (e *EventEmitter) sendStreamResp(hdr *pb.RequestHeader, code 
*grpc.StatusCode) error {
-       return e.emitter.Send(&pb.SimpleMessage{
-               Header:  hdr,
-               Content: code.ToJSONString(),
-       })
+type MessageContext struct {
+       MsgRandomNo      string
+       SubscriptionMode pb.Subscription_SubscriptionItem_SubscriptionMode
+       GrpcType         consts.GRPCType
+       ConsumerGroup    string
+       Event            *cloudv2.Event
+       TopicConfig      ConsumerGroupTopicOption
 }
diff --git a/eventmesh-server-go/runtime/core/protocol/grpc/message_handler.go 
b/eventmesh-server-go/runtime/core/protocol/grpc/consumer/message_handler.go
similarity index 82%
rename from eventmesh-server-go/runtime/core/protocol/grpc/message_handler.go
rename to 
eventmesh-server-go/runtime/core/protocol/grpc/consumer/message_handler.go
index dc24139fd..989f65d0b 100644
--- a/eventmesh-server-go/runtime/core/protocol/grpc/message_handler.go
+++ b/eventmesh-server-go/runtime/core/protocol/grpc/consumer/message_handler.go
@@ -13,10 +13,11 @@
 // See the License for the specific language governing permissions and
 // limitations under the License.
 
-package grpc
+package consumer
 
 import (
        "github.com/apache/incubator-eventmesh/eventmesh-server-go/log"
+       
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/consts"
        "github.com/pkg/errors"
        "sync"
        "time"
@@ -28,18 +29,22 @@ var (
        ErrRequestReachMaxThreshold = errors.New("request reach the max 
threshold")
 )
 
-type MessageHandler struct {
+type MessageHandler interface {
+       Handler(mctx *MessageContext) error
+}
+
+type messageHandler struct {
        //pool *ants.Pool
        // waitingRequests waiting to request
        // key to consumerGroup value to []*PushRequest
        waitingRequests *sync.Map
 }
 
-func NewMessageHandler(consumerGroup string) (*MessageHandler, error) {
+func NewMessageHandler(consumerGroup string) (MessageHandler, error) {
        wr := new(sync.Map)
        // TODO need goroutine safe in []*Request{}
        wr.Store(consumerGroup, []*Request{})
-       hdl := &MessageHandler{
+       hdl := &messageHandler{
                //pool:            p,
                waitingRequests: wr,
        }
@@ -47,13 +52,13 @@ func NewMessageHandler(consumerGroup string) 
(*MessageHandler, error) {
        return hdl, nil
 }
 
-func (m *MessageHandler) checkTimeout() {
+func (m *messageHandler) checkTimeout() {
        tk := time.NewTicker(time.Second)
        for range tk.C {
                m.waitingRequests.Range(func(key, value interface{}) bool {
                        reqs := value.([]*Request)
                        for _, req := range reqs {
-                               if req.timeout() {
+                               if req.Timeout() {
 
                                }
                        }
@@ -62,7 +67,7 @@ func (m *MessageHandler) checkTimeout() {
        }
 }
 
-func (m *MessageHandler) Handler(mctx *MessageContext) error {
+func (m *messageHandler) Handler(mctx *MessageContext) error {
        if m.Size() > ConsumerGroupWaitingRequestThreshold {
                log.Warnf("too many request, reject and send back to MQ, 
group:%v, threshold:%v",
                        mctx.ConsumerGroup, 
ConsumerGroupWaitingRequestThreshold)
@@ -72,7 +77,7 @@ func (m *MessageHandler) Handler(mctx *MessageContext) error {
        var (
                try func() error
        )
-       if mctx.GrpcType == WEBHOOK {
+       if mctx.GrpcType == consts.WEBHOOK {
                req, err := NewWebhookRequest(mctx)
                if err != nil {
                        return err
@@ -93,7 +98,7 @@ func (m *MessageHandler) Handler(mctx *MessageContext) error {
        return nil
 }
 
-func (m *MessageHandler) Size() int {
+func (m *messageHandler) Size() int {
        count := 0
        m.waitingRequests.Range(func(key, value any) bool {
                count++
diff --git a/eventmesh-server-go/runtime/core/protocol/grpc/request.go 
b/eventmesh-server-go/runtime/core/protocol/grpc/consumer/message_request.go
similarity index 95%
rename from eventmesh-server-go/runtime/core/protocol/grpc/request.go
rename to 
eventmesh-server-go/runtime/core/protocol/grpc/consumer/message_request.go
index 52cdde3c7..e32d4dbf6 100644
--- a/eventmesh-server-go/runtime/core/protocol/grpc/request.go
+++ b/eventmesh-server-go/runtime/core/protocol/grpc/consumer/message_request.go
@@ -13,10 +13,23 @@
 // See the License for the specific language governing permissions and
 // limitations under the License.
 
-package grpc
+package consumer
 
 import (
        "fmt"
+       "io"
+       "math/rand"
+       "net/http"
+       "net/url"
+       "sync"
+       "time"
+
+       cloudv2 "github.com/cloudevents/sdk-go/v2"
+       jsoniter "github.com/json-iterator/go"
+       "github.com/liyue201/gostl/ds/set"
+       "github.com/pkg/errors"
+       "go.uber.org/atomic"
+
        "github.com/apache/incubator-eventmesh/eventmesh-server-go/config"
        "github.com/apache/incubator-eventmesh/eventmesh-server-go/log"
        
"github.com/apache/incubator-eventmesh/eventmesh-server-go/pkg/common/protocol/grpc"
@@ -24,24 +37,20 @@ import (
        "github.com/apache/incubator-eventmesh/eventmesh-server-go/plugin"
        
"github.com/apache/incubator-eventmesh/eventmesh-server-go/plugin/protocol"
        
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/consts"
+       
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/core/protocol/grpc/retry"
        
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/proto/pb"
-       cloudv2 "github.com/cloudevents/sdk-go/v2"
-       jsoniter "github.com/json-iterator/go"
-       "github.com/liyue201/gostl/ds/set"
-       "github.com/pkg/errors"
-       "go.uber.org/atomic"
-       "io/ioutil"
-       "math/rand"
-       "net/http"
-       "net/url"
-       "sync"
-       "time"
 )
 
 var (
        ErrNoProtocolFound = errors.New("no protocol type found in event 
message")
 
        defaultWebhookTimeout = time.Second * 5
+
+       jsonPool = sync.Pool{New: func() interface{} {
+               return jsoniter.Config{
+                       EscapeHTML: true,
+               }.Froze()
+       }}
 )
 
 type Response struct {
@@ -50,7 +59,7 @@ type Response struct {
 }
 
 type Request struct {
-       *Context
+       *retry.Retry
 
        MessageContext *MessageContext
        CreateTime     time.Time
@@ -86,7 +95,7 @@ func eventToSimpleMessage(ev *cloudv2.Event) 
(*pb.SimpleMessage, error) {
        return msg.(*grpc.SimpleMessageWrapper).SimpleMessage, nil
 }
 
-func (r *Request) timeout() bool {
+func (r *Request) Timeout() bool {
        return true
 }
 
@@ -182,7 +191,7 @@ func NewWebhookRequest(mctx *MessageContext) 
(*WebhookRequest, error) {
                                        resp.StatusCode, 
hr.SimpleMessage.Topic, hr.SimpleMessage.SeqNum, hr.SimpleMessage.UniqueId)
                                continue
                        }
-                       buf, err := ioutil.ReadAll(resp.Body)
+                       buf, err := io.ReadAll(resp.Body)
                        if err != nil {
                                log.Warnf("err:%v in read response 
url:%v|topic={}|bizSeqNo={}|uniqueId={}",
                                        err, hr.SimpleMessage.Topic, 
hr.SimpleMessage.SeqNum, hr.SimpleMessage.UniqueId)
diff --git 
a/eventmesh-server-go/runtime/core/protocol/grpc/consumer/mocks/consumer_group_option.go
 
b/eventmesh-server-go/runtime/core/protocol/grpc/consumer/mocks/consumer_group_option.go
new file mode 100644
index 000000000..50e85a135
--- /dev/null
+++ 
b/eventmesh-server-go/runtime/core/protocol/grpc/consumer/mocks/consumer_group_option.go
@@ -0,0 +1,193 @@
+// Code generated by MockGen. DO NOT EDIT.
+// Source: 
github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/core/protocol/grpc/consumer
 (interfaces: ConsumerGroupTopicOption)
+
+// Package mocks is a generated GoMock package.
+package mocks
+
+import (
+       reflect "reflect"
+       sync "sync"
+
+       consts 
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/consts"
+       consumer 
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/core/protocol/grpc/consumer"
+       pb 
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/proto/pb"
+       gomock "github.com/golang/mock/gomock"
+       set "github.com/liyue201/gostl/ds/set"
+)
+
+// MockConsumerGroupTopicOption is a mock of ConsumerGroupTopicOption 
interface.
+type MockConsumerGroupTopicOption struct {
+       ctrl     *gomock.Controller
+       recorder *MockConsumerGroupTopicOptionMockRecorder
+}
+
+// MockConsumerGroupTopicOptionMockRecorder is the mock recorder for 
MockConsumerGroupTopicOption.
+type MockConsumerGroupTopicOptionMockRecorder struct {
+       mock *MockConsumerGroupTopicOption
+}
+
+// NewMockConsumerGroupTopicOption creates a new mock instance.
+func NewMockConsumerGroupTopicOption(ctrl *gomock.Controller) 
*MockConsumerGroupTopicOption {
+       mock := &MockConsumerGroupTopicOption{ctrl: ctrl}
+       mock.recorder = &MockConsumerGroupTopicOptionMockRecorder{mock}
+       return mock
+}
+
+// EXPECT returns an object that allows the caller to indicate expected use.
+func (m *MockConsumerGroupTopicOption) EXPECT() 
*MockConsumerGroupTopicOptionMockRecorder {
+       return m.recorder
+}
+
+// AllEmiters mocks base method.
+func (m *MockConsumerGroupTopicOption) AllEmiters() *set.Set {
+       m.ctrl.T.Helper()
+       ret := m.ctrl.Call(m, "AllEmiters")
+       ret0, _ := ret[0].(*set.Set)
+       return ret0
+}
+
+// AllEmiters indicates an expected call of AllEmiters.
+func (mr *MockConsumerGroupTopicOptionMockRecorder) AllEmiters() *gomock.Call {
+       mr.mock.ctrl.T.Helper()
+       return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "AllEmiters", 
reflect.TypeOf((*MockConsumerGroupTopicOption)(nil).AllEmiters))
+}
+
+// AllURLs mocks base method.
+func (m *MockConsumerGroupTopicOption) AllURLs() *set.Set {
+       m.ctrl.T.Helper()
+       ret := m.ctrl.Call(m, "AllURLs")
+       ret0, _ := ret[0].(*set.Set)
+       return ret0
+}
+
+// AllURLs indicates an expected call of AllURLs.
+func (mr *MockConsumerGroupTopicOptionMockRecorder) AllURLs() *gomock.Call {
+       mr.mock.ctrl.T.Helper()
+       return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "AllURLs", 
reflect.TypeOf((*MockConsumerGroupTopicOption)(nil).AllURLs))
+}
+
+// ConsumerGroup mocks base method.
+func (m *MockConsumerGroupTopicOption) ConsumerGroup() string {
+       m.ctrl.T.Helper()
+       ret := m.ctrl.Call(m, "ConsumerGroup")
+       ret0, _ := ret[0].(string)
+       return ret0
+}
+
+// ConsumerGroup indicates an expected call of ConsumerGroup.
+func (mr *MockConsumerGroupTopicOptionMockRecorder) ConsumerGroup() 
*gomock.Call {
+       mr.mock.ctrl.T.Helper()
+       return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "ConsumerGroup", 
reflect.TypeOf((*MockConsumerGroupTopicOption)(nil).ConsumerGroup))
+}
+
+// DeregisterClient mocks base method.
+func (m *MockConsumerGroupTopicOption) DeregisterClient() 
consumer.DeregisterClient {
+       m.ctrl.T.Helper()
+       ret := m.ctrl.Call(m, "DeregisterClient")
+       ret0, _ := ret[0].(consumer.DeregisterClient)
+       return ret0
+}
+
+// DeregisterClient indicates an expected call of DeregisterClient.
+func (mr *MockConsumerGroupTopicOptionMockRecorder) DeregisterClient() 
*gomock.Call {
+       mr.mock.ctrl.T.Helper()
+       return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, 
"DeregisterClient", 
reflect.TypeOf((*MockConsumerGroupTopicOption)(nil).DeregisterClient))
+}
+
+// GRPCType mocks base method.
+func (m *MockConsumerGroupTopicOption) GRPCType() consts.GRPCType {
+       m.ctrl.T.Helper()
+       ret := m.ctrl.Call(m, "GRPCType")
+       ret0, _ := ret[0].(consts.GRPCType)
+       return ret0
+}
+
+// GRPCType indicates an expected call of GRPCType.
+func (mr *MockConsumerGroupTopicOptionMockRecorder) GRPCType() *gomock.Call {
+       mr.mock.ctrl.T.Helper()
+       return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "GRPCType", 
reflect.TypeOf((*MockConsumerGroupTopicOption)(nil).GRPCType))
+}
+
+// IDCEmiters mocks base method.
+func (m *MockConsumerGroupTopicOption) IDCEmiters() *sync.Map {
+       m.ctrl.T.Helper()
+       ret := m.ctrl.Call(m, "IDCEmiters")
+       ret0, _ := ret[0].(*sync.Map)
+       return ret0
+}
+
+// IDCEmiters indicates an expected call of IDCEmiters.
+func (mr *MockConsumerGroupTopicOptionMockRecorder) IDCEmiters() *gomock.Call {
+       mr.mock.ctrl.T.Helper()
+       return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "IDCEmiters", 
reflect.TypeOf((*MockConsumerGroupTopicOption)(nil).IDCEmiters))
+}
+
+// IDCURLs mocks base method.
+func (m *MockConsumerGroupTopicOption) IDCURLs() *sync.Map {
+       m.ctrl.T.Helper()
+       ret := m.ctrl.Call(m, "IDCURLs")
+       ret0, _ := ret[0].(*sync.Map)
+       return ret0
+}
+
+// IDCURLs indicates an expected call of IDCURLs.
+func (mr *MockConsumerGroupTopicOptionMockRecorder) IDCURLs() *gomock.Call {
+       mr.mock.ctrl.T.Helper()
+       return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "IDCURLs", 
reflect.TypeOf((*MockConsumerGroupTopicOption)(nil).IDCURLs))
+}
+
+// RegisterClient mocks base method.
+func (m *MockConsumerGroupTopicOption) RegisterClient() 
consumer.RegisterClient {
+       m.ctrl.T.Helper()
+       ret := m.ctrl.Call(m, "RegisterClient")
+       ret0, _ := ret[0].(consumer.RegisterClient)
+       return ret0
+}
+
+// RegisterClient indicates an expected call of RegisterClient.
+func (mr *MockConsumerGroupTopicOptionMockRecorder) RegisterClient() 
*gomock.Call {
+       mr.mock.ctrl.T.Helper()
+       return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "RegisterClient", 
reflect.TypeOf((*MockConsumerGroupTopicOption)(nil).RegisterClient))
+}
+
+// Size mocks base method.
+func (m *MockConsumerGroupTopicOption) Size() int {
+       m.ctrl.T.Helper()
+       ret := m.ctrl.Call(m, "Size")
+       ret0, _ := ret[0].(int)
+       return ret0
+}
+
+// Size indicates an expected call of Size.
+func (mr *MockConsumerGroupTopicOptionMockRecorder) Size() *gomock.Call {
+       mr.mock.ctrl.T.Helper()
+       return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "Size", 
reflect.TypeOf((*MockConsumerGroupTopicOption)(nil).Size))
+}
+
+// SubscriptionMode mocks base method.
+func (m *MockConsumerGroupTopicOption) SubscriptionMode() 
pb.Subscription_SubscriptionItem_SubscriptionMode {
+       m.ctrl.T.Helper()
+       ret := m.ctrl.Call(m, "SubscriptionMode")
+       ret0, _ := ret[0].(pb.Subscription_SubscriptionItem_SubscriptionMode)
+       return ret0
+}
+
+// SubscriptionMode indicates an expected call of SubscriptionMode.
+func (mr *MockConsumerGroupTopicOptionMockRecorder) SubscriptionMode() 
*gomock.Call {
+       mr.mock.ctrl.T.Helper()
+       return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, 
"SubscriptionMode", 
reflect.TypeOf((*MockConsumerGroupTopicOption)(nil).SubscriptionMode))
+}
+
+// Topic mocks base method.
+func (m *MockConsumerGroupTopicOption) Topic() string {
+       m.ctrl.T.Helper()
+       ret := m.ctrl.Call(m, "Topic")
+       ret0, _ := ret[0].(string)
+       return ret0
+}
+
+// Topic indicates an expected call of Topic.
+func (mr *MockConsumerGroupTopicOptionMockRecorder) Topic() *gomock.Call {
+       mr.mock.ctrl.T.Helper()
+       return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "Topic", 
reflect.TypeOf((*MockConsumerGroupTopicOption)(nil).Topic))
+}
diff --git 
a/eventmesh-server-go/runtime/core/protocol/grpc/consumer/mocks/consumer_manager.go
 
b/eventmesh-server-go/runtime/core/protocol/grpc/consumer/mocks/consumer_manager.go
new file mode 100644
index 000000000..81a43559b
--- /dev/null
+++ 
b/eventmesh-server-go/runtime/core/protocol/grpc/consumer/mocks/consumer_manager.go
@@ -0,0 +1,132 @@
+// Code generated by MockGen. DO NOT EDIT.
+// Source: 
github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/core/protocol/grpc/consumer
 (interfaces: ConsumerManager)
+
+// Package mocks is a generated GoMock package.
+package mocks
+
+import (
+       reflect "reflect"
+
+       consumer 
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/core/protocol/grpc/consumer"
+       gomock "github.com/golang/mock/gomock"
+)
+
+// MockConsumerManager is a mock of ConsumerManager interface.
+type MockConsumerManager struct {
+       ctrl     *gomock.Controller
+       recorder *MockConsumerManagerMockRecorder
+}
+
+// MockConsumerManagerMockRecorder is the mock recorder for 
MockConsumerManager.
+type MockConsumerManagerMockRecorder struct {
+       mock *MockConsumerManager
+}
+
+// NewMockConsumerManager creates a new mock instance.
+func NewMockConsumerManager(ctrl *gomock.Controller) *MockConsumerManager {
+       mock := &MockConsumerManager{ctrl: ctrl}
+       mock.recorder = &MockConsumerManagerMockRecorder{mock}
+       return mock
+}
+
+// EXPECT returns an object that allows the caller to indicate expected use.
+func (m *MockConsumerManager) EXPECT() *MockConsumerManagerMockRecorder {
+       return m.recorder
+}
+
+// DeRegisterClient mocks base method.
+func (m *MockConsumerManager) DeRegisterClient(arg0 *consumer.GroupClient) 
error {
+       m.ctrl.T.Helper()
+       ret := m.ctrl.Call(m, "DeRegisterClient", arg0)
+       ret0, _ := ret[0].(error)
+       return ret0
+}
+
+// DeRegisterClient indicates an expected call of DeRegisterClient.
+func (mr *MockConsumerManagerMockRecorder) DeRegisterClient(arg0 interface{}) 
*gomock.Call {
+       mr.mock.ctrl.T.Helper()
+       return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, 
"DeRegisterClient", 
reflect.TypeOf((*MockConsumerManager)(nil).DeRegisterClient), arg0)
+}
+
+// GetConsumer mocks base method.
+func (m *MockConsumerManager) GetConsumer(arg0 string) 
(consumer.EventMeshConsumer, error) {
+       m.ctrl.T.Helper()
+       ret := m.ctrl.Call(m, "GetConsumer", arg0)
+       ret0, _ := ret[0].(consumer.EventMeshConsumer)
+       ret1, _ := ret[1].(error)
+       return ret0, ret1
+}
+
+// GetConsumer indicates an expected call of GetConsumer.
+func (mr *MockConsumerManagerMockRecorder) GetConsumer(arg0 interface{}) 
*gomock.Call {
+       mr.mock.ctrl.T.Helper()
+       return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "GetConsumer", 
reflect.TypeOf((*MockConsumerManager)(nil).GetConsumer), arg0)
+}
+
+// RegisterClient mocks base method.
+func (m *MockConsumerManager) RegisterClient(arg0 *consumer.GroupClient) error 
{
+       m.ctrl.T.Helper()
+       ret := m.ctrl.Call(m, "RegisterClient", arg0)
+       ret0, _ := ret[0].(error)
+       return ret0
+}
+
+// RegisterClient indicates an expected call of RegisterClient.
+func (mr *MockConsumerManagerMockRecorder) RegisterClient(arg0 interface{}) 
*gomock.Call {
+       mr.mock.ctrl.T.Helper()
+       return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "RegisterClient", 
reflect.TypeOf((*MockConsumerManager)(nil).RegisterClient), arg0)
+}
+
+// RestartConsumer mocks base method.
+func (m *MockConsumerManager) RestartConsumer(arg0 string) error {
+       m.ctrl.T.Helper()
+       ret := m.ctrl.Call(m, "RestartConsumer", arg0)
+       ret0, _ := ret[0].(error)
+       return ret0
+}
+
+// RestartConsumer indicates an expected call of RestartConsumer.
+func (mr *MockConsumerManagerMockRecorder) RestartConsumer(arg0 interface{}) 
*gomock.Call {
+       mr.mock.ctrl.T.Helper()
+       return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, 
"RestartConsumer", reflect.TypeOf((*MockConsumerManager)(nil).RestartConsumer), 
arg0)
+}
+
+// Start mocks base method.
+func (m *MockConsumerManager) Start() error {
+       m.ctrl.T.Helper()
+       ret := m.ctrl.Call(m, "Start")
+       ret0, _ := ret[0].(error)
+       return ret0
+}
+
+// Start indicates an expected call of Start.
+func (mr *MockConsumerManagerMockRecorder) Start() *gomock.Call {
+       mr.mock.ctrl.T.Helper()
+       return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "Start", 
reflect.TypeOf((*MockConsumerManager)(nil).Start))
+}
+
+// Stop mocks base method.
+func (m *MockConsumerManager) Stop() error {
+       m.ctrl.T.Helper()
+       ret := m.ctrl.Call(m, "Stop")
+       ret0, _ := ret[0].(error)
+       return ret0
+}
+
+// Stop indicates an expected call of Stop.
+func (mr *MockConsumerManagerMockRecorder) Stop() *gomock.Call {
+       mr.mock.ctrl.T.Helper()
+       return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "Stop", 
reflect.TypeOf((*MockConsumerManager)(nil).Stop))
+}
+
+// UpdateClientTime mocks base method.
+func (m *MockConsumerManager) UpdateClientTime(arg0 *consumer.GroupClient) {
+       m.ctrl.T.Helper()
+       m.ctrl.Call(m, "UpdateClientTime", arg0)
+}
+
+// UpdateClientTime indicates an expected call of UpdateClientTime.
+func (mr *MockConsumerManagerMockRecorder) UpdateClientTime(arg0 interface{}) 
*gomock.Call {
+       mr.mock.ctrl.T.Helper()
+       return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, 
"UpdateClientTime", 
reflect.TypeOf((*MockConsumerManager)(nil).UpdateClientTime), arg0)
+}
diff --git 
a/eventmesh-server-go/runtime/core/protocol/grpc/consumer_service_test.go 
b/eventmesh-server-go/runtime/core/protocol/grpc/consumer_service_test.go
deleted file mode 100644
index eddeb17f2..000000000
--- a/eventmesh-server-go/runtime/core/protocol/grpc/consumer_service_test.go
+++ /dev/null
@@ -1,154 +0,0 @@
-// 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 grpc
-
-import (
-       "context"
-       "testing"
-       "time"
-
-       "github.com/stretchr/testify/assert"
-
-       "github.com/apache/incubator-eventmesh/eventmesh-server-go/pkg/util"
-       
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/proto/pb"
-)
-
-//     {
-//        "header": {
-//            "env": "11",
-//            "region": "sh",
-//            "idc": "test-idc",
-//            "ip": "169.254.45.15",
-//            "pid": "12345",
-//            "sys": "test",
-//            "username": "username",
-//            "password": "password",
-//            "language": "go",
-//            "protocolType": "cloudevents",
-//            "protocolVersion": "1.0",
-//            "protocolDesc": "nil"
-//        },
-//        "subscriptionItems": [
-//            {
-//                "topic":"grpc-topic",
-//                "mode":"CLUSTERING",
-//                "type":"SYNC"
-//            }
-//        ],
-//        "consumerGroup": "test-grpc-group",
-//        "url": "http://127.0.0.1:18080";
-//     }
-func Test_Subscribe(t *testing.T) {
-       cli := newTestClient(t)
-       assert.NotNil(t, cli)
-       resp, err := cli.consumerClient.Subscribe(context.TODO(), 
&pb.Subscription{
-               Header: &pb.RequestHeader{
-                       Env:             "grpc-env",
-                       Region:          "sh",
-                       Idc:             "idc-sh",
-                       Ip:              util.GetIP(),
-                       Pid:             util.PID(),
-                       Sys:             "grpc-sys",
-                       Username:        "grpc-username",
-                       Password:        "grpc-passwd",
-                       Language:        "Go",
-                       ProtocolType:    "cloudevents",
-                       ProtocolVersion: "1.0",
-                       ProtocolDesc:    "cloudevents",
-               },
-               ConsumerGroup: "grpc-stream-consumergroup",
-               SubscriptionItems: []*pb.Subscription_SubscriptionItem{
-                       {
-                               Topic: _testWebhookTopic,
-                               Mode:  
pb.Subscription_SubscriptionItem_CLUSTERING,
-                               Type:  pb.Subscription_SubscriptionItem_SYNC,
-                       },
-               },
-               Url: "http://127.0.0.1:18080/onmessage";,
-       })
-       assert.NoError(t, err)
-       assert.NotNil(t, cli)
-       t.Log(resp.String())
-       assert.Equal(t, resp.RespCode, 0)
-}
-
-func Test_unsubscribe(t *testing.T) {
-       cli := newTestClient(t)
-       assert.NotNil(t, cli)
-       resp, err := cli.consumerClient.Unsubscribe(context.TODO(), 
&pb.Subscription{
-               Header: &pb.RequestHeader{
-                       Env:             "grpc-env",
-                       Region:          "sh",
-                       Idc:             "idc-sh",
-                       Ip:              util.GetIP(),
-                       Pid:             util.PID(),
-                       Sys:             "grpc-sys",
-                       Username:        "grpc-username",
-                       Password:        "grpc-passwd",
-                       Language:        "Go",
-                       ProtocolType:    "cloudevents",
-                       ProtocolVersion: "1.0",
-                       ProtocolDesc:    "cloudevents",
-               },
-               ConsumerGroup: "grpc-stream-consumergroup",
-               SubscriptionItems: []*pb.Subscription_SubscriptionItem{
-                       {
-                               Topic: _testWebhookTopic,
-                               Mode:  
pb.Subscription_SubscriptionItem_CLUSTERING,
-                               Type:  pb.Subscription_SubscriptionItem_SYNC,
-                       },
-               },
-               Url: "http://127.0.0.1:18080/onmessage";,
-       })
-       assert.NoError(t, err)
-       assert.NotNil(t, cli)
-       t.Log(resp.String())
-       assert.Equal(t, resp.RespCode, 0)
-}
-
-func Test_subscribeStream(t *testing.T) {
-       cli := newTestClient(t)
-       assert.NotNil(t, cli)
-
-       stream, err := cli.consumerClient.SubscribeStream(context.TODO())
-       assert.NoError(t, err)
-       err = stream.Send(&pb.Subscription{
-               Header: &pb.RequestHeader{
-                       Env:             "grpc-env",
-                       Region:          "sh",
-                       Idc:             "idc-sh",
-                       Ip:              util.GetIP(),
-                       Pid:             util.PID(),
-                       Sys:             "grpc-sys",
-                       Username:        "grpc-username",
-                       Password:        "grpc-passwd",
-                       Language:        "Go",
-                       ProtocolType:    "cloudevents",
-                       ProtocolVersion: "1.0",
-                       ProtocolDesc:    "cloudevents",
-               },
-               ConsumerGroup: "grpc-stream-consumergroup",
-               SubscriptionItems: []*pb.Subscription_SubscriptionItem{
-                       {
-                               Topic: _testStreamTopic,
-                               Mode:  
pb.Subscription_SubscriptionItem_CLUSTERING,
-                               Type:  pb.Subscription_SubscriptionItem_ASYNC,
-                       },
-               },
-       })
-       assert.NoError(t, err)
-       time.Sleep(time.Hour)
-}
diff --git a/eventmesh-server-go/runtime/core/protocol/grpc/context.go 
b/eventmesh-server-go/runtime/core/protocol/grpc/context.go
deleted file mode 100644
index 04655b012..000000000
--- a/eventmesh-server-go/runtime/core/protocol/grpc/context.go
+++ /dev/null
@@ -1,101 +0,0 @@
-// 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 grpc
-
-import (
-       "context"
-       "github.com/apache/incubator-eventmesh/eventmesh-server-go/config"
-       "github.com/apache/incubator-eventmesh/eventmesh-server-go/log"
-       
"github.com/apache/incubator-eventmesh/eventmesh-server-go/pkg/naming/registry"
-       
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/proto/pb"
-       cloudv2 "github.com/cloudevents/sdk-go/v2"
-       "golang.org/x/time/rate"
-       "time"
-)
-
-// GRPCContext grpc server api, used to handle the client
-type GRPCContext struct {
-       ConsumerMgr *ConsumerManager
-       ProducerMgr *ProducerManager
-       RateLimiter *rate.Limiter
-       Registry    registry.Registry
-}
-
-// New create new grpc server
-func New() (*GRPCContext, error) {
-       log.Infof("create new grpc serer")
-       msgReqPerSeconds := 
config.GlobalConfig().Server.GRPCOption.MsgReqNumPerSecond
-       limiter := rate.NewLimiter(rate.Limit(msgReqPerSeconds), 10)
-       cmgr, err := NewConsumerManager()
-       if err != nil {
-               return nil, err
-       }
-       pmgr, err := NewProducerManager()
-       if err != nil {
-               return nil, err
-       }
-
-       registryName := config.GlobalConfig().Server.GRPCOption.RegistryName
-       regis := registry.Get(registryName)
-       return &GRPCContext{
-               ConsumerMgr: cmgr,
-               ProducerMgr: pmgr,
-               RateLimiter: limiter,
-               Registry:    regis,
-       }, nil
-}
-
-func (g *GRPCContext) Start() error {
-       if err := g.ProducerMgr.Start(); err != nil {
-               return err
-       }
-       if err := g.ConsumerMgr.Start(); err != nil {
-               return err
-       }
-       //if err := g.Registry.Start(); err != nil {
-       //      return err
-       //}
-       return nil
-}
-
-//
-//func (g *GRPCContext) SendResp(code *grpc.StatusCode) {
-//     resp := &pb.Response{
-//             RespCode: code.RetCode,
-//             RespMsg:  code.ErrMsg,
-//             RespTime: fmt.Sprintf("%v", time.Now().UnixMilli()),
-//     }
-//}
-
-type MessageContext struct {
-       MsgRandomNo      string
-       SubscriptionMode pb.Subscription_SubscriptionItem_SubscriptionMode
-       GrpcType         GRPCType
-       ConsumerGroup    string
-       Event            *cloudv2.Event
-       TopicConfig      ConsumerGroupTopicOption
-       // channel for server
-       Consumer *EventMeshConsumer
-}
-
-// SendMessageContext context in produce message
-type SendMessageContext struct {
-       Ctx         context.Context
-       Event       *cloudv2.Event
-       BizSeqNO    string
-       ProducerAPI *EventMeshProducer
-       CreateTime  time.Time
-}
diff --git a/eventmesh-server-go/runtime/core/protocol/grpc/emitter.go 
b/eventmesh-server-go/runtime/core/protocol/grpc/emitter/emitter.go
similarity index 69%
rename from eventmesh-server-go/runtime/core/protocol/grpc/emitter.go
rename to eventmesh-server-go/runtime/core/protocol/grpc/emitter/emitter.go
index 7ccc3a065..da3fba878 100644
--- a/eventmesh-server-go/runtime/core/protocol/grpc/emitter.go
+++ b/eventmesh-server-go/runtime/core/protocol/grpc/emitter/emitter.go
@@ -13,18 +13,29 @@
 // See the License for the specific language governing permissions and
 // limitations under the License.
 
-package grpc
+package emitter
 
 import (
        
"github.com/apache/incubator-eventmesh/eventmesh-server-go/pkg/common/protocol/grpc"
        
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/proto/pb"
 )
 
-type EventEmitter struct {
+//go:generate mockgen -destination ./mocks/emitter.go -package mocks 
github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/core/protocol/grpc/emitter
 EventEmitter
+type EventEmitter interface {
+       SendStreamResp(hdr *pb.RequestHeader, code *grpc.StatusCode) error
+}
+
+type eventEmitter struct {
        emitter pb.ConsumerService_SubscribeStreamServer
 }
 
-func (e *EventEmitter) sendStreamResp(hdr *pb.RequestHeader, code 
*grpc.StatusCode) error {
+func NewEventEmitter(stream pb.ConsumerService_SubscribeStreamServer) 
EventEmitter {
+       return &eventEmitter{
+               emitter: stream,
+       }
+}
+
+func (e *eventEmitter) SendStreamResp(hdr *pb.RequestHeader, code 
*grpc.StatusCode) error {
        return e.emitter.Send(&pb.SimpleMessage{
                Header:  hdr,
                Content: code.ToJSONString(),
diff --git 
a/eventmesh-server-go/runtime/core/protocol/grpc/emitter/mocks/emitter.go 
b/eventmesh-server-go/runtime/core/protocol/grpc/emitter/mocks/emitter.go
new file mode 100644
index 000000000..8d52afb44
--- /dev/null
+++ b/eventmesh-server-go/runtime/core/protocol/grpc/emitter/mocks/emitter.go
@@ -0,0 +1,50 @@
+// Code generated by MockGen. DO NOT EDIT.
+// Source: 
github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/core/protocol/grpc/emitter
 (interfaces: EventEmitter)
+
+// Package mocks is a generated GoMock package.
+package mocks
+
+import (
+       reflect "reflect"
+
+       grpc 
"github.com/apache/incubator-eventmesh/eventmesh-server-go/pkg/common/protocol/grpc"
+       pb 
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/proto/pb"
+       gomock "github.com/golang/mock/gomock"
+)
+
+// MockEventEmitter is a mock of EventEmitter interface.
+type MockEventEmitter struct {
+       ctrl     *gomock.Controller
+       recorder *MockEventEmitterMockRecorder
+}
+
+// MockEventEmitterMockRecorder is the mock recorder for MockEventEmitter.
+type MockEventEmitterMockRecorder struct {
+       mock *MockEventEmitter
+}
+
+// NewMockEventEmitter creates a new mock instance.
+func NewMockEventEmitter(ctrl *gomock.Controller) *MockEventEmitter {
+       mock := &MockEventEmitter{ctrl: ctrl}
+       mock.recorder = &MockEventEmitterMockRecorder{mock}
+       return mock
+}
+
+// EXPECT returns an object that allows the caller to indicate expected use.
+func (m *MockEventEmitter) EXPECT() *MockEventEmitterMockRecorder {
+       return m.recorder
+}
+
+// SendStreamResp mocks base method.
+func (m *MockEventEmitter) SendStreamResp(arg0 *pb.RequestHeader, arg1 
*grpc.StatusCode) error {
+       m.ctrl.T.Helper()
+       ret := m.ctrl.Call(m, "SendStreamResp", arg0, arg1)
+       ret0, _ := ret[0].(error)
+       return ret0
+}
+
+// SendStreamResp indicates an expected call of SendStreamResp.
+func (mr *MockEventEmitterMockRecorder) SendStreamResp(arg0, arg1 interface{}) 
*gomock.Call {
+       mr.mock.ctrl.T.Helper()
+       return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "SendStreamResp", 
reflect.TypeOf((*MockEventEmitter)(nil).SendStreamResp), arg0, arg1)
+}
diff --git a/eventmesh-server-go/runtime/core/protocol/grpc/fake_client.go 
b/eventmesh-server-go/runtime/core/protocol/grpc/fake_client.go
deleted file mode 100644
index 49ea3fd1a..000000000
--- a/eventmesh-server-go/runtime/core/protocol/grpc/fake_client.go
+++ /dev/null
@@ -1,113 +0,0 @@
-// 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 grpc
-
-import (
-       "context"
-       "fmt"
-       
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/proto/pb"
-       "github.com/gin-gonic/gin"
-       "github.com/stretchr/testify/assert"
-       "google.golang.org/grpc"
-       "net/http"
-       "testing"
-)
-
-var (
-       _testStreamTopic  = "test-grpc-stream-topic"
-       _testWebhookTopic = "test-grpc-webhook-topic"
-)
-
-// testClient test client to do the grpc service test
-type testClient struct {
-       consumerClient  pb.ConsumerServiceClient
-       producerClient  pb.PublisherServiceClient
-       heartbeatClient pb.HeartbeatServiceClient
-}
-
-func newTestClient(t *testing.T) *testClient {
-       clientCoon, err := grpc.Dial("127.0.0.1:10010", grpc.WithInsecure())
-       assert.NoError(t, err)
-       consumer := pb.NewConsumerServiceClient(clientCoon)
-       producer := pb.NewPublisherServiceClient(clientCoon)
-       heartheat := pb.NewHeartbeatServiceClient(clientCoon)
-       return &testClient{
-               consumerClient:  consumer,
-               producerClient:  producer,
-               heartbeatClient: heartheat,
-       }
-}
-
-func (c *testClient) createGRPCServer(t *testing.T) {
-
-}
-
-func (c *testClient) createGRPCClient(t *testing.T) {
-
-}
-
-func (c *testClient) startWebhookServer(t *testing.T) error {
-       router := gin.Default()
-
-       // TestResponse indicate the http response
-       type TestResponse struct {
-               RetCode string `json:"retCode"`
-               ErrMsg  string `json:"errMsg"`
-       }
-
-       router.Any("/*anypath", func(c *gin.Context) {
-               c.JSON(http.StatusOK, &TestResponse{
-                       RetCode: "0",
-                       ErrMsg:  "OK",
-               })
-       })
-       go func() {
-               if err := router.Run(fmt.Sprintf(":%d", 18080)); err != nil {
-                       panic(err)
-               }
-       }()
-       return nil
-}
-
-func Subscribe(ctx context.Context, sub *pb.Subscription) (*pb.Response, 
error) {
-       return nil, nil
-}
-
-func SubscribeStream(sss pb.ConsumerService_SubscribeStreamServer) error {
-       return nil
-}
-
-func Unsubscribe(ctx context.Context, unsub *pb.Subscription) (*pb.Response, 
error) {
-       return nil, nil
-}
-
-// producer service
-func Publish(ctx context.Context, msg *pb.SimpleMessage) (*pb.Response, error) 
{
-       return nil, nil
-}
-
-func RequestReply(ctx context.Context, msg *pb.SimpleMessage) 
(*pb.SimpleMessage, error) {
-       return nil, nil
-}
-
-func BatchPublish(ctx context.Context, msg *pb.BatchMessage) (*pb.Response, 
error) {
-       return nil, nil
-}
-
-// heartbeat service
-func Heartbeat(ctx context.Context, msg *pb.Heartbeat) (*pb.Response, error) {
-       return nil, nil
-}
diff --git a/eventmesh-server-go/runtime/core/protocol/grpc/generate_mocks.sh 
b/eventmesh-server-go/runtime/core/protocol/grpc/generate_mocks.sh
new file mode 100644
index 000000000..9590b7def
--- /dev/null
+++ b/eventmesh-server-go/runtime/core/protocol/grpc/generate_mocks.sh
@@ -0,0 +1,12 @@
+#!/usr/bin/env bash
+
+mockgen.exe -package=mocks -destination=./mocks/heartbeat_service.go 
github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/proto/pb 
HeartbeatServiceServer
+mockgen.exe -package=mocks -destination=./mocks/consumer_service.go 
github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/proto/pb 
ConsumerServiceServer
+mockgen.exe -package=mocks -destination=./mocks/producer_service.go 
github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/proto/pb 
PublisherServiceServer
+
+mockgen.exe -package=mocks -destination=./mocks/consumer_manager.go 
github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/core/protocol/grpc
 ConsumerManager
+mockgen.exe -package=mocks -destination=./mocks/consumer_mesh.go 
github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/core/protocol/grpc
 EventMeshConsumer
+mockgen.exe -package=mocks -destination=./mocks/message_handler.go 
github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/core/protocol/grpc
 MessageHandler
+mockgen.exe -package=mocks -destination=./mocks/processor.go 
github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/core/protocol/grpc
 Processor
+mockgen.exe -package=mocks -destination=./mocks/producer_manager.go 
github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/core/protocol/grpc
 ProducerManager
+mockgen.exe -package=mocks -destination=./mocks/producer_mesh.go 
github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/core/protocol/grpc
 EventMeshProducer
diff --git 
a/eventmesh-server-go/runtime/core/protocol/grpc/heartbeat/heartbeat_processor.go
 
b/eventmesh-server-go/runtime/core/protocol/grpc/heartbeat/heartbeat_processor.go
new file mode 100644
index 000000000..9d99605b7
--- /dev/null
+++ 
b/eventmesh-server-go/runtime/core/protocol/grpc/heartbeat/heartbeat_processor.go
@@ -0,0 +1,61 @@
+package heartbeat
+
+import (
+       "fmt"
+       "github.com/apache/incubator-eventmesh/eventmesh-server-go/log"
+       
"github.com/apache/incubator-eventmesh/eventmesh-server-go/pkg/common/protocol/grpc"
+       
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/core/protocol/grpc/consumer"
+       
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/core/protocol/grpc/validator"
+       
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/proto/pb"
+       "time"
+)
+
+type Processor interface {
+       Heartbeat(consumerMgr consumer.ConsumerManager, msg *pb.Heartbeat) 
(*pb.Response, error)
+}
+
+type processor struct {
+}
+
+func NewProcessor() Processor {
+       return &processor{}
+}
+
+func (p *processor) Heartbeat(consumerMgr consumer.ConsumerManager, msg 
*pb.Heartbeat) (*pb.Response, error) {
+       hdr := msg.Header
+       if err := validator.ValidateHeader(hdr); err != nil {
+               log.Warnf("invalid header:%v", err)
+               return buildPBResponse(grpc.EVENTMESH_PROTOCOL_HEADER_ERR), err
+       }
+       if err := validator.ValidateHeartBeat(msg); err != nil {
+               log.Warnf("invalid body:%v", err)
+               return buildPBResponse(grpc.EVENTMESH_PROTOCOL_BODY_ERR), err
+       }
+       if msg.ClientType != pb.Heartbeat_SUB {
+               log.Warnf("client type err, not sub")
+               return buildPBResponse(grpc.EVENTMESH_Heartbeat_Protocol_ERR), 
fmt.Errorf("protocol not sub")
+       }
+       consumerGroup := msg.ConsumerGroup
+       for _, item := range msg.HeartbeatItems {
+               cli := &consumer.GroupClient{
+                       ENV:           hdr.Env,
+                       IDC:           hdr.Idc,
+                       SYS:           hdr.Sys,
+                       IP:            hdr.Ip,
+                       PID:           hdr.Pid,
+                       ConsumerGroup: consumerGroup,
+                       Topic:         item.Topic,
+                       LastUPTime:    time.Now(),
+               }
+               consumerMgr.UpdateClientTime(cli)
+       }
+       return buildPBResponse(grpc.SUCCESS), nil
+}
+
+func buildPBResponse(code *grpc.StatusCode) *pb.Response {
+       return &pb.Response{
+               RespCode: code.RetCode,
+               RespMsg:  code.ErrMsg,
+               RespTime: fmt.Sprintf("%v", time.Now().UnixMilli()),
+       }
+}
diff --git 
a/eventmesh-server-go/runtime/core/protocol/grpc/heartbeat_service.go 
b/eventmesh-server-go/runtime/core/protocol/grpc/heartbeat/heartbeat_service.go
similarity index 82%
rename from eventmesh-server-go/runtime/core/protocol/grpc/heartbeat_service.go
rename to 
eventmesh-server-go/runtime/core/protocol/grpc/heartbeat/heartbeat_service.go
index 64d540b32..9991bcbd9 100644
--- a/eventmesh-server-go/runtime/core/protocol/grpc/heartbeat_service.go
+++ 
b/eventmesh-server-go/runtime/core/protocol/grpc/heartbeat/heartbeat_service.go
@@ -13,12 +13,13 @@
 // See the License for the specific language governing permissions and
 // limitations under the License.
 
-package grpc
+package heartbeat
 
 import (
        "context"
        "github.com/apache/incubator-eventmesh/eventmesh-server-go/config"
        "github.com/apache/incubator-eventmesh/eventmesh-server-go/log"
+       
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/core/protocol/grpc/consumer"
        
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/proto/pb"
        "github.com/panjf2000/ants/v2"
        "time"
@@ -26,19 +27,19 @@ import (
 
 type HeartbeatService struct {
        pb.UnimplementedHeartbeatServiceServer
-       gctx *GRPCContext
-       pool *ants.Pool
+       consumerMgr consumer.ConsumerManager
+       pool        *ants.Pool
 }
 
-func NewHeartbeatServiceServer(gctx *GRPCContext) (*HeartbeatService, error) {
+func NewHeartbeatServiceServer(consumerMgr consumer.ConsumerManager) 
(*HeartbeatService, error) {
        sp := config.GlobalConfig().Server.GRPCOption.SubscribePoolSize
        pl, err := ants.NewPool(sp)
        if err != nil {
                return nil, err
        }
        return &HeartbeatService{
-               gctx: gctx,
-               pool: pl,
+               consumerMgr: consumerMgr,
+               pool:        pl,
        }, nil
 }
 
@@ -51,7 +52,7 @@ func (h *HeartbeatService) Heartbeat(ctx context.Context, hb 
*pb.Heartbeat) (*pb
                err     error
        )
        h.pool.Submit(func() {
-               resp, err = ProcessHeartbeat(h.gctx, hb)
+               resp, err = NewProcessor().Heartbeat(h.consumerMgr, hb)
                errChan <- err
        })
        select {
diff --git 
a/eventmesh-server-go/runtime/core/protocol/grpc/mocks/consumer_service.go 
b/eventmesh-server-go/runtime/core/protocol/grpc/mocks/consumer_service.go
new file mode 100644
index 000000000..4623f0cdf
--- /dev/null
+++ b/eventmesh-server-go/runtime/core/protocol/grpc/mocks/consumer_service.go
@@ -0,0 +1,92 @@
+// Code generated by MockGen. DO NOT EDIT.
+// Source: 
github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/proto/pb 
(interfaces: ConsumerServiceServer)
+
+// Package mocks is a generated GoMock package.
+package mocks
+
+import (
+       context "context"
+       reflect "reflect"
+
+       pb 
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/proto/pb"
+       gomock "github.com/golang/mock/gomock"
+)
+
+// MockConsumerServiceServer is a mock of ConsumerServiceServer interface.
+type MockConsumerServiceServer struct {
+       ctrl     *gomock.Controller
+       recorder *MockConsumerServiceServerMockRecorder
+}
+
+// MockConsumerServiceServerMockRecorder is the mock recorder for 
MockConsumerServiceServer.
+type MockConsumerServiceServerMockRecorder struct {
+       mock *MockConsumerServiceServer
+}
+
+// NewMockConsumerServiceServer creates a new mock instance.
+func NewMockConsumerServiceServer(ctrl *gomock.Controller) 
*MockConsumerServiceServer {
+       mock := &MockConsumerServiceServer{ctrl: ctrl}
+       mock.recorder = &MockConsumerServiceServerMockRecorder{mock}
+       return mock
+}
+
+// EXPECT returns an object that allows the caller to indicate expected use.
+func (m *MockConsumerServiceServer) EXPECT() 
*MockConsumerServiceServerMockRecorder {
+       return m.recorder
+}
+
+// Subscribe mocks base method.
+func (m *MockConsumerServiceServer) Subscribe(arg0 context.Context, arg1 
*pb.Subscription) (*pb.Response, error) {
+       m.ctrl.T.Helper()
+       ret := m.ctrl.Call(m, "Subscribe", arg0, arg1)
+       ret0, _ := ret[0].(*pb.Response)
+       ret1, _ := ret[1].(error)
+       return ret0, ret1
+}
+
+// Subscribe indicates an expected call of Subscribe.
+func (mr *MockConsumerServiceServerMockRecorder) Subscribe(arg0, arg1 
interface{}) *gomock.Call {
+       mr.mock.ctrl.T.Helper()
+       return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "Subscribe", 
reflect.TypeOf((*MockConsumerServiceServer)(nil).Subscribe), arg0, arg1)
+}
+
+// SubscribeStream mocks base method.
+func (m *MockConsumerServiceServer) SubscribeStream(arg0 
pb.ConsumerService_SubscribeStreamServer) error {
+       m.ctrl.T.Helper()
+       ret := m.ctrl.Call(m, "SubscribeStream", arg0)
+       ret0, _ := ret[0].(error)
+       return ret0
+}
+
+// SubscribeStream indicates an expected call of SubscribeStream.
+func (mr *MockConsumerServiceServerMockRecorder) SubscribeStream(arg0 
interface{}) *gomock.Call {
+       mr.mock.ctrl.T.Helper()
+       return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, 
"SubscribeStream", 
reflect.TypeOf((*MockConsumerServiceServer)(nil).SubscribeStream), arg0)
+}
+
+// Unsubscribe mocks base method.
+func (m *MockConsumerServiceServer) Unsubscribe(arg0 context.Context, arg1 
*pb.Subscription) (*pb.Response, error) {
+       m.ctrl.T.Helper()
+       ret := m.ctrl.Call(m, "Unsubscribe", arg0, arg1)
+       ret0, _ := ret[0].(*pb.Response)
+       ret1, _ := ret[1].(error)
+       return ret0, ret1
+}
+
+// Unsubscribe indicates an expected call of Unsubscribe.
+func (mr *MockConsumerServiceServerMockRecorder) Unsubscribe(arg0, arg1 
interface{}) *gomock.Call {
+       mr.mock.ctrl.T.Helper()
+       return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "Unsubscribe", 
reflect.TypeOf((*MockConsumerServiceServer)(nil).Unsubscribe), arg0, arg1)
+}
+
+// mustEmbedUnimplementedConsumerServiceServer mocks base method.
+func (m *MockConsumerServiceServer) 
mustEmbedUnimplementedConsumerServiceServer() {
+       m.ctrl.T.Helper()
+       m.ctrl.Call(m, "mustEmbedUnimplementedConsumerServiceServer")
+}
+
+// mustEmbedUnimplementedConsumerServiceServer indicates an expected call of 
mustEmbedUnimplementedConsumerServiceServer.
+func (mr *MockConsumerServiceServerMockRecorder) 
mustEmbedUnimplementedConsumerServiceServer() *gomock.Call {
+       mr.mock.ctrl.T.Helper()
+       return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, 
"mustEmbedUnimplementedConsumerServiceServer", 
reflect.TypeOf((*MockConsumerServiceServer)(nil).mustEmbedUnimplementedConsumerServiceServer))
+}
diff --git a/eventmesh-server-go/runtime/core/protocol/grpc/producer_group.go 
b/eventmesh-server-go/runtime/core/protocol/grpc/producer/message_context.go
similarity index 72%
copy from eventmesh-server-go/runtime/core/protocol/grpc/producer_group.go
copy to 
eventmesh-server-go/runtime/core/protocol/grpc/producer/message_context.go
index 519840507..fdfb84dd6 100644
--- a/eventmesh-server-go/runtime/core/protocol/grpc/producer_group.go
+++ b/eventmesh-server-go/runtime/core/protocol/grpc/producer/message_context.go
@@ -13,8 +13,19 @@
 // See the License for the specific language governing permissions and
 // limitations under the License.
 
-package grpc
+package producer
 
-type ProducerGroupConfig struct {
-       GroupName string `json:"groupName"`
+import (
+       "context"
+       cloudv2 "github.com/cloudevents/sdk-go/v2"
+       "time"
+)
+
+// SendMessageContext context in produce message
+type SendMessageContext struct {
+       Ctx         context.Context
+       Event       *cloudv2.Event
+       BizSeqNO    string
+       ProducerAPI EventMeshProducer
+       CreateTime  time.Time
 }
diff --git a/eventmesh-server-go/runtime/core/protocol/grpc/producer_group.go 
b/eventmesh-server-go/runtime/core/protocol/grpc/producer/producer_group.go
similarity index 98%
rename from eventmesh-server-go/runtime/core/protocol/grpc/producer_group.go
rename to 
eventmesh-server-go/runtime/core/protocol/grpc/producer/producer_group.go
index 519840507..bb877f397 100644
--- a/eventmesh-server-go/runtime/core/protocol/grpc/producer_group.go
+++ b/eventmesh-server-go/runtime/core/protocol/grpc/producer/producer_group.go
@@ -13,7 +13,7 @@
 // See the License for the specific language governing permissions and
 // limitations under the License.
 
-package grpc
+package producer
 
 type ProducerGroupConfig struct {
        GroupName string `json:"groupName"`
diff --git a/eventmesh-server-go/runtime/core/protocol/grpc/producer_manager.go 
b/eventmesh-server-go/runtime/core/protocol/grpc/producer/producer_manager.go
similarity index 69%
rename from eventmesh-server-go/runtime/core/protocol/grpc/producer_manager.go
rename to 
eventmesh-server-go/runtime/core/protocol/grpc/producer/producer_manager.go
index 49f2b6383..17e6e3cf6 100644
--- a/eventmesh-server-go/runtime/core/protocol/grpc/producer_manager.go
+++ 
b/eventmesh-server-go/runtime/core/protocol/grpc/producer/producer_manager.go
@@ -13,29 +13,36 @@
 // See the License for the specific language governing permissions and
 // limitations under the License.
 
-package grpc
+package producer
 
 import (
        "github.com/apache/incubator-eventmesh/eventmesh-server-go/log"
        "sync"
 )
 
-// ProducerManager manger for all producer
-type ProducerManager struct {
+type ProducerManager interface {
+       GetProducer(groupName string) (EventMeshProducer, error)
+       CreateProducer(producerGroupConfig *ProducerGroupConfig) 
(EventMeshProducer, error)
+       Start() error
+       Shutdown() error
+}
+
+// producerManager manger for all producer
+type producerManager struct {
        // EventMeshProducers {groupName, *EventMeshProducer}
        EventMeshProducers *sync.Map
 }
 
-func NewProducerManager() (*ProducerManager, error) {
-       return &ProducerManager{
+func NewProducerManager() (ProducerManager, error) {
+       return &producerManager{
                EventMeshProducers: new(sync.Map),
        }, nil
 }
 
-func (m *ProducerManager) GetProducer(groupName string) (*EventMeshProducer, 
error) {
+func (m *producerManager) GetProducer(groupName string) (EventMeshProducer, 
error) {
        p, ok := m.EventMeshProducers.Load(groupName)
        if ok {
-               return p.(*EventMeshProducer), nil
+               return p.(EventMeshProducer), nil
        }
        pgc := &ProducerGroupConfig{GroupName: groupName}
        pg, err := m.CreateProducer(pgc)
@@ -45,10 +52,10 @@ func (m *ProducerManager) GetProducer(groupName string) 
(*EventMeshProducer, err
        return pg, nil
 }
 
-func (m *ProducerManager) CreateProducer(producerGroupConfig 
*ProducerGroupConfig) (*EventMeshProducer, error) {
+func (m *producerManager) CreateProducer(producerGroupConfig 
*ProducerGroupConfig) (EventMeshProducer, error) {
        val, ok := m.EventMeshProducers.Load(producerGroupConfig.GroupName)
        if ok {
-               return val.(*EventMeshProducer), nil
+               return val.(EventMeshProducer), nil
        }
        pg, err := NewEventMeshProducer(producerGroupConfig)
        if err != nil {
@@ -58,16 +65,16 @@ func (m *ProducerManager) 
CreateProducer(producerGroupConfig *ProducerGroupConfi
        return pg, nil
 }
 
-func (m *ProducerManager) Start() error {
+func (m *producerManager) Start() error {
        log.Infof("start producer manager")
        return nil
 }
 
-func (m *ProducerManager) Shutdown() error {
+func (m *producerManager) Shutdown() error {
        log.Infof("shutdown producer manager")
 
        m.EventMeshProducers.Range(func(key, value any) bool {
-               pg := value.(*EventMeshProducer)
+               pg := value.(EventMeshProducer)
                if err := pg.Shutdown(); err != nil {
                        log.Infof("shutdown eventmesh producer:%v, err:%v", 
key, err)
                }
diff --git a/eventmesh-server-go/runtime/core/protocol/grpc/producer_mesh.go 
b/eventmesh-server-go/runtime/core/protocol/grpc/producer/producer_mesh.go
similarity index 75%
rename from eventmesh-server-go/runtime/core/protocol/grpc/producer_mesh.go
rename to 
eventmesh-server-go/runtime/core/protocol/grpc/producer/producer_mesh.go
index ddf925e90..b9f21890f 100644
--- a/eventmesh-server-go/runtime/core/protocol/grpc/producer_mesh.go
+++ b/eventmesh-server-go/runtime/core/protocol/grpc/producer/producer_mesh.go
@@ -13,7 +13,7 @@
 // See the License for the specific language governing permissions and
 // limitations under the License.
 
-package grpc
+package producer
 
 import (
        "fmt"
@@ -26,13 +26,23 @@ import (
        "time"
 )
 
-type EventMeshProducer struct {
+type EventMeshProducer interface {
+       Send(sctx SendMessageContext, callback *connector.SendCallback) error
+       Request(sctx SendMessageContext, callback 
*connector.RequestReplyCallback, timeout time.Duration) error
+       Reply(sctx SendMessageContext, callback *connector.SendCallback) error
+       Start() error
+       Shutdown() error
+       Status() consts.ServiceState
+       String() string
+}
+
+type eventMeshProducer struct {
        cfg          *ProducerGroupConfig
        producer     *wrapper.Producer
        ServiceState consts.ServiceState
 }
 
-func NewEventMeshProducer(cfg *ProducerGroupConfig) (*EventMeshProducer, 
error) {
+func NewEventMeshProducer(cfg *ProducerGroupConfig) (EventMeshProducer, error) 
{
        pm, err := wrapper.NewProducer()
        if err != nil {
                return nil, err
@@ -48,7 +58,7 @@ func NewEventMeshProducer(cfg *ProducerGroupConfig) 
(*EventMeshProducer, error)
                return nil, err
        }
 
-       p := &EventMeshProducer{
+       p := &eventMeshProducer{
                cfg:          cfg,
                producer:     pm,
                ServiceState: consts.INITED,
@@ -56,19 +66,19 @@ func NewEventMeshProducer(cfg *ProducerGroupConfig) 
(*EventMeshProducer, error)
        return p, nil
 }
 
-func (e *EventMeshProducer) Send(sctx SendMessageContext, callback 
*connector.SendCallback) error {
+func (e *eventMeshProducer) Send(sctx SendMessageContext, callback 
*connector.SendCallback) error {
        return e.producer.Send(sctx.Ctx, sctx.Event, callback)
 }
 
-func (e *EventMeshProducer) Request(sctx SendMessageContext, callback 
*connector.RequestReplyCallback, timeout time.Duration) error {
+func (e *eventMeshProducer) Request(sctx SendMessageContext, callback 
*connector.RequestReplyCallback, timeout time.Duration) error {
        return e.producer.Request(sctx.Ctx, sctx.Event, callback, timeout)
 }
 
-func (e *EventMeshProducer) Reply(sctx SendMessageContext, callback 
*connector.SendCallback) error {
+func (e *eventMeshProducer) Reply(sctx SendMessageContext, callback 
*connector.SendCallback) error {
        return e.producer.Reply(sctx.Ctx, sctx.Event, callback)
 }
 
-func (e *EventMeshProducer) Start() error {
+func (e *eventMeshProducer) Start() error {
        if e.ServiceState == "" || e.ServiceState == consts.RUNNING {
                return nil
        }
@@ -80,7 +90,7 @@ func (e *EventMeshProducer) Start() error {
        return nil
 }
 
-func (e *EventMeshProducer) Shutdown() error {
+func (e *eventMeshProducer) Shutdown() error {
        if e.ServiceState == "" || e.ServiceState == consts.INITED {
                return nil
        }
@@ -91,10 +101,10 @@ func (e *EventMeshProducer) Shutdown() error {
        return nil
 }
 
-func (e *EventMeshProducer) Status() consts.ServiceState {
+func (e *eventMeshProducer) Status() consts.ServiceState {
        return e.ServiceState
 }
 
-func (e *EventMeshProducer) String() string {
+func (e *eventMeshProducer) String() string {
        return fmt.Sprintf("eventMeshProducer, status:%s,  groupName:%s", 
e.ServiceState, e.cfg.GroupName)
 }
diff --git a/eventmesh-server-go/runtime/core/protocol/grpc/processor.go 
b/eventmesh-server-go/runtime/core/protocol/grpc/producer/producer_processor.go
similarity index 55%
rename from eventmesh-server-go/runtime/core/protocol/grpc/processor.go
rename to 
eventmesh-server-go/runtime/core/protocol/grpc/producer/producer_processor.go
index 7560e2cee..90f17b03b 100644
--- a/eventmesh-server-go/runtime/core/protocol/grpc/processor.go
+++ 
b/eventmesh-server-go/runtime/core/protocol/grpc/producer/producer_processor.go
@@ -13,11 +13,12 @@
 // See the License for the specific language governing permissions and
 // limitations under the License.
 
-package grpc
+package producer
 
 import (
        "context"
        "fmt"
+
        ce "github.com/cloudevents/sdk-go/v2"
        jsoniter "github.com/json-iterator/go"
        "sync"
@@ -29,6 +30,8 @@ import (
        
"github.com/apache/incubator-eventmesh/eventmesh-server-go/plugin/connector"
        
"github.com/apache/incubator-eventmesh/eventmesh-server-go/plugin/protocol"
        
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/consts"
+       
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/core/protocol/grpc/emitter"
+       
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/core/protocol/grpc/validator"
        
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/proto/pb"
 )
 
@@ -42,182 +45,27 @@ var (
        }}
 )
 
-// SubscribeProcessor Subscribe process subscribe message
-func ProcessSubscribe(gctx *GRPCContext, msg *pb.Subscription) (*pb.Response, 
error) {
-       hdr := msg.Header
-       if err := ValidateHeader(hdr); err != nil {
-               log.Warnf("invalid header:%v", err)
-               return buildPBResponse(grpc.EVENTMESH_PROTOCOL_HEADER_ERR), err
-       }
-       if err := ValidateSubscription(WEBHOOK, msg); err != nil {
-               log.Warnf("invalid body:%v", err)
-               return buildPBResponse(grpc.EVENTMESH_PROTOCOL_BODY_ERR), err
-       }
-       cmgr := gctx.ConsumerMgr
-       consumerGroup := msg.ConsumerGroup
-       url := msg.Url
-       items := msg.SubscriptionItems
-       var newClients []*GroupClient
-       for _, item := range items {
-               newClients = append(newClients, &GroupClient{
-                       ENV:              hdr.Env,
-                       IDC:              hdr.Idc,
-                       SYS:              hdr.Sys,
-                       IP:               hdr.Ip,
-                       PID:              hdr.Pid,
-                       ConsumerGroup:    consumerGroup,
-                       Topic:            item.Topic,
-                       SubscriptionMode: item.Mode,
-                       GRPCType:         WEBHOOK,
-                       URL:              url,
-                       LastUPTime:       time.Now(),
-               })
-       }
-       for _, cli := range newClients {
-               if err := cmgr.RegisterClient(cli); err != nil {
-                       return 
buildPBResponse(grpc.EVENTMESH_Subscribe_Register_ERR), err
-               }
-       }
-       meshConsumer, err := cmgr.GetConsumer(consumerGroup)
-       if err != nil {
-               return buildPBResponse(grpc.EVENTMESH_Consumer_NotFound_ERR), 
err
-       }
-       requireRestart := false
-       for _, cli := range newClients {
-               if meshConsumer.RegisterClient(cli) {
-                       requireRestart = true
-               }
-       }
-       if requireRestart {
-               log.Infof("ConsumerGroup %v topic info changed, restart 
EventMesh Consumer", consumerGroup)
-               if err := cmgr.restartConsumer(consumerGroup); err != nil {
-                       return 
buildPBResponse(grpc.EVENTMESH_Consumer_NotFound_ERR), err
-               }
-       } else {
-               log.Warnf("EventMesh consumer [%v] didn't restart.", 
consumerGroup)
-       }
-       return buildPBResponse(grpc.SUCCESS), nil
+type Processor interface {
+       AsyncMessage(ctx context.Context, producerMgr ProducerManager, msg 
*pb.SimpleMessage) (*pb.Response, error)
+       ReplyMessage(ctx context.Context, producerMgr ProducerManager, emiter 
emitter.EventEmitter, msg *pb.SimpleMessage) error
+       RequestReplyMessage(ctx context.Context, producerMgr ProducerManager, 
msg *pb.SimpleMessage) (*pb.SimpleMessage, error)
+       BatchPublish(ctx context.Context, producerMgr ProducerManager, msg 
*pb.BatchMessage) (*pb.Response, error)
 }
 
-func ProcessUnSubscribe(gctx *GRPCContext, msg *pb.Subscription) 
(*pb.Response, error) {
-       hdr := msg.Header
-       if err := ValidateHeader(hdr); err != nil {
-               log.Warnf("invalid header:%v", err)
-               return buildPBResponse(grpc.EVENTMESH_PROTOCOL_HEADER_ERR), err
-       }
-       if err := ValidateSubscription(WEBHOOK, msg); err != nil {
-               log.Warnf("invalid body:%v", err)
-               return buildPBResponse(grpc.EVENTMESH_PROTOCOL_BODY_ERR), err
-       }
-       cmgr := gctx.ConsumerMgr
-       consumerGroup := msg.ConsumerGroup
-       url := msg.Url
-       items := msg.SubscriptionItems
-       var removeClients []*GroupClient
-       for _, item := range items {
-               removeClients = append(removeClients, &GroupClient{
-                       ENV:              hdr.Env,
-                       IDC:              hdr.Idc,
-                       SYS:              hdr.Sys,
-                       IP:               hdr.Ip,
-                       PID:              hdr.Pid,
-                       ConsumerGroup:    consumerGroup,
-                       Topic:            item.Topic,
-                       SubscriptionMode: item.Mode,
-                       GRPCType:         WEBHOOK,
-                       URL:              url,
-                       LastUPTime:       time.Now(),
-               })
-       }
-       for _, cli := range removeClients {
-               if err := cmgr.DeRegisterClient(cli); err != nil {
-                       return 
buildPBResponse(grpc.EVENTMESH_Subscribe_Register_ERR), err
-               }
-       }
-       meshConsumer, err := cmgr.GetConsumer(consumerGroup)
-       if err != nil {
-               return buildPBResponse(grpc.EVENTMESH_Consumer_NotFound_ERR), 
err
-       }
-       requireRestart := false
-       for _, cli := range removeClients {
-               if meshConsumer.DeRegisterClient(cli) {
-                       requireRestart = true
-               }
-       }
-       if requireRestart {
-               log.Infof("ConsumerGroup %v topic info changed, restart 
EventMesh Consumer", consumerGroup)
-               if err := cmgr.restartConsumer(consumerGroup); err != nil {
-                       return 
buildPBResponse(grpc.EVENTMESH_Consumer_NotFound_ERR), err
-               }
-       } else {
-               log.Warnf("EventMesh consumer [%v] didn't restart.", 
consumerGroup)
-       }
-       return buildPBResponse(grpc.SUCCESS), nil
+type processor struct {
 }
 
-func ProcessSubscribeStream(ctx context.Context, gctx *GRPCContext, emiter 
*EventEmitter, msg *pb.Subscription) error {
-       hdr := msg.Header
-       if err := ValidateHeader(hdr); err != nil {
-               log.Warnf("invalid header:%v", err)
-               emiter.sendStreamResp(hdr, grpc.EVENTMESH_PROTOCOL_HEADER_ERR)
-               return err
-       }
-       if err := ValidateSubscription(STREAM, msg); err != nil {
-               log.Warnf("invalid body:%v", err)
-               emiter.sendStreamResp(hdr, grpc.EVENTMESH_PROTOCOL_BODY_ERR)
-               return err
-       }
-       cmgr := gctx.ConsumerMgr
-       consumerGroup := msg.ConsumerGroup
-       var clients []*GroupClient
-       for _, item := range msg.SubscriptionItems {
-               clients = append(clients, &GroupClient{
-                       ENV:              hdr.Env,
-                       IDC:              hdr.Idc,
-                       SYS:              hdr.Sys,
-                       IP:               hdr.Ip,
-                       PID:              hdr.Pid,
-                       ConsumerGroup:    consumerGroup,
-                       Topic:            item.Topic,
-                       SubscriptionMode: item.Mode,
-                       GRPCType:         STREAM,
-                       LastUPTime:       time.Now(),
-                       Emiter:           emiter,
-               })
-       }
-       for _, cli := range clients {
-               if err := cmgr.RegisterClient(cli); err != nil {
-                       return err
-               }
-       }
-       meshConsumer, err := cmgr.GetConsumer(consumerGroup)
-       if err != nil {
-               return err
-       }
-       requireRestart := false
-       for _, cli := range clients {
-               if meshConsumer.RegisterClient(cli) {
-                       requireRestart = true
-               }
-       }
-       if requireRestart {
-               log.Infof("ConsumerGroup %v topic info changed, restart 
EventMesh Consumer", consumerGroup)
-               return cmgr.restartConsumer(consumerGroup)
-       } else {
-               log.Warnf("EventMesh consumer [%v] didn't restart.", 
consumerGroup)
-       }
-
-       return nil
+func NewProcessor() Processor {
+       return &processor{}
 }
 
-// ProcessAsyncMessage process async message
-func ProcessAsyncMessage(ctx context.Context, gctx *GRPCContext, msg 
*pb.SimpleMessage) (*pb.Response, error) {
+func (p *processor) AsyncMessage(ctx context.Context, producerMgr 
ProducerManager, msg *pb.SimpleMessage) (*pb.Response, error) {
        hdr := msg.Header
-       if err := ValidateHeader(hdr); err != nil {
+       if err := validator.ValidateHeader(hdr); err != nil {
                log.Warnf("invalid header:%v", err)
                return buildPBResponse(grpc.EVENTMESH_PROTOCOL_HEADER_ERR), err
        }
-       if err := ValidateMessage(msg); err != nil {
+       if err := validator.ValidateMessage(msg); err != nil {
                log.Warnf("invalid body:%v", err)
                return buildPBResponse(grpc.EVENTMESH_PROTOCOL_BODY_ERR), err
        }
@@ -238,7 +86,7 @@ func ProcessAsyncMessage(ctx context.Context, gctx 
*GRPCContext, msg *pb.SimpleM
        if err != nil {
                return buildPBResponse(grpc.EVENTMESH_Transfer_Protocol_ERR), 
err
        }
-       ep, err := gctx.ProducerMgr.GetProducer(pg)
+       ep, err := producerMgr.GetProducer(pg)
        if err != nil {
                return buildPBResponse(grpc.EVENTMESH_PROTOCOL_BODY_ERR), err
        }
@@ -269,16 +117,16 @@ func ProcessAsyncMessage(ctx context.Context, gctx 
*GRPCContext, msg *pb.SimpleM
        return buildPBResponse(code), nil
 }
 
-func ProcessReplyMessage(ctx context.Context, gctx *GRPCContext, emiter 
*EventEmitter, msg *pb.SimpleMessage) error {
+func (p *processor) ReplyMessage(ctx context.Context, producerMgr 
ProducerManager, emiter emitter.EventEmitter, msg *pb.SimpleMessage) error {
        hdr := msg.Header
-       if err := ValidateHeader(hdr); err != nil {
+       if err := validator.ValidateHeader(hdr); err != nil {
                log.Warnf("invalid header:%v", err)
-               emiter.sendStreamResp(hdr, grpc.EVENTMESH_PROTOCOL_HEADER_ERR)
+               emiter.SendStreamResp(hdr, grpc.EVENTMESH_PROTOCOL_HEADER_ERR)
                return err
        }
-       if err := ValidateMessage(msg); err != nil {
+       if err := validator.ValidateMessage(msg); err != nil {
                log.Warnf("invalid body:%v", err)
-               emiter.sendStreamResp(hdr, grpc.EVENTMESH_PROTOCOL_BODY_ERR)
+               emiter.SendStreamResp(hdr, grpc.EVENTMESH_PROTOCOL_BODY_ERR)
                return err
        }
        seqNum := msg.SeqNum
@@ -291,19 +139,19 @@ func ProcessReplyMessage(ctx context.Context, gctx 
*GRPCContext, emiter *EventEm
        adp := plugin.Get(plugin.Protocol, protocolType).(protocol.Adapter)
        if adp == nil {
                log.Warnf("protocol plugin not found:%v", protocolType)
-               emiter.sendStreamResp(hdr, grpc.EVENTMESH_Plugin_NotFound_ERR)
+               emiter.SendStreamResp(hdr, grpc.EVENTMESH_Plugin_NotFound_ERR)
                return ErrProtocolPluginNotFound
        }
        cevt, err := adp.ToCloudEvent(&grpc.SimpleMessageWrapper{SimpleMessage: 
msg})
        if err != nil {
                log.Warnf("transfer to cloud event msg err:%v", err)
-               emiter.sendStreamResp(hdr, grpc.EVENTMESH_Transfer_Protocol_ERR)
+               emiter.SendStreamResp(hdr, grpc.EVENTMESH_Transfer_Protocol_ERR)
                return err
        }
-       emProducer, err := gctx.ProducerMgr.GetProducer(producerGroup)
+       emProducer, err := producerMgr.GetProducer(producerGroup)
        if err != nil {
                log.Warnf("no eventmesh producer found, err:%v, group:%v", err, 
producerGroup)
-               emiter.sendStreamResp(hdr, 
grpc.EVENTMESH_Producer_Group_NotFound_ERR)
+               emiter.SendStreamResp(hdr, 
grpc.EVENTMESH_Producer_Group_NotFound_ERR)
                return err
        }
        start := time.Now()
@@ -321,7 +169,7 @@ func ProcessReplyMessage(ctx context.Context, gctx 
*GRPCContext, emiter *EventEm
                                        time.Now().Sub(start).Milliseconds(), 
replyTopic, seqNum, uniqID)
                        },
                        OnError: func(result *connector.ErrorResult) {
-                               emiter.sendStreamResp(hdr, 
grpc.EVENTMESH_REPLY_MSG_ERR)
+                               emiter.SendStreamResp(hdr, 
grpc.EVENTMESH_REPLY_MSG_ERR)
                                
log.Errorf("message|mq2eventmesh|REPLY|ReplyToServer|send2MQCost=%vms|topic=%v|bizSeqNo=%v|uniqueId=%v",
                                        time.Now().Sub(start).Milliseconds(), 
replyTopic, seqNum, uniqID, result.Err)
                        },
@@ -329,49 +177,17 @@ func ProcessReplyMessage(ctx context.Context, gctx 
*GRPCContext, emiter *EventEm
        )
 }
 
-func ProcessHeartbeat(gctx *GRPCContext, msg *pb.Heartbeat) (*pb.Response, 
error) {
-       hdr := msg.Header
-       if err := ValidateHeader(hdr); err != nil {
-               log.Warnf("invalid header:%v", err)
-               return buildPBResponse(grpc.EVENTMESH_PROTOCOL_HEADER_ERR), err
-       }
-       if err := ValidateHeartBeat(msg); err != nil {
-               log.Warnf("invalid body:%v", err)
-               return buildPBResponse(grpc.EVENTMESH_PROTOCOL_BODY_ERR), err
-       }
-       if msg.ClientType != pb.Heartbeat_SUB {
-               log.Warnf("client type err, not sub")
-               return buildPBResponse(grpc.EVENTMESH_Heartbeat_Protocol_ERR), 
fmt.Errorf("protocol not sub")
-       }
-       cmgr := gctx.ConsumerMgr
-       consumerGroup := msg.ConsumerGroup
-       for _, item := range msg.HeartbeatItems {
-               cli := &GroupClient{
-                       ENV:           hdr.Env,
-                       IDC:           hdr.Idc,
-                       SYS:           hdr.Sys,
-                       IP:            hdr.Ip,
-                       PID:           hdr.Pid,
-                       ConsumerGroup: consumerGroup,
-                       Topic:         item.Topic,
-                       LastUPTime:    time.Now(),
-               }
-               cmgr.UpdateClientTime(cli)
-       }
-       return buildPBResponse(grpc.SUCCESS), nil
-}
-
-func ProcessRequestReplyMessage(ctx context.Context, gctx *GRPCContext, msg 
*pb.SimpleMessage) (*pb.SimpleMessage, error) {
+func (p *processor) RequestReplyMessage(ctx context.Context, producerMgr 
ProducerManager, msg *pb.SimpleMessage) (*pb.SimpleMessage, error) {
        var (
                err  error
                resp *pb.SimpleMessage
                hdr  = msg.Header
        )
-       if err = ValidateHeader(hdr); err != nil {
+       if err = validator.ValidateHeader(hdr); err != nil {
                log.Warnf("invalid header:%v", err)
                return buildPBSimpleMessage(hdr, 
grpc.EVENTMESH_PROTOCOL_HEADER_ERR), err
        }
-       if err = ValidateMessage(msg); err != nil {
+       if err = validator.ValidateMessage(msg); err != nil {
                log.Warnf("invalid body:%v", err)
                return buildPBSimpleMessage(hdr, 
grpc.EVENTMESH_PROTOCOL_BODY_ERR), err
        }
@@ -392,7 +208,7 @@ func ProcessRequestReplyMessage(ctx context.Context, gctx 
*GRPCContext, msg *pb.
        producerGroup := msg.ProducerGroup
        ttl, _ := StringToDuration(msg.Ttl)
        start := time.Now()
-       ep, err := gctx.ProducerMgr.GetProducer(producerGroup)
+       ep, err := producerMgr.GetProducer(producerGroup)
        if err != nil {
                return buildPBSimpleMessage(hdr, 
grpc.EVENTMESH_PROTOCOL_BODY_ERR), err
        }
@@ -431,16 +247,16 @@ func ProcessRequestReplyMessage(ctx context.Context, gctx 
*GRPCContext, msg *pb.
        return resp, err
 }
 
-func ProcessBatchPublish(ctx context.Context, gctx *GRPCContext, msg 
*pb.BatchMessage) (*pb.Response, error) {
+func (p *processor) BatchPublish(ctx context.Context, producerMgr 
ProducerManager, msg *pb.BatchMessage) (*pb.Response, error) {
        var (
                err error
                hdr = msg.Header
        )
-       if err = ValidateHeader(hdr); err != nil {
+       if err = validator.ValidateHeader(hdr); err != nil {
                log.Warnf("invalid header:%v", err)
                return buildPBResponse(grpc.EVENTMESH_PROTOCOL_HEADER_ERR), err
        }
-       if err = ValidateBatchMessage(msg); err != nil {
+       if err = validator.ValidateBatchMessage(msg); err != nil {
                log.Warnf("invalid body:%v", err)
                return buildPBResponse(grpc.EVENTMESH_PROTOCOL_BODY_ERR), err
        }
@@ -457,7 +273,7 @@ func ProcessBatchPublish(ctx context.Context, gctx 
*GRPCContext, msg *pb.BatchMe
        }
        topic := msg.Topic
        producerGroup := msg.ProducerGroup
-       ep, err := gctx.ProducerMgr.GetProducer(producerGroup)
+       ep, err := producerMgr.GetProducer(producerGroup)
        if err != nil {
                return buildPBResponse(grpc.EVENTMESH_PROTOCOL_BODY_ERR), err
        }
diff --git a/eventmesh-server-go/runtime/core/protocol/grpc/producer_service.go 
b/eventmesh-server-go/runtime/core/protocol/grpc/producer/producer_service.go
similarity index 88%
rename from eventmesh-server-go/runtime/core/protocol/grpc/producer_service.go
rename to 
eventmesh-server-go/runtime/core/protocol/grpc/producer/producer_service.go
index 191f3df81..46b6561ed 100644
--- a/eventmesh-server-go/runtime/core/protocol/grpc/producer_service.go
+++ 
b/eventmesh-server-go/runtime/core/protocol/grpc/producer/producer_service.go
@@ -13,7 +13,7 @@
 // See the License for the specific language governing permissions and
 // limitations under the License.
 
-package grpc
+package producer
 
 import (
        "context"
@@ -28,19 +28,19 @@ var defaultAsyncTimeout = time.Second * 5
 
 type ProducerService struct {
        pb.UnimplementedPublisherServiceServer
-       gctx     *GRPCContext
-       sendPool *ants.Pool
+       producerMgr ProducerManager
+       sendPool    *ants.Pool
 }
 
-func NewProducerServiceServer(gctx *GRPCContext) (*ProducerService, error) {
+func NewProducerServiceServer(producerMgr ProducerManager) (*ProducerService, 
error) {
        ps := config.GlobalConfig().Server.GRPCOption.SendPoolSize
        pl, err := ants.NewPool(ps)
        if err != nil {
                return nil, err
        }
        return &ProducerService{
-               gctx:     gctx,
-               sendPool: pl,
+               producerMgr: producerMgr,
+               sendPool:    pl,
        }, nil
 }
 
@@ -55,7 +55,7 @@ func (p *ProducerService) Publish(ctx context.Context, msg 
*pb.SimpleMessage) (*
                err     error
        )
        p.sendPool.Submit(func() {
-               resp, err = ProcessAsyncMessage(ctx, p.gctx, msg)
+               resp, err = NewProcessor().AsyncMessage(ctx, p.producerMgr, msg)
                errChan <- err
        })
        select {
@@ -80,7 +80,7 @@ func (p *ProducerService) RequestReply(ctx context.Context, 
msg *pb.SimpleMessag
                err     error
        )
        p.sendPool.Submit(func() {
-               resp, err = ProcessRequestReplyMessage(ctx, p.gctx, msg)
+               resp, err = NewProcessor().RequestReplyMessage(ctx, 
p.producerMgr, msg)
                errChan <- err
        })
        select {
@@ -106,7 +106,7 @@ func (p *ProducerService) BatchPublish(ctx context.Context, 
msg *pb.BatchMessage
                err     error
        )
        p.sendPool.Submit(func() {
-               resp, err = ProcessBatchPublish(ctx, p.gctx, msg)
+               resp, err = NewProcessor().BatchPublish(ctx, p.producerMgr, msg)
                errChan <- err
        })
        select {
diff --git 
a/eventmesh-server-go/runtime/core/protocol/grpc/producer_service_test.go 
b/eventmesh-server-go/runtime/core/protocol/grpc/producer/producer_service_test.go
similarity index 99%
rename from 
eventmesh-server-go/runtime/core/protocol/grpc/producer_service_test.go
rename to 
eventmesh-server-go/runtime/core/protocol/grpc/producer/producer_service_test.go
index a07fad86f..0c800fcc2 100644
--- a/eventmesh-server-go/runtime/core/protocol/grpc/producer_service_test.go
+++ 
b/eventmesh-server-go/runtime/core/protocol/grpc/producer/producer_service_test.go
@@ -13,7 +13,7 @@
 // See the License for the specific language governing permissions and
 // limitations under the License.
 
-package grpc
+package producer
 
 import (
        "context"
diff --git a/eventmesh-server-go/runtime/core/protocol/grpc/retry.go 
b/eventmesh-server-go/runtime/core/protocol/grpc/retry/retry.go
similarity index 87%
rename from eventmesh-server-go/runtime/core/protocol/grpc/retry.go
rename to eventmesh-server-go/runtime/core/protocol/grpc/retry/retry.go
index a42bfa3f7..7f8ae47d2 100644
--- a/eventmesh-server-go/runtime/core/protocol/grpc/retry.go
+++ b/eventmesh-server-go/runtime/core/protocol/grpc/retry/retry.go
@@ -13,21 +13,21 @@
 // See the License for the specific language governing permissions and
 // limitations under the License.
 
-package grpc
+package retry
 
 import "time"
 
-type Context struct {
+type Retry struct {
        RetryTimes  int
        ExecuteTime time.Time
        Do          func() error
 }
 
-func (c *Context) SetDelay(delay time.Duration) *Context {
+func (c *Retry) SetDelay(delay time.Duration) *Retry {
        c.ExecuteTime = time.Now().Add(delay)
        return c
 }
 
-func (c *Context) GetDelay() time.Duration {
+func (c *Retry) GetDelay() time.Duration {
        return c.ExecuteTime.Sub(time.Now())
 }
diff --git a/eventmesh-server-go/runtime/core/protocol/grpc/validator.go 
b/eventmesh-server-go/runtime/core/protocol/grpc/validator/validator.go
similarity index 95%
rename from eventmesh-server-go/runtime/core/protocol/grpc/validator.go
rename to eventmesh-server-go/runtime/core/protocol/grpc/validator/validator.go
index aa6367b5c..3f2e9ffa7 100644
--- a/eventmesh-server-go/runtime/core/protocol/grpc/validator.go
+++ b/eventmesh-server-go/runtime/core/protocol/grpc/validator/validator.go
@@ -13,10 +13,11 @@
 // See the License for the specific language governing permissions and
 // limitations under the License.
 
-package grpc
+package validator
 
 import (
        "fmt"
+       
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/consts"
        
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/proto/pb"
 )
 
@@ -101,8 +102,8 @@ func ValidateMessage(msg *pb.SimpleMessage) error {
        return nil
 }
 
-func ValidateSubscription(stype GRPCType, msg *pb.Subscription) error {
-       if stype == WEBHOOK && msg.Url == "" {
+func ValidateSubscription(stype consts.GRPCType, msg *pb.Subscription) error {
+       if stype == consts.WEBHOOK && msg.Url == "" {
                return ErrSubscriptionNoURL
        }
 
diff --git a/eventmesh-server-go/runtime/emserver/grpc.go 
b/eventmesh-server-go/runtime/emserver/grpc.go
index 1c824435e..30858a725 100644
--- a/eventmesh-server-go/runtime/emserver/grpc.go
+++ b/eventmesh-server-go/runtime/emserver/grpc.go
@@ -18,7 +18,9 @@ package emserver
 import (
        "fmt"
        "github.com/apache/incubator-eventmesh/eventmesh-server-go/config"
-       grpc2 
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/core/protocol/grpc"
+       
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/core/protocol/grpc/consumer"
+       
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/core/protocol/grpc/heartbeat"
+       
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/core/protocol/grpc/producer"
        
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/proto/pb"
        "google.golang.org/grpc"
        "google.golang.org/grpc/credentials"
@@ -45,20 +47,30 @@ func NewGRPCServer(opt *config.GRPCOption) (GracefulServer, 
error) {
                err error
        )
 
-       grpcCtx, err := grpc2.New()
+       //msgReqPerSeconds := 
config.GlobalConfig().Server.GRPCOption.MsgReqNumPerSecond
+       //limiter := rate.NewLimiter(rate.Limit(msgReqPerSeconds), 10)
+
+       consumerMgr, err := consumer.NewConsumerManager()
        if err != nil {
                return nil, err
        }
+       producerMgr, err := producer.NewProducerManager()
+       if err != nil {
+               return nil, err
+       }
+
+       //registryName := config.GlobalConfig().Server.GRPCOption.RegistryName
+       //regis := registry.Get(registryName)
 
-       consumerSVC, err := grpc2.NewConsumerServiceServer(grpcCtx)
+       consumerSVC, err := consumer.NewConsumerServiceServer(consumerMgr)
        if err != nil {
                return nil, err
        }
-       producerSVC, err := grpc2.NewProducerServiceServer(grpcCtx)
+       producerSVC, err := producer.NewProducerServiceServer(producerMgr)
        if err != nil {
                return nil, err
        }
-       hbSVC, err := grpc2.NewHeartbeatServiceServer(grpcCtx)
+       hbSVC, err := heartbeat.NewHeartbeatServiceServer(consumerMgr)
        if err != nil {
                return nil, err
        }
diff --git a/eventmesh-server-go/runtime/emserver/grpc_test.go 
b/eventmesh-server-go/runtime/emserver/grpc_test.go
index 7eb6e599e..329b93fe2 100644
--- a/eventmesh-server-go/runtime/emserver/grpc_test.go
+++ b/eventmesh-server-go/runtime/emserver/grpc_test.go
@@ -16,53 +16,24 @@
 package emserver
 
 import (
-       "github.com/apache/incubator-eventmesh/eventmesh-server-go/config"
+       
"github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/emserver/mocks"
+       "github.com/golang/mock/gomock"
        "github.com/stretchr/testify/assert"
-       "gopkg.in/yaml.v3"
        "testing"
-       "time"
 )
 
-func Test_NewGRPCServer(t *testing.T) {
-       staticCfg := `server:
-  grpc:
-    port: 10010
-    tls:
-      enable-secure: false
-      ca: ""
-      certfile: ""
-      keyfile: ""
-    pprof:
-      port: 10011
-    send-pool-size: 10
-    subscribe-pool-size: 10
-    retry-pool-size: 10
-    push-message-pool-size: 10
-    reply-pool-size: 10
-    msg-req-num-per-second: 5
-    cluster: "test"
-    idc: "idc1"
-    session-expired-in-mills: 5s
-    send-message-timeout: 5s`
-       cfg := &config.Config{}
-       err := yaml.Unmarshal([]byte(staticCfg), cfg)
+func Test_Serve(t *testing.T) {
+       mockCtl := gomock.NewController(t)
+       mockSvr := mocks.NewMockGracefulServer(mockCtl)
+       mockSvr.EXPECT().Serve().Return(nil).Times(1)
+       err := mockSvr.Serve()
        assert.NoError(t, err)
-       config.SetGlobalConfig(cfg)
-
-       t.Run("create plain server", func(t *testing.T) {
-               svr, err := NewGRPCServer(cfg.Server.GRPCOption)
-               assert.NoError(t, err)
-               assert.NotNil(t, svr)
-       })
+}
 
-       t.Run("boot grpc srever", func(t *testing.T) {
-               svr, err := NewGRPCServer(cfg.Server.GRPCOption)
-               assert.NoError(t, err)
-               assert.NotNil(t, svr)
-               go func() {
-                       assert.NoError(t, svr.Serve())
-               }()
-               time.Sleep(3 * time.Second)
-               svr.Stop()
-       })
+func Test_Stop(t *testing.T) {
+       mockCtl := gomock.NewController(t)
+       mockSvr := mocks.NewMockGracefulServer(mockCtl)
+       mockSvr.EXPECT().Stop().Return(nil).Times(1)
+       err := mockSvr.Stop()
+       assert.NoError(t, err)
 }
diff --git a/eventmesh-server-go/runtime/emserver/mocks/mock_graceful.go 
b/eventmesh-server-go/runtime/emserver/mocks/mock_graceful.go
new file mode 100644
index 000000000..0d8bd216f
--- /dev/null
+++ b/eventmesh-server-go/runtime/emserver/mocks/mock_graceful.go
@@ -0,0 +1,62 @@
+// Code generated by MockGen. DO NOT EDIT.
+// Source: 
github.com/apache/incubator-eventmesh/eventmesh-server-go/runtime/emserver 
(interfaces: GracefulServer)
+
+// Package mocks is a generated GoMock package.
+package mocks
+
+import (
+       reflect "reflect"
+
+       gomock "github.com/golang/mock/gomock"
+)
+
+// MockGracefulServer is a mock of GracefulServer interface.
+type MockGracefulServer struct {
+       ctrl     *gomock.Controller
+       recorder *MockGracefulServerMockRecorder
+}
+
+// MockGracefulServerMockRecorder is the mock recorder for MockGracefulServer.
+type MockGracefulServerMockRecorder struct {
+       mock *MockGracefulServer
+}
+
+// NewMockGracefulServer creates a new mock instance.
+func NewMockGracefulServer(ctrl *gomock.Controller) *MockGracefulServer {
+       mock := &MockGracefulServer{ctrl: ctrl}
+       mock.recorder = &MockGracefulServerMockRecorder{mock}
+       return mock
+}
+
+// EXPECT returns an object that allows the caller to indicate expected use.
+func (m *MockGracefulServer) EXPECT() *MockGracefulServerMockRecorder {
+       return m.recorder
+}
+
+// Serve mocks base method.
+func (m *MockGracefulServer) Serve() error {
+       m.ctrl.T.Helper()
+       ret := m.ctrl.Call(m, "Serve")
+       ret0, _ := ret[0].(error)
+       return ret0
+}
+
+// Serve indicates an expected call of Serve.
+func (mr *MockGracefulServerMockRecorder) Serve() *gomock.Call {
+       mr.mock.ctrl.T.Helper()
+       return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "Serve", 
reflect.TypeOf((*MockGracefulServer)(nil).Serve))
+}
+
+// Stop mocks base method.
+func (m *MockGracefulServer) Stop() error {
+       m.ctrl.T.Helper()
+       ret := m.ctrl.Call(m, "Stop")
+       ret0, _ := ret[0].(error)
+       return ret0
+}
+
+// Stop indicates an expected call of Stop.
+func (mr *MockGracefulServerMockRecorder) Stop() *gomock.Call {
+       mr.mock.ctrl.T.Helper()
+       return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "Stop", 
reflect.TypeOf((*MockGracefulServer)(nil).Stop))
+}


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


Reply via email to