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]