You are viewing a plain text version of this content. The canonical link for it is here.
Posted to commits@eventmesh.apache.org by wa...@apache.org on 2022/08/22 11:03:30 UTC
[incubator-eventmesh] 03/10: add standalone
This is an automated email from the ASF dual-hosted git repository.
walleliu pushed a commit to branch eventmesh-server-go
in repository https://gitbox.apache.org/repos/asf/incubator-eventmesh.git
commit 3f0d41a362554d0e940cb93e6fa85d4430d1859f
Author: walleliu <li...@163.com>
AuthorDate: Thu Aug 18 19:32:38 2022 +0800
add standalone
---
eventmesh-server-go/go.mod | 6 +-
eventmesh-server-go/go.sum | 16 +-
eventmesh-server-go/pkg/connector/action.go | 9 +
.../{go.mod => pkg/connector/consumer.go} | 40 +-
.../{go.mod => pkg/connector/lifecycle.go} | 31 +-
.../{go.mod => pkg/connector/listener.go} | 28 +-
eventmesh-server-go/pkg/connector/message_queue.go | 165 +++
.../{go.mod => pkg/connector/properties.go} | 29 +-
eventmesh-server-go/pkg/connector/publisher.go | 62 +
.../pkg/connector/standalone/broker.go | 66 +
.../connector/standalone/message_entity.go} | 34 +-
eventmesh-server-go/pkg/runtime/emserver/tcp.go | 14 +-
eventmesh-server-go/pkg/runtime/proto/pb/README.md | 20 +
.../pkg/runtime/proto/pb/eventmesh-client.pb.go | 1485 ++++++++++++++++++++
.../pkg/runtime/proto/pb/eventmesh-client.proto | 165 +++
.../runtime/proto/pb/eventmesh-client_grpc.pb.go | 479 +++++++
16 files changed, 2518 insertions(+), 131 deletions(-)
diff --git a/eventmesh-server-go/go.mod b/eventmesh-server-go/go.mod
index 62d5e28a..878a9cdc 100644
--- a/eventmesh-server-go/go.mod
+++ b/eventmesh-server-go/go.mod
@@ -18,6 +18,7 @@ module github.com/apache/incubator-eventmesh/eventmesh-server-go
go 1.16
require (
+ github.com/cloudevents/sdk-go/v2 v2.11.0
github.com/gogf/gf v1.16.9
github.com/hashicorp/go-multierror v1.1.1
github.com/lestrrat-go/strftime v1.0.6
@@ -33,7 +34,10 @@ require (
github.com/BurntSushi/toml v1.2.0 // indirect
github.com/gin-contrib/pprof v1.4.0
github.com/gin-gonic/gin v1.8.1
- github.com/panjf2000/gnet/v2 v2.1.1
github.com/unrolled/secure v1.12.0
+ go.uber.org/atomic v1.9.0 // indirect
go.uber.org/fx v1.18.1
+ go.uber.org/multierr v1.7.0 // indirect
+ golang.org/x/sys v0.0.0-20220224120231-95c6836cb0e7 // indirect
+ google.golang.org/protobuf v1.28.0
)
diff --git a/eventmesh-server-go/go.sum b/eventmesh-server-go/go.sum
index ad21fd2f..a23464ee 100644
--- a/eventmesh-server-go/go.sum
+++ b/eventmesh-server-go/go.sum
@@ -25,6 +25,8 @@ github.com/cespare/xxhash/v2 v2.1.1/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XL
github.com/clbanning/mxj v1.8.5-0.20200714211355-ff02cfb8ea28 h1:LdXxtjzvZYhhUaonAaAKArG3pyC67kGL3YY+6hGG8G4=
github.com/clbanning/mxj v1.8.5-0.20200714211355-ff02cfb8ea28/go.mod h1:BVjHeAH+rl9rs6f+QIpeRl0tfu10SXn1pUSa5PVGJng=
github.com/client9/misspell v0.3.4/go.mod h1:qj6jICC3Q7zFZvVWo7KLAzC3yx5G7kyvSDkc90ppPyw=
+github.com/cloudevents/sdk-go/v2 v2.6.0 h1:yp6zLEvhXSi6P25zzfgORgFI0quG2/NXoH9QoHzvKn8=
+github.com/cloudevents/sdk-go/v2 v2.6.0/go.mod h1:nlXhgFkf0uTopxmRXalyMwS2LG70cRGPrxzmjJgSG0U=
github.com/cncf/udpa/go v0.0.0-20201120205902-5459f2c99403/go.mod h1:WmhPx2Nbnhtbo57+VJT5O0JRkEi1Wbu0z5j0R8u5Hbk=
github.com/cpuguy83/go-md2man/v2 v2.0.2/go.mod h1:tgQtvFlXSQOSOSIRvRPT7W67SCa46tRHOmNcaadrF8o=
github.com/creack/pty v1.1.9/go.mod h1:oKZEueFk5CKHvIhNR5MUki03XCEU+Q6VDXinZuGJ33E=
@@ -100,6 +102,8 @@ github.com/google/go-cmp v0.5.6 h1:BKbKCqvP6I+rmFHt06ZmyQtvB8xAkWdhFyr0ZUNZcxQ=
github.com/google/go-cmp v0.5.6/go.mod h1:v8dTdLbMG2kIc/vJvl+f65V22dbkXbowE6jgT/gNBxE=
github.com/google/gofuzz v1.0.0/go.mod h1:dBl0BpW6vV/+mYPU4Po3pmUjxk6FQPldtuIdl/M65Eg=
github.com/google/renameio v0.1.0/go.mod h1:KWCgfxg9yswjAJkECMjeO8J8rahYeXnNhOm40UhjYkI=
+github.com/google/uuid v1.1.1/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo=
+github.com/google/uuid v1.1.2 h1:EVhdT+1Kseyi1/pUmXKaFxYsDNy9RQYkMWRH68J/W7Y=
github.com/google/uuid v1.1.2/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo=
github.com/gopherjs/gopherjs v0.0.0-20181017120253-0766667cb4d1 h1:EGx4pi6eqNxGaHF6qqu48+N2wcFQ5qg5FXgOdqsJ5d8=
github.com/gopherjs/gopherjs v0.0.0-20181017120253-0766667cb4d1/go.mod h1:wJfORRmW1u3UXTncJ5qlYoELFm8eSnnEO6hX4iZ3EWY=
@@ -164,12 +168,9 @@ github.com/mwitkow/go-conntrack v0.0.0-20161129095857-cc309e4a2223/go.mod h1:qRW
github.com/mwitkow/go-conntrack v0.0.0-20190716064945-2f068394615f/go.mod h1:qRWi+5nqEBWmkhHvq77mSJWrCKwh8bxhgT7d/eI7P4U=
github.com/nacos-group/nacos-sdk-go/v2 v2.0.3 h1:wNaSObv2wNFyMQird3vPc4XaaeNYPT6s8Oqu7GxCbRw=
github.com/nacos-group/nacos-sdk-go/v2 v2.0.3/go.mod h1:SlhyCAv961LcZ198XpKfPEQqlJWt2HkL1fDLas0uy/w=
+github.com/niemeyer/pretty v0.0.0-20200227124842-a10e7caefd8e/go.mod h1:zD1mROLANZcx1PVRCS0qkT7pwLkGfwJo4zjcN/Tysno=
github.com/olekukonko/tablewriter v0.0.5 h1:P2Ga83D34wi1o9J6Wh1mRuqd4mF/x/lgBS7N7AbDhec=
github.com/olekukonko/tablewriter v0.0.5/go.mod h1:hPp6KlRPjbx+hW8ykQs1w3UBbZlj6HuIJcUGPhkA7kY=
-github.com/panjf2000/ants/v2 v2.4.8 h1:JgTbolX6K6RreZ4+bfctI0Ifs+3mrE5BIHudQxUDQ9k=
-github.com/panjf2000/ants/v2 v2.4.8/go.mod h1:f6F0NZVFsGCp5A7QW/Zj/m92atWwOkY0OIhFxRNFr4A=
-github.com/panjf2000/gnet/v2 v2.1.1 h1:YK2Vs7nnAecnm/k7cccW2FkhlefgXUIIUww6Y6cChOc=
-github.com/panjf2000/gnet/v2 v2.1.1/go.mod h1:unWr2B4jF0DQPJH3GsXBGQiDcAamM6+Pf5FiK705kc4=
github.com/pelletier/go-toml/v2 v2.0.1 h1:8e3L2cCQzLFi2CR4g7vGFuFxX7Jl1kKX8gW+iV0GUKU=
github.com/pelletier/go-toml/v2 v2.0.1/go.mod h1:r9LEWfGN8R5k0VXJ+0BkIe7MYkRdwZOjgMj2KwnJFUo=
github.com/pkg/diff v0.0.0-20210226163009-20ebb0f2a09e/go.mod h1:pJLUxLENpZxwdsKMEsNbx1VGcRFpLqf3715MtcvvzbA=
@@ -237,6 +238,7 @@ go.opentelemetry.io/otel v1.0.0 h1:qTTn6x71GVBvoafHK/yaRUmFzI4LcONZD0/kXxl5PHI=
go.opentelemetry.io/otel v1.0.0/go.mod h1:AjRVh9A5/5DE7S+mZtTR6t8vpKKryam+0lREnfmS4cg=
go.opentelemetry.io/otel/trace v1.0.0 h1:TSBr8GTEtKevYMG/2d21M989r5WJYVimhTHBKVEZuh4=
go.opentelemetry.io/otel/trace v1.0.0/go.mod h1:PXTWqayeFUlJV1YDNhsJYB184+IvAH814St6o6ajzIs=
+go.uber.org/atomic v1.4.0/go.mod h1:gD2HeocX3+yG+ygLZcrzQJaqmWj9AIm7n08wl/qW/PE=
go.uber.org/atomic v1.6.0/go.mod h1:sABNBOSYdrvTF6hTgEIbc7YasKWGhgEQZyfxyTvoXHQ=
go.uber.org/atomic v1.7.0/go.mod h1:fEN4uk6kAWBTFdckzkM89CLk9XfWZrxpCo0nPH17wJc=
go.uber.org/atomic v1.9.0 h1:ECmE8Bn/WFTYwEW/bpKD3M8VtR/zQVbavAoalC1PYyE=
@@ -247,14 +249,15 @@ go.uber.org/fx v1.18.1 h1:I7VWkdv4iKcbpH7KVSi9Fe1LGmpJv+pbBIb9NidPb+E=
go.uber.org/fx v1.18.1/go.mod h1:g0V1KMQ66zIRk8bLu3Ea5Jt2w/cHlOIp4wdRsgh0JaY=
go.uber.org/goleak v1.1.11 h1:wy28qYRKZgnJTxGxvye5/wgWr1EKjmUDGYox5mGlRlI=
go.uber.org/goleak v1.1.11/go.mod h1:cwTWslyiVhfpKIDGSZEM2HlOvcqm+tG4zioyIeLoqMQ=
+go.uber.org/multierr v1.1.0/go.mod h1:wR5kodmAFQ0UK8QlbwjlSNy0Z68gJhDJUG5sjR94q/0=
go.uber.org/multierr v1.5.0/go.mod h1:FeouvMocqHpRaaGuG9EjoKcStLC43Zu/fmqdUMPcKYU=
go.uber.org/multierr v1.6.0/go.mod h1:cdWPpRnG4AhwMwsgIHip0KRBQjJy5kYEpYjJxpXp9iU=
go.uber.org/multierr v1.7.0 h1:zaiO/rmgFjbmCXdSYJWQcdvOCsthmdaHfr3Gm2Kx4Ec=
go.uber.org/multierr v1.7.0/go.mod h1:7EAYxJLBy9rStEaz58O2t4Uvip6FSURkq8/ppBp95ak=
go.uber.org/tools v0.0.0-20190618225709-2cfd321de3ee/go.mod h1:vJERXedbb3MVM5f9Ejo0C68/HhF8uaILCdgjnY+goOA=
+go.uber.org/zap v1.10.0/go.mod h1:vwi/ZaCAaUcBkycHslxD9B2zi4UTXhF60s6SWpuDF0Q=
go.uber.org/zap v1.15.0/go.mod h1:Mb2vm2krFEG5DV0W9qcHBYFtp/Wku1cvYaqPsS/WYfc=
go.uber.org/zap v1.16.0/go.mod h1:MA8QOfq0BHJwdXa996Y4dYkAqRKB8/1K1QMMZVaNZjQ=
-go.uber.org/zap v1.21.0/go.mod h1:wjWOCqI0f2ZZrJF/UufIOkiC8ii6tm1iqIsLo76RfJw=
go.uber.org/zap v1.22.0 h1:Zcye5DUgBloQ9BaT4qc9BnjOFog5TvBSAGkJ3Nf70c0=
go.uber.org/zap v1.22.0/go.mod h1:H4siCOZOrAolnUPJEkfaSjDqyP+BDS0DdDWzwcgt3+U=
golang.org/x/crypto v0.0.0-20180904163835-0709b304e793/go.mod h1:6SG95UA2DQfeDnfUPMdvaQW0Q7yPrPDi9nlGo2tz2b4=
@@ -371,6 +374,7 @@ gopkg.in/alecthomas/kingpin.v2 v2.2.6/go.mod h1:FMv+mEhP44yOT+4EoQTLFTRgOQ1FBLks
gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
gopkg.in/check.v1 v1.0.0-20180628173108-788fd7840127/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
gopkg.in/check.v1 v1.0.0-20190902080502-41f04d3bba15/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
+gopkg.in/check.v1 v1.0.0-20200227125254-8fa46927fb4f/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c h1:Hei/4ADfdWqJk1ZMxUNpqntNwaWcugrBjAiHlqqRiVk=
gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c/go.mod h1:JHkPIbrfpd72SG/EVd6muEfDQjcINNoR0C8j2r3qZ4Q=
gopkg.in/errgo.v2 v2.1.0/go.mod h1:hNsd1EY+bozCKY1Ytp96fpM3vjJbqLJn88ws8XvfDNI=
@@ -382,8 +386,6 @@ gopkg.in/yaml.v2 v2.2.1/go.mod h1:hI93XBmqTisBFMUTm0b8Fm+jr3Dg1NNxqwp+5A1VGuI=
gopkg.in/yaml.v2 v2.2.2/go.mod h1:hI93XBmqTisBFMUTm0b8Fm+jr3Dg1NNxqwp+5A1VGuI=
gopkg.in/yaml.v2 v2.2.4/go.mod h1:hI93XBmqTisBFMUTm0b8Fm+jr3Dg1NNxqwp+5A1VGuI=
gopkg.in/yaml.v2 v2.2.5/go.mod h1:hI93XBmqTisBFMUTm0b8Fm+jr3Dg1NNxqwp+5A1VGuI=
-gopkg.in/yaml.v2 v2.2.7/go.mod h1:hI93XBmqTisBFMUTm0b8Fm+jr3Dg1NNxqwp+5A1VGuI=
-gopkg.in/yaml.v2 v2.2.8/go.mod h1:hI93XBmqTisBFMUTm0b8Fm+jr3Dg1NNxqwp+5A1VGuI=
gopkg.in/yaml.v2 v2.3.0/go.mod h1:hI93XBmqTisBFMUTm0b8Fm+jr3Dg1NNxqwp+5A1VGuI=
gopkg.in/yaml.v2 v2.4.0 h1:D8xgwECY7CYvx+Y2n4sBz93Jn9JRvxdiyyo8CTfuKaY=
gopkg.in/yaml.v2 v2.4.0/go.mod h1:RDklbk79AGWmwhnvt/jBztapEOGDOx6ZbXqjP6csGnQ=
diff --git a/eventmesh-server-go/pkg/connector/action.go b/eventmesh-server-go/pkg/connector/action.go
new file mode 100644
index 00000000..c92fb721
--- /dev/null
+++ b/eventmesh-server-go/pkg/connector/action.go
@@ -0,0 +1,9 @@
+package connector
+
+type EventMeshAction string
+
+const (
+ CommitMessage EventMeshAction = "CommitMessage"
+ ReconsumeLater EventMeshAction = "ReconsumeLater"
+ ManualAck EventMeshAction = "ManualAck"
+)
diff --git a/eventmesh-server-go/go.mod b/eventmesh-server-go/pkg/connector/consumer.go
similarity index 56%
copy from eventmesh-server-go/go.mod
copy to eventmesh-server-go/pkg/connector/consumer.go
index 62d5e28a..563959fd 100644
--- a/eventmesh-server-go/go.mod
+++ b/eventmesh-server-go/pkg/connector/consumer.go
@@ -13,27 +13,23 @@
// See the License for the specific language governing permissions and
// limitations under the License.
-module github.com/apache/incubator-eventmesh/eventmesh-server-go
-
-go 1.16
-
-require (
- github.com/gogf/gf v1.16.9
- github.com/hashicorp/go-multierror v1.1.1
- github.com/lestrrat-go/strftime v1.0.6
- github.com/nacos-group/nacos-sdk-go/v2 v2.0.3
- github.com/pkg/errors v0.9.1
- github.com/spf13/cobra v1.5.0
- go.uber.org/zap v1.22.0
- google.golang.org/grpc v1.36.1
- gopkg.in/yaml.v3 v3.0.1
-)
+package connector
-require (
- github.com/BurntSushi/toml v1.2.0 // indirect
- github.com/gin-contrib/pprof v1.4.0
- github.com/gin-gonic/gin v1.8.1
- github.com/panjf2000/gnet/v2 v2.1.1
- github.com/unrolled/secure v1.12.0
- go.uber.org/fx v1.18.1
+import (
+ cloudevents "github.com/cloudevents/sdk-go/v2"
)
+
+// Consumer message for eventmesh standalone connector
+type Consumer interface {
+ LifeCycle
+
+ Initialize(*Properties) error
+
+ UpdateOffset([]cloudevents.Event)
+
+ Subscribe(string)
+
+ UnSubscribe(string)
+
+ RegisterEventListener(EventListener)
+}
diff --git a/eventmesh-server-go/go.mod b/eventmesh-server-go/pkg/connector/lifecycle.go
similarity index 56%
copy from eventmesh-server-go/go.mod
copy to eventmesh-server-go/pkg/connector/lifecycle.go
index 62d5e28a..0dd04c6a 100644
--- a/eventmesh-server-go/go.mod
+++ b/eventmesh-server-go/pkg/connector/lifecycle.go
@@ -13,27 +13,12 @@
// See the License for the specific language governing permissions and
// limitations under the License.
-module github.com/apache/incubator-eventmesh/eventmesh-server-go
+package connector
-go 1.16
-
-require (
- github.com/gogf/gf v1.16.9
- github.com/hashicorp/go-multierror v1.1.1
- github.com/lestrrat-go/strftime v1.0.6
- github.com/nacos-group/nacos-sdk-go/v2 v2.0.3
- github.com/pkg/errors v0.9.1
- github.com/spf13/cobra v1.5.0
- go.uber.org/zap v1.22.0
- google.golang.org/grpc v1.36.1
- gopkg.in/yaml.v3 v3.0.1
-)
-
-require (
- github.com/BurntSushi/toml v1.2.0 // indirect
- github.com/gin-contrib/pprof v1.4.0
- github.com/gin-gonic/gin v1.8.1
- github.com/panjf2000/gnet/v2 v2.1.1
- github.com/unrolled/secure v1.12.0
- go.uber.org/fx v1.18.1
-)
+// LifeCycle defines a lifecycle interface for a OMS related service endpoint,
+type LifeCycle interface {
+ IsStarted() bool
+ IsClosed() bool
+ Start()
+ Shutdown()
+}
diff --git a/eventmesh-server-go/go.mod b/eventmesh-server-go/pkg/connector/listener.go
similarity index 56%
copy from eventmesh-server-go/go.mod
copy to eventmesh-server-go/pkg/connector/listener.go
index 62d5e28a..32f3bdef 100644
--- a/eventmesh-server-go/go.mod
+++ b/eventmesh-server-go/pkg/connector/listener.go
@@ -13,27 +13,13 @@
// See the License for the specific language governing permissions and
// limitations under the License.
-module github.com/apache/incubator-eventmesh/eventmesh-server-go
+package connector
-go 1.16
-
-require (
- github.com/gogf/gf v1.16.9
- github.com/hashicorp/go-multierror v1.1.1
- github.com/lestrrat-go/strftime v1.0.6
- github.com/nacos-group/nacos-sdk-go/v2 v2.0.3
- github.com/pkg/errors v0.9.1
- github.com/spf13/cobra v1.5.0
- go.uber.org/zap v1.22.0
- google.golang.org/grpc v1.36.1
- gopkg.in/yaml.v3 v3.0.1
+import (
+ cloudevents "github.com/cloudevents/sdk-go/v2"
)
-require (
- github.com/BurntSushi/toml v1.2.0 // indirect
- github.com/gin-contrib/pprof v1.4.0
- github.com/gin-gonic/gin v1.8.1
- github.com/panjf2000/gnet/v2 v2.1.1
- github.com/unrolled/secure v1.12.0
- go.uber.org/fx v1.18.1
-)
+// EventListener listener to consume the cloudevents message
+type EventListener interface {
+ Consume(cloudevents.Event)
+}
diff --git a/eventmesh-server-go/pkg/connector/message_queue.go b/eventmesh-server-go/pkg/connector/message_queue.go
new file mode 100644
index 00000000..d933d9f4
--- /dev/null
+++ b/eventmesh-server-go/pkg/connector/message_queue.go
@@ -0,0 +1,165 @@
+// 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 standalone
+
+var (
+ // ErrMessageDeleted message has been deleted with invalid offset
+ ErrMessageDeleted = fmt.Errorf("message has been deleted")
+)
+
+// MessageQueue is a block queue, can get entity by offset.
+// The queue is a FIFO data structure.
+type MessageQueue struct {
+ items []*MessageEntity
+ takeIndex int
+ putIndex int
+ count int
+ capacity int
+ lock *sync.Mutex
+}
+
+// NewDefaultMessageQueue create new message queue with
+// capacity size 2<<10
+func NewDefaultMessageQueue() *MessageQueue {
+ return NewMessageQueueWithCapacity(2 << 10)
+}
+
+// NewMessageQueueWithCapacity crate message queue with
+// given capacity
+func NewMessageQueueWithCapacity(capacity int) *MessageQueue {
+ return &MessageQueue{
+ items: make(*MessageEntity, capacity),
+ lock: new(sync.Mutex),
+ capacity: capacity,
+ }
+}
+
+// Put insert the message at the tail of this queue,
+// waiting for space to become available if the queue is full
+func (m *MessageQueue) Put(msg *MessageEntity) error {
+ m.lock.Lock()
+ defer m.lock.UnLock()
+ enqueue(msg)
+ return nil
+}
+
+// Take Get the first message at this queue,
+// waiting for the message is available if the queue is empty,
+// this method will not remove the message
+func (m *MessageQueue) Take() *MessageEntity {
+ m.lock.Lock()
+ defer m.lock.UnLock()
+ return dequeue()
+}
+
+// Peek Get the first message at this queue,
+// if the queue is empty return null immediately
+func (m *MessageQueue) Peek() *MessageEntity {
+ m.lock.Lock()
+ defer m.lock.UnLock()
+ return m.ItemAt(m.takeIndex)
+}
+
+// GetHead Get the head in this queue
+func (m *MessageQueue) GetHead() *MessageEntity {
+ return m.Peek()
+}
+
+// GetTail Get the tail in this queue
+func (m *MessageQueue) GetTail() *MessageEntity {
+ m.lock.Lock()
+ defer m.lock.UnLock()
+
+ if m.count == 0 {
+ return nil
+ }
+
+ tailIndex := m.putIndex - 1
+ if tailIndex < 0 {
+ tailIndex += len(m.items)
+ }
+
+ return m.ItemAt(tailIndex)
+}
+
+// Get the message by offset, since the offset is increment,
+// so we can get the first message in this queue
+// and calculate the index of this offset
+func (m *MessageQueue) GetByOffset(offset long) (*MessageEntity, error) {
+ m.lock.Lock()
+ defer m.lock.UnLock()
+
+ head := m.GetHead()
+ if head == nil {
+ return nil, nil
+ }
+
+ if head.Offset > offset {
+ return nil, ErrMessageDeleted
+ }
+
+ tail := m.GetTail()
+ if tail == nil || tail.Offset < offset {
+ return nil, nil
+ }
+
+ offsetDis := head.Offset - offset
+ offsetIndex := m.takeIndex - offsetDis
+ if offsetIndex < 0 {
+ offsetIndex += len(m.items)
+ }
+
+ return m.items[offsetIndex], nil
+}
+
+func (m *MessageQueue) RemoveHead() {
+ m.lock.Lock()
+ defer m.lock.UnLock()
+ if m.count = 0 {
+ return
+ }
+ m.items[m.takeIndex] = nil
+ m.takeIndex++
+ if m.takeIndex == len(m.items) {
+ m.takeIndex = 0
+ }
+}
+
+func (m *MessageQueue) ItemAt(idx int) *MessageEntity {
+ return m.items[idx]
+}
+
+func (m *MessageQueue) GetSize() int {
+ return m.count
+}
+
+func (m *MessageQueue) enqueue(msg *MessageEntity) {
+ m.items = append(m.items, msg)
+ putIndex++
+ if putIndex == len(m.items) {
+ putIndex = 0
+ }
+ m.count++
+}
+
+func (m *MessageQueue) dequeue() *MessageEntity {
+ item := m.items[m.takeIndex]
+ m.takeIndex++
+ if m.takeIndex == len(m.items) {
+ m.takeIndex == 0
+ }
+ return item
+}
diff --git a/eventmesh-server-go/go.mod b/eventmesh-server-go/pkg/connector/properties.go
similarity index 56%
copy from eventmesh-server-go/go.mod
copy to eventmesh-server-go/pkg/connector/properties.go
index 62d5e28a..86523a1e 100644
--- a/eventmesh-server-go/go.mod
+++ b/eventmesh-server-go/pkg/connector/properties.go
@@ -13,27 +13,10 @@
// See the License for the specific language governing permissions and
// limitations under the License.
-module github.com/apache/incubator-eventmesh/eventmesh-server-go
+package connector
-go 1.16
-
-require (
- github.com/gogf/gf v1.16.9
- github.com/hashicorp/go-multierror v1.1.1
- github.com/lestrrat-go/strftime v1.0.6
- github.com/nacos-group/nacos-sdk-go/v2 v2.0.3
- github.com/pkg/errors v0.9.1
- github.com/spf13/cobra v1.5.0
- go.uber.org/zap v1.22.0
- google.golang.org/grpc v1.36.1
- gopkg.in/yaml.v3 v3.0.1
-)
-
-require (
- github.com/BurntSushi/toml v1.2.0 // indirect
- github.com/gin-contrib/pprof v1.4.0
- github.com/gin-gonic/gin v1.8.1
- github.com/panjf2000/gnet/v2 v2.1.1
- github.com/unrolled/secure v1.12.0
- go.uber.org/fx v1.18.1
-)
+// Properties represents a persistent set of properties.
+// The Properties can be saved to a stream
+// or loaded from a stream. Each key and its corresponding value in
+// the property list is a string.
+type Properties map[string]interface{}
diff --git a/eventmesh-server-go/pkg/connector/publisher.go b/eventmesh-server-go/pkg/connector/publisher.go
new file mode 100644
index 00000000..d54f870c
--- /dev/null
+++ b/eventmesh-server-go/pkg/connector/publisher.go
@@ -0,0 +1,62 @@
+// 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 connector
+
+import (
+ cloudevents "github.com/cloudevents/sdk-go/v2"
+)
+
+type SendResult struct {
+ MessageId string `json:"messageId"`
+ Mopic string `json:"topic"`
+}
+
+type SendErrResult struct {
+ *SendResult
+ Err error `json:"err"`
+}
+
+// SendCallback callback for send message to eventmesh
+type SendCallback interface {
+ OnSuccess(*SendResult)
+
+ OnError(*SendErrResult)
+}
+
+type RequestReplyCallback interface {
+ OnSuccess(*SendResult)
+
+ OnError(*SendErrResult)
+}
+
+// Consumer message for eventmesh standalone connector
+type Producer interface {
+ LifeCycle
+
+ Initialize(*Properties) error
+
+ Publish(cloudevents.Event, SendCallback) error
+
+ SendOneway(cloudevents.Event) error
+
+ Request(cloudevents.Event, RequestReplyCallback, time.Duration) error
+
+ Reply(cloudevents.Event, SendCallback) error
+
+ CheckTopicExist(string) bool
+
+ SetExtFields();
+}
diff --git a/eventmesh-server-go/pkg/connector/standalone/broker.go b/eventmesh-server-go/pkg/connector/standalone/broker.go
new file mode 100644
index 00000000..050777c9
--- /dev/null
+++ b/eventmesh-server-go/pkg/connector/standalone/broker.go
@@ -0,0 +1,66 @@
+// 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 standalone
+
+var (
+ // messageStoreWindow msg ttl,
+ // If the currentTimeMills - messageCreateTimeMills >= MESSAGE_STORE_WINDOW,
+ // then the message will be clear, default to 1 hour
+ messageStoreWindow = time.Hour
+)
+
+// Broker used to store event, it just support standalone mode,
+// you shouldn't use this module in production environment
+type Broker struct {
+ // messageContainer store the topic and the queue
+ // key = TopicMetadata value = MessageQueue
+ messageContainer *sync.Map
+ // offsetMap store the offset for topic
+ // key = TopicMetadata value = atomic.Long
+ offsetMap *sync.Map
+}
+
+func NewBroker(ctx context.Context) *Broker {
+ b := &Broker{
+ messageContainer: new(sync.Map),
+ offsetMap: new(sync.Map),
+ }
+ return b
+}
+
+func (b *Broker) startHistoryMessageCleanTask(ctx context.Context) {
+ go func() {
+ cleanTicker := time.NewTicker(time.Second)
+ for {
+ select {
+ case <-ctx.Done():
+ return
+ case <-cleanTicker.C:
+ b.messageContainer.Range(func(k, v interface{}) bool {
+ now := time.Now()
+ currentMsg := v.(*MessageQueue).GetHead()
+ if currentMsg == nil {
+ return
+ }
+ if now.Sub(currentMsg.Create) > messageStoreWindow {
+ v.(*MessageQueue).RemoveHead()
+ }
+ return true
+ })
+ }
+ }
+ }()
+}
diff --git a/eventmesh-server-go/go.mod b/eventmesh-server-go/pkg/connector/standalone/message_entity.go
similarity index 56%
copy from eventmesh-server-go/go.mod
copy to eventmesh-server-go/pkg/connector/standalone/message_entity.go
index 62d5e28a..52216855 100644
--- a/eventmesh-server-go/go.mod
+++ b/eventmesh-server-go/pkg/connector/standalone/message_entity.go
@@ -13,27 +13,19 @@
// See the License for the specific language governing permissions and
// limitations under the License.
-module github.com/apache/incubator-eventmesh/eventmesh-server-go
+package standalone
-go 1.16
-
-require (
- github.com/gogf/gf v1.16.9
- github.com/hashicorp/go-multierror v1.1.1
- github.com/lestrrat-go/strftime v1.0.6
- github.com/nacos-group/nacos-sdk-go/v2 v2.0.3
- github.com/pkg/errors v0.9.1
- github.com/spf13/cobra v1.5.0
- go.uber.org/zap v1.22.0
- google.golang.org/grpc v1.36.1
- gopkg.in/yaml.v3 v3.0.1
+import (
+ cloudevents "github.com/cloudevents/sdk-go/v2"
)
-require (
- github.com/BurntSushi/toml v1.2.0 // indirect
- github.com/gin-contrib/pprof v1.4.0
- github.com/gin-gonic/gin v1.8.1
- github.com/panjf2000/gnet/v2 v2.1.1
- github.com/unrolled/secure v1.12.0
- go.uber.org/fx v1.18.1
-)
+type TopicMetadata struct {
+ TopicName string `json:"topicName"`
+}
+
+type MessageEntity struct {
+ TopicMetadata *TopicMetadata `json:"topicMetadata`
+ Message cloudevents.Event `json:"message"`
+ Offset long `json:"offset"`
+ createTime time.Time `json:"createTime"`
+}
diff --git a/eventmesh-server-go/pkg/runtime/emserver/tcp.go b/eventmesh-server-go/pkg/runtime/emserver/tcp.go
index 781f4215..7433c91c 100644
--- a/eventmesh-server-go/pkg/runtime/emserver/tcp.go
+++ b/eventmesh-server-go/pkg/runtime/emserver/tcp.go
@@ -16,15 +16,11 @@
package emserver
import (
- "fmt"
"github.com/apache/incubator-eventmesh/eventmesh-server-go/config"
- "github.com/panjf2000/gnet/v2"
)
type TCPServer struct {
tcpOption *config.TCPOption
- eng gnet.Engine
- gnet.BuiltinEventEngine
}
func NewTCPServer(opt *config.TCPOption) (GracefulServer, error) {
@@ -34,17 +30,9 @@ func NewTCPServer(opt *config.TCPOption) (GracefulServer, error) {
}
func (t *TCPServer) Serve() error {
- return gnet.Run(t, fmt.Sprintf("tcp://:%d", t.tcpOption.Port), gnet.WithMulticore(t.tcpOption.Multicore))
+ return nil
}
func (t *TCPServer) Stop() error {
return nil
}
-
-func (t *TCPServer) OnBoot(eng gnet.Engine) gnet.Action {
- return gnet.None
-}
-
-func (t *TCPServer) OnTraffic(c gnet.Conn) gnet.Action {
- return gnet.None
-}
diff --git a/eventmesh-server-go/pkg/runtime/proto/pb/README.md b/eventmesh-server-go/pkg/runtime/proto/pb/README.md
new file mode 100644
index 00000000..8210f63c
--- /dev/null
+++ b/eventmesh-server-go/pkg/runtime/proto/pb/README.md
@@ -0,0 +1,20 @@
+how to generate proto go files
+---
+generate go files by the protoc-gen-go owned by Google
+
+1. install the latest protoc-gen-go and protoc-gen-go-grpc
+```
+go install google.golang.org/protobuf/cmd/protoc-gen-go
+go install google.golang.org/grpc/cmd/protoc-gen-go-grpc
+```
+2. run command
+```
+protoc --go_out=. eventmesh-client.proto
+protoc --go-grpc_out=. eventmesh-client.proto
+```
+
+if you use the latest version protoc-gen-go, and generate by the old command, you will got these error:
+```
+--go_out: protoc-gen-go: plugins are not supported; use 'protoc --go-grpc_out=...' to generate gRPC
+```
+
diff --git a/eventmesh-server-go/pkg/runtime/proto/pb/eventmesh-client.pb.go b/eventmesh-server-go/pkg/runtime/proto/pb/eventmesh-client.pb.go
new file mode 100644
index 00000000..d5ef31c2
--- /dev/null
+++ b/eventmesh-server-go/pkg/runtime/proto/pb/eventmesh-client.pb.go
@@ -0,0 +1,1485 @@
+//
+// 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.
+
+// Code generated by protoc-gen-go. DO NOT EDIT.
+// versions:
+// protoc-gen-go v1.28.0
+// protoc v3.19.4
+// source: eventmesh-client.proto
+
+//package eventmesh.common.protocol.grpc;
+//
+//option java_multiple_files = true;
+//option java_package = "org.apache.eventmesh.common.protocol.grpc.protos";
+//option java_outer_classname = "EventmeshGrpc";
+
+// make sure the protoc and protoc-gen-go is installed on your machine, and has set
+// its directory into path
+// download protoc: https://github.com/protocolbuffers/protobuf/releases
+// install protoc-gen-go: go install google.golang.org/protobuf/cmd/protoc-gen-go
+// generate go code by protoc: protoc --go_out=. eventmesh-client.proto
+
+package pb
+
+import (
+ protoreflect "google.golang.org/protobuf/reflect/protoreflect"
+ protoimpl "google.golang.org/protobuf/runtime/protoimpl"
+ reflect "reflect"
+ sync "sync"
+)
+
+const (
+ // Verify that this generated code is sufficiently up-to-date.
+ _ = protoimpl.EnforceVersion(20 - protoimpl.MinVersion)
+ // Verify that runtime/protoimpl is sufficiently up-to-date.
+ _ = protoimpl.EnforceVersion(protoimpl.MaxVersion - 20)
+)
+
+type Subscription_SubscriptionItem_SubscriptionMode int32
+
+const (
+ Subscription_SubscriptionItem_CLUSTERING Subscription_SubscriptionItem_SubscriptionMode = 0
+ Subscription_SubscriptionItem_BROADCASTING Subscription_SubscriptionItem_SubscriptionMode = 1
+)
+
+// Enum value maps for Subscription_SubscriptionItem_SubscriptionMode.
+var (
+ Subscription_SubscriptionItem_SubscriptionMode_name = map[int32]string{
+ 0: "CLUSTERING",
+ 1: "BROADCASTING",
+ }
+ Subscription_SubscriptionItem_SubscriptionMode_value = map[string]int32{
+ "CLUSTERING": 0,
+ "BROADCASTING": 1,
+ }
+)
+
+func (x Subscription_SubscriptionItem_SubscriptionMode) Enum() *Subscription_SubscriptionItem_SubscriptionMode {
+ p := new(Subscription_SubscriptionItem_SubscriptionMode)
+ *p = x
+ return p
+}
+
+func (x Subscription_SubscriptionItem_SubscriptionMode) String() string {
+ return protoimpl.X.EnumStringOf(x.Descriptor(), protoreflect.EnumNumber(x))
+}
+
+func (Subscription_SubscriptionItem_SubscriptionMode) Descriptor() protoreflect.EnumDescriptor {
+ return file_eventmesh_client_proto_enumTypes[0].Descriptor()
+}
+
+func (Subscription_SubscriptionItem_SubscriptionMode) Type() protoreflect.EnumType {
+ return &file_eventmesh_client_proto_enumTypes[0]
+}
+
+func (x Subscription_SubscriptionItem_SubscriptionMode) Number() protoreflect.EnumNumber {
+ return protoreflect.EnumNumber(x)
+}
+
+// Deprecated: Use Subscription_SubscriptionItem_SubscriptionMode.Descriptor instead.
+func (Subscription_SubscriptionItem_SubscriptionMode) EnumDescriptor() ([]byte, []int) {
+ return file_eventmesh_client_proto_rawDescGZIP(), []int{4, 0, 0}
+}
+
+type Subscription_SubscriptionItem_SubscriptionType int32
+
+const (
+ Subscription_SubscriptionItem_ASYNC Subscription_SubscriptionItem_SubscriptionType = 0
+ Subscription_SubscriptionItem_SYNC Subscription_SubscriptionItem_SubscriptionType = 1
+)
+
+// Enum value maps for Subscription_SubscriptionItem_SubscriptionType.
+var (
+ Subscription_SubscriptionItem_SubscriptionType_name = map[int32]string{
+ 0: "ASYNC",
+ 1: "SYNC",
+ }
+ Subscription_SubscriptionItem_SubscriptionType_value = map[string]int32{
+ "ASYNC": 0,
+ "SYNC": 1,
+ }
+)
+
+func (x Subscription_SubscriptionItem_SubscriptionType) Enum() *Subscription_SubscriptionItem_SubscriptionType {
+ p := new(Subscription_SubscriptionItem_SubscriptionType)
+ *p = x
+ return p
+}
+
+func (x Subscription_SubscriptionItem_SubscriptionType) String() string {
+ return protoimpl.X.EnumStringOf(x.Descriptor(), protoreflect.EnumNumber(x))
+}
+
+func (Subscription_SubscriptionItem_SubscriptionType) Descriptor() protoreflect.EnumDescriptor {
+ return file_eventmesh_client_proto_enumTypes[1].Descriptor()
+}
+
+func (Subscription_SubscriptionItem_SubscriptionType) Type() protoreflect.EnumType {
+ return &file_eventmesh_client_proto_enumTypes[1]
+}
+
+func (x Subscription_SubscriptionItem_SubscriptionType) Number() protoreflect.EnumNumber {
+ return protoreflect.EnumNumber(x)
+}
+
+// Deprecated: Use Subscription_SubscriptionItem_SubscriptionType.Descriptor instead.
+func (Subscription_SubscriptionItem_SubscriptionType) EnumDescriptor() ([]byte, []int) {
+ return file_eventmesh_client_proto_rawDescGZIP(), []int{4, 0, 1}
+}
+
+type Heartbeat_ClientType int32
+
+const (
+ Heartbeat_PUB Heartbeat_ClientType = 0
+ Heartbeat_SUB Heartbeat_ClientType = 1
+)
+
+// Enum value maps for Heartbeat_ClientType.
+var (
+ Heartbeat_ClientType_name = map[int32]string{
+ 0: "PUB",
+ 1: "SUB",
+ }
+ Heartbeat_ClientType_value = map[string]int32{
+ "PUB": 0,
+ "SUB": 1,
+ }
+)
+
+func (x Heartbeat_ClientType) Enum() *Heartbeat_ClientType {
+ p := new(Heartbeat_ClientType)
+ *p = x
+ return p
+}
+
+func (x Heartbeat_ClientType) String() string {
+ return protoimpl.X.EnumStringOf(x.Descriptor(), protoreflect.EnumNumber(x))
+}
+
+func (Heartbeat_ClientType) Descriptor() protoreflect.EnumDescriptor {
+ return file_eventmesh_client_proto_enumTypes[2].Descriptor()
+}
+
+func (Heartbeat_ClientType) Type() protoreflect.EnumType {
+ return &file_eventmesh_client_proto_enumTypes[2]
+}
+
+func (x Heartbeat_ClientType) Number() protoreflect.EnumNumber {
+ return protoreflect.EnumNumber(x)
+}
+
+// Deprecated: Use Heartbeat_ClientType.Descriptor instead.
+func (Heartbeat_ClientType) EnumDescriptor() ([]byte, []int) {
+ return file_eventmesh_client_proto_rawDescGZIP(), []int{5, 0}
+}
+
+type RequestHeader struct {
+ state protoimpl.MessageState
+ sizeCache protoimpl.SizeCache
+ unknownFields protoimpl.UnknownFields
+
+ Env string `protobuf:"bytes,1,opt,name=env,proto3" json:"env,omitempty"`
+ Region string `protobuf:"bytes,2,opt,name=region,proto3" json:"region,omitempty"`
+ Idc string `protobuf:"bytes,3,opt,name=idc,proto3" json:"idc,omitempty"`
+ Ip string `protobuf:"bytes,4,opt,name=ip,proto3" json:"ip,omitempty"`
+ Pid string `protobuf:"bytes,5,opt,name=pid,proto3" json:"pid,omitempty"`
+ Sys string `protobuf:"bytes,6,opt,name=sys,proto3" json:"sys,omitempty"`
+ Username string `protobuf:"bytes,7,opt,name=username,proto3" json:"username,omitempty"`
+ Password string `protobuf:"bytes,8,opt,name=password,proto3" json:"password,omitempty"`
+ Language string `protobuf:"bytes,9,opt,name=language,proto3" json:"language,omitempty"`
+ ProtocolType string `protobuf:"bytes,10,opt,name=protocolType,proto3" json:"protocolType,omitempty"`
+ ProtocolVersion string `protobuf:"bytes,11,opt,name=protocolVersion,proto3" json:"protocolVersion,omitempty"`
+ ProtocolDesc string `protobuf:"bytes,12,opt,name=protocolDesc,proto3" json:"protocolDesc,omitempty"`
+}
+
+func (x *RequestHeader) Reset() {
+ *x = RequestHeader{}
+ if protoimpl.UnsafeEnabled {
+ mi := &file_eventmesh_client_proto_msgTypes[0]
+ ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
+ ms.StoreMessageInfo(mi)
+ }
+}
+
+func (x *RequestHeader) String() string {
+ return protoimpl.X.MessageStringOf(x)
+}
+
+func (*RequestHeader) ProtoMessage() {}
+
+func (x *RequestHeader) ProtoReflect() protoreflect.Message {
+ mi := &file_eventmesh_client_proto_msgTypes[0]
+ if protoimpl.UnsafeEnabled && x != nil {
+ ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
+ if ms.LoadMessageInfo() == nil {
+ ms.StoreMessageInfo(mi)
+ }
+ return ms
+ }
+ return mi.MessageOf(x)
+}
+
+// Deprecated: Use RequestHeader.ProtoReflect.Descriptor instead.
+func (*RequestHeader) Descriptor() ([]byte, []int) {
+ return file_eventmesh_client_proto_rawDescGZIP(), []int{0}
+}
+
+func (x *RequestHeader) GetEnv() string {
+ if x != nil {
+ return x.Env
+ }
+ return ""
+}
+
+func (x *RequestHeader) GetRegion() string {
+ if x != nil {
+ return x.Region
+ }
+ return ""
+}
+
+func (x *RequestHeader) GetIdc() string {
+ if x != nil {
+ return x.Idc
+ }
+ return ""
+}
+
+func (x *RequestHeader) GetIp() string {
+ if x != nil {
+ return x.Ip
+ }
+ return ""
+}
+
+func (x *RequestHeader) GetPid() string {
+ if x != nil {
+ return x.Pid
+ }
+ return ""
+}
+
+func (x *RequestHeader) GetSys() string {
+ if x != nil {
+ return x.Sys
+ }
+ return ""
+}
+
+func (x *RequestHeader) GetUsername() string {
+ if x != nil {
+ return x.Username
+ }
+ return ""
+}
+
+func (x *RequestHeader) GetPassword() string {
+ if x != nil {
+ return x.Password
+ }
+ return ""
+}
+
+func (x *RequestHeader) GetLanguage() string {
+ if x != nil {
+ return x.Language
+ }
+ return ""
+}
+
+func (x *RequestHeader) GetProtocolType() string {
+ if x != nil {
+ return x.ProtocolType
+ }
+ return ""
+}
+
+func (x *RequestHeader) GetProtocolVersion() string {
+ if x != nil {
+ return x.ProtocolVersion
+ }
+ return ""
+}
+
+func (x *RequestHeader) GetProtocolDesc() string {
+ if x != nil {
+ return x.ProtocolDesc
+ }
+ return ""
+}
+
+type SimpleMessage struct {
+ state protoimpl.MessageState
+ sizeCache protoimpl.SizeCache
+ unknownFields protoimpl.UnknownFields
+
+ Header *RequestHeader `protobuf:"bytes,1,opt,name=header,proto3" json:"header,omitempty"`
+ ProducerGroup string `protobuf:"bytes,2,opt,name=producerGroup,proto3" json:"producerGroup,omitempty"`
+ Topic string `protobuf:"bytes,3,opt,name=topic,proto3" json:"topic,omitempty"`
+ Content string `protobuf:"bytes,4,opt,name=content,proto3" json:"content,omitempty"`
+ Ttl string `protobuf:"bytes,5,opt,name=ttl,proto3" json:"ttl,omitempty"`
+ UniqueId string `protobuf:"bytes,6,opt,name=uniqueId,proto3" json:"uniqueId,omitempty"`
+ SeqNum string `protobuf:"bytes,7,opt,name=seqNum,proto3" json:"seqNum,omitempty"`
+ Tag string `protobuf:"bytes,8,opt,name=tag,proto3" json:"tag,omitempty"`
+ Properties map[string]string `protobuf:"bytes,9,rep,name=properties,proto3" json:"properties,omitempty" protobuf_key:"bytes,1,opt,name=key,proto3" protobuf_val:"bytes,2,opt,name=value,proto3"`
+}
+
+func (x *SimpleMessage) Reset() {
+ *x = SimpleMessage{}
+ if protoimpl.UnsafeEnabled {
+ mi := &file_eventmesh_client_proto_msgTypes[1]
+ ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
+ ms.StoreMessageInfo(mi)
+ }
+}
+
+func (x *SimpleMessage) String() string {
+ return protoimpl.X.MessageStringOf(x)
+}
+
+func (*SimpleMessage) ProtoMessage() {}
+
+func (x *SimpleMessage) ProtoReflect() protoreflect.Message {
+ mi := &file_eventmesh_client_proto_msgTypes[1]
+ if protoimpl.UnsafeEnabled && x != nil {
+ ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
+ if ms.LoadMessageInfo() == nil {
+ ms.StoreMessageInfo(mi)
+ }
+ return ms
+ }
+ return mi.MessageOf(x)
+}
+
+// Deprecated: Use SimpleMessage.ProtoReflect.Descriptor instead.
+func (*SimpleMessage) Descriptor() ([]byte, []int) {
+ return file_eventmesh_client_proto_rawDescGZIP(), []int{1}
+}
+
+func (x *SimpleMessage) GetHeader() *RequestHeader {
+ if x != nil {
+ return x.Header
+ }
+ return nil
+}
+
+func (x *SimpleMessage) GetProducerGroup() string {
+ if x != nil {
+ return x.ProducerGroup
+ }
+ return ""
+}
+
+func (x *SimpleMessage) GetTopic() string {
+ if x != nil {
+ return x.Topic
+ }
+ return ""
+}
+
+func (x *SimpleMessage) GetContent() string {
+ if x != nil {
+ return x.Content
+ }
+ return ""
+}
+
+func (x *SimpleMessage) GetTtl() string {
+ if x != nil {
+ return x.Ttl
+ }
+ return ""
+}
+
+func (x *SimpleMessage) GetUniqueId() string {
+ if x != nil {
+ return x.UniqueId
+ }
+ return ""
+}
+
+func (x *SimpleMessage) GetSeqNum() string {
+ if x != nil {
+ return x.SeqNum
+ }
+ return ""
+}
+
+func (x *SimpleMessage) GetTag() string {
+ if x != nil {
+ return x.Tag
+ }
+ return ""
+}
+
+func (x *SimpleMessage) GetProperties() map[string]string {
+ if x != nil {
+ return x.Properties
+ }
+ return nil
+}
+
+type BatchMessage struct {
+ state protoimpl.MessageState
+ sizeCache protoimpl.SizeCache
+ unknownFields protoimpl.UnknownFields
+
+ Header *RequestHeader `protobuf:"bytes,1,opt,name=header,proto3" json:"header,omitempty"`
+ ProducerGroup string `protobuf:"bytes,2,opt,name=producerGroup,proto3" json:"producerGroup,omitempty"`
+ Topic string `protobuf:"bytes,3,opt,name=topic,proto3" json:"topic,omitempty"`
+ MessageItem []*BatchMessage_MessageItem `protobuf:"bytes,4,rep,name=messageItem,proto3" json:"messageItem,omitempty"`
+}
+
+func (x *BatchMessage) Reset() {
+ *x = BatchMessage{}
+ if protoimpl.UnsafeEnabled {
+ mi := &file_eventmesh_client_proto_msgTypes[2]
+ ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
+ ms.StoreMessageInfo(mi)
+ }
+}
+
+func (x *BatchMessage) String() string {
+ return protoimpl.X.MessageStringOf(x)
+}
+
+func (*BatchMessage) ProtoMessage() {}
+
+func (x *BatchMessage) ProtoReflect() protoreflect.Message {
+ mi := &file_eventmesh_client_proto_msgTypes[2]
+ if protoimpl.UnsafeEnabled && x != nil {
+ ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
+ if ms.LoadMessageInfo() == nil {
+ ms.StoreMessageInfo(mi)
+ }
+ return ms
+ }
+ return mi.MessageOf(x)
+}
+
+// Deprecated: Use BatchMessage.ProtoReflect.Descriptor instead.
+func (*BatchMessage) Descriptor() ([]byte, []int) {
+ return file_eventmesh_client_proto_rawDescGZIP(), []int{2}
+}
+
+func (x *BatchMessage) GetHeader() *RequestHeader {
+ if x != nil {
+ return x.Header
+ }
+ return nil
+}
+
+func (x *BatchMessage) GetProducerGroup() string {
+ if x != nil {
+ return x.ProducerGroup
+ }
+ return ""
+}
+
+func (x *BatchMessage) GetTopic() string {
+ if x != nil {
+ return x.Topic
+ }
+ return ""
+}
+
+func (x *BatchMessage) GetMessageItem() []*BatchMessage_MessageItem {
+ if x != nil {
+ return x.MessageItem
+ }
+ return nil
+}
+
+type Response struct {
+ state protoimpl.MessageState
+ sizeCache protoimpl.SizeCache
+ unknownFields protoimpl.UnknownFields
+
+ RespCode string `protobuf:"bytes,1,opt,name=respCode,proto3" json:"respCode,omitempty"`
+ RespMsg string `protobuf:"bytes,2,opt,name=respMsg,proto3" json:"respMsg,omitempty"`
+ RespTime string `protobuf:"bytes,3,opt,name=respTime,proto3" json:"respTime,omitempty"`
+}
+
+func (x *Response) Reset() {
+ *x = Response{}
+ if protoimpl.UnsafeEnabled {
+ mi := &file_eventmesh_client_proto_msgTypes[3]
+ ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
+ ms.StoreMessageInfo(mi)
+ }
+}
+
+func (x *Response) String() string {
+ return protoimpl.X.MessageStringOf(x)
+}
+
+func (*Response) ProtoMessage() {}
+
+func (x *Response) ProtoReflect() protoreflect.Message {
+ mi := &file_eventmesh_client_proto_msgTypes[3]
+ if protoimpl.UnsafeEnabled && x != nil {
+ ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
+ if ms.LoadMessageInfo() == nil {
+ ms.StoreMessageInfo(mi)
+ }
+ return ms
+ }
+ return mi.MessageOf(x)
+}
+
+// Deprecated: Use Response.ProtoReflect.Descriptor instead.
+func (*Response) Descriptor() ([]byte, []int) {
+ return file_eventmesh_client_proto_rawDescGZIP(), []int{3}
+}
+
+func (x *Response) GetRespCode() string {
+ if x != nil {
+ return x.RespCode
+ }
+ return ""
+}
+
+func (x *Response) GetRespMsg() string {
+ if x != nil {
+ return x.RespMsg
+ }
+ return ""
+}
+
+func (x *Response) GetRespTime() string {
+ if x != nil {
+ return x.RespTime
+ }
+ return ""
+}
+
+type Subscription struct {
+ state protoimpl.MessageState
+ sizeCache protoimpl.SizeCache
+ unknownFields protoimpl.UnknownFields
+
+ Header *RequestHeader `protobuf:"bytes,1,opt,name=header,proto3" json:"header,omitempty"`
+ ConsumerGroup string `protobuf:"bytes,2,opt,name=consumerGroup,proto3" json:"consumerGroup,omitempty"`
+ SubscriptionItems []*Subscription_SubscriptionItem `protobuf:"bytes,3,rep,name=subscriptionItems,proto3" json:"subscriptionItems,omitempty"`
+ Url string `protobuf:"bytes,4,opt,name=url,proto3" json:"url,omitempty"`
+ Reply *Subscription_Reply `protobuf:"bytes,5,opt,name=reply,proto3" json:"reply,omitempty"`
+}
+
+func (x *Subscription) Reset() {
+ *x = Subscription{}
+ if protoimpl.UnsafeEnabled {
+ mi := &file_eventmesh_client_proto_msgTypes[4]
+ ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
+ ms.StoreMessageInfo(mi)
+ }
+}
+
+func (x *Subscription) String() string {
+ return protoimpl.X.MessageStringOf(x)
+}
+
+func (*Subscription) ProtoMessage() {}
+
+func (x *Subscription) ProtoReflect() protoreflect.Message {
+ mi := &file_eventmesh_client_proto_msgTypes[4]
+ if protoimpl.UnsafeEnabled && x != nil {
+ ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
+ if ms.LoadMessageInfo() == nil {
+ ms.StoreMessageInfo(mi)
+ }
+ return ms
+ }
+ return mi.MessageOf(x)
+}
+
+// Deprecated: Use Subscription.ProtoReflect.Descriptor instead.
+func (*Subscription) Descriptor() ([]byte, []int) {
+ return file_eventmesh_client_proto_rawDescGZIP(), []int{4}
+}
+
+func (x *Subscription) GetHeader() *RequestHeader {
+ if x != nil {
+ return x.Header
+ }
+ return nil
+}
+
+func (x *Subscription) GetConsumerGroup() string {
+ if x != nil {
+ return x.ConsumerGroup
+ }
+ return ""
+}
+
+func (x *Subscription) GetSubscriptionItems() []*Subscription_SubscriptionItem {
+ if x != nil {
+ return x.SubscriptionItems
+ }
+ return nil
+}
+
+func (x *Subscription) GetUrl() string {
+ if x != nil {
+ return x.Url
+ }
+ return ""
+}
+
+func (x *Subscription) GetReply() *Subscription_Reply {
+ if x != nil {
+ return x.Reply
+ }
+ return nil
+}
+
+type Heartbeat struct {
+ state protoimpl.MessageState
+ sizeCache protoimpl.SizeCache
+ unknownFields protoimpl.UnknownFields
+
+ Header *RequestHeader `protobuf:"bytes,1,opt,name=header,proto3" json:"header,omitempty"`
+ ClientType Heartbeat_ClientType `protobuf:"varint,2,opt,name=clientType,proto3,enum=eventmesh.common.protocol.grpc.Heartbeat_ClientType" json:"clientType,omitempty"`
+ ProducerGroup string `protobuf:"bytes,3,opt,name=producerGroup,proto3" json:"producerGroup,omitempty"`
+ ConsumerGroup string `protobuf:"bytes,4,opt,name=consumerGroup,proto3" json:"consumerGroup,omitempty"`
+ HeartbeatItems []*Heartbeat_HeartbeatItem `protobuf:"bytes,5,rep,name=heartbeatItems,proto3" json:"heartbeatItems,omitempty"`
+}
+
+func (x *Heartbeat) Reset() {
+ *x = Heartbeat{}
+ if protoimpl.UnsafeEnabled {
+ mi := &file_eventmesh_client_proto_msgTypes[5]
+ ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
+ ms.StoreMessageInfo(mi)
+ }
+}
+
+func (x *Heartbeat) String() string {
+ return protoimpl.X.MessageStringOf(x)
+}
+
+func (*Heartbeat) ProtoMessage() {}
+
+func (x *Heartbeat) ProtoReflect() protoreflect.Message {
+ mi := &file_eventmesh_client_proto_msgTypes[5]
+ if protoimpl.UnsafeEnabled && x != nil {
+ ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
+ if ms.LoadMessageInfo() == nil {
+ ms.StoreMessageInfo(mi)
+ }
+ return ms
+ }
+ return mi.MessageOf(x)
+}
+
+// Deprecated: Use Heartbeat.ProtoReflect.Descriptor instead.
+func (*Heartbeat) Descriptor() ([]byte, []int) {
+ return file_eventmesh_client_proto_rawDescGZIP(), []int{5}
+}
+
+func (x *Heartbeat) GetHeader() *RequestHeader {
+ if x != nil {
+ return x.Header
+ }
+ return nil
+}
+
+func (x *Heartbeat) GetClientType() Heartbeat_ClientType {
+ if x != nil {
+ return x.ClientType
+ }
+ return Heartbeat_PUB
+}
+
+func (x *Heartbeat) GetProducerGroup() string {
+ if x != nil {
+ return x.ProducerGroup
+ }
+ return ""
+}
+
+func (x *Heartbeat) GetConsumerGroup() string {
+ if x != nil {
+ return x.ConsumerGroup
+ }
+ return ""
+}
+
+func (x *Heartbeat) GetHeartbeatItems() []*Heartbeat_HeartbeatItem {
+ if x != nil {
+ return x.HeartbeatItems
+ }
+ return nil
+}
+
+type BatchMessage_MessageItem struct {
+ state protoimpl.MessageState
+ sizeCache protoimpl.SizeCache
+ unknownFields protoimpl.UnknownFields
+
+ Content string `protobuf:"bytes,1,opt,name=content,proto3" json:"content,omitempty"`
+ Ttl string `protobuf:"bytes,2,opt,name=ttl,proto3" json:"ttl,omitempty"`
+ UniqueId string `protobuf:"bytes,3,opt,name=uniqueId,proto3" json:"uniqueId,omitempty"`
+ SeqNum string `protobuf:"bytes,4,opt,name=seqNum,proto3" json:"seqNum,omitempty"`
+ Tag string `protobuf:"bytes,5,opt,name=tag,proto3" json:"tag,omitempty"`
+ Properties map[string]string `protobuf:"bytes,6,rep,name=properties,proto3" json:"properties,omitempty" protobuf_key:"bytes,1,opt,name=key,proto3" protobuf_val:"bytes,2,opt,name=value,proto3"`
+}
+
+func (x *BatchMessage_MessageItem) Reset() {
+ *x = BatchMessage_MessageItem{}
+ if protoimpl.UnsafeEnabled {
+ mi := &file_eventmesh_client_proto_msgTypes[7]
+ ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
+ ms.StoreMessageInfo(mi)
+ }
+}
+
+func (x *BatchMessage_MessageItem) String() string {
+ return protoimpl.X.MessageStringOf(x)
+}
+
+func (*BatchMessage_MessageItem) ProtoMessage() {}
+
+func (x *BatchMessage_MessageItem) ProtoReflect() protoreflect.Message {
+ mi := &file_eventmesh_client_proto_msgTypes[7]
+ if protoimpl.UnsafeEnabled && x != nil {
+ ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
+ if ms.LoadMessageInfo() == nil {
+ ms.StoreMessageInfo(mi)
+ }
+ return ms
+ }
+ return mi.MessageOf(x)
+}
+
+// Deprecated: Use BatchMessage_MessageItem.ProtoReflect.Descriptor instead.
+func (*BatchMessage_MessageItem) Descriptor() ([]byte, []int) {
+ return file_eventmesh_client_proto_rawDescGZIP(), []int{2, 0}
+}
+
+func (x *BatchMessage_MessageItem) GetContent() string {
+ if x != nil {
+ return x.Content
+ }
+ return ""
+}
+
+func (x *BatchMessage_MessageItem) GetTtl() string {
+ if x != nil {
+ return x.Ttl
+ }
+ return ""
+}
+
+func (x *BatchMessage_MessageItem) GetUniqueId() string {
+ if x != nil {
+ return x.UniqueId
+ }
+ return ""
+}
+
+func (x *BatchMessage_MessageItem) GetSeqNum() string {
+ if x != nil {
+ return x.SeqNum
+ }
+ return ""
+}
+
+func (x *BatchMessage_MessageItem) GetTag() string {
+ if x != nil {
+ return x.Tag
+ }
+ return ""
+}
+
+func (x *BatchMessage_MessageItem) GetProperties() map[string]string {
+ if x != nil {
+ return x.Properties
+ }
+ return nil
+}
+
+type Subscription_SubscriptionItem struct {
+ state protoimpl.MessageState
+ sizeCache protoimpl.SizeCache
+ unknownFields protoimpl.UnknownFields
+
+ Topic string `protobuf:"bytes,1,opt,name=topic,proto3" json:"topic,omitempty"`
+ Mode Subscription_SubscriptionItem_SubscriptionMode `protobuf:"varint,2,opt,name=mode,proto3,enum=eventmesh.common.protocol.grpc.Subscription_SubscriptionItem_SubscriptionMode" json:"mode,omitempty"`
+ Type Subscription_SubscriptionItem_SubscriptionType `protobuf:"varint,3,opt,name=type,proto3,enum=eventmesh.common.protocol.grpc.Subscription_SubscriptionItem_SubscriptionType" json:"type,omitempty"`
+}
+
+func (x *Subscription_SubscriptionItem) Reset() {
+ *x = Subscription_SubscriptionItem{}
+ if protoimpl.UnsafeEnabled {
+ mi := &file_eventmesh_client_proto_msgTypes[9]
+ ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
+ ms.StoreMessageInfo(mi)
+ }
+}
+
+func (x *Subscription_SubscriptionItem) String() string {
+ return protoimpl.X.MessageStringOf(x)
+}
+
+func (*Subscription_SubscriptionItem) ProtoMessage() {}
+
+func (x *Subscription_SubscriptionItem) ProtoReflect() protoreflect.Message {
+ mi := &file_eventmesh_client_proto_msgTypes[9]
+ if protoimpl.UnsafeEnabled && x != nil {
+ ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
+ if ms.LoadMessageInfo() == nil {
+ ms.StoreMessageInfo(mi)
+ }
+ return ms
+ }
+ return mi.MessageOf(x)
+}
+
+// Deprecated: Use Subscription_SubscriptionItem.ProtoReflect.Descriptor instead.
+func (*Subscription_SubscriptionItem) Descriptor() ([]byte, []int) {
+ return file_eventmesh_client_proto_rawDescGZIP(), []int{4, 0}
+}
+
+func (x *Subscription_SubscriptionItem) GetTopic() string {
+ if x != nil {
+ return x.Topic
+ }
+ return ""
+}
+
+func (x *Subscription_SubscriptionItem) GetMode() Subscription_SubscriptionItem_SubscriptionMode {
+ if x != nil {
+ return x.Mode
+ }
+ return Subscription_SubscriptionItem_CLUSTERING
+}
+
+func (x *Subscription_SubscriptionItem) GetType() Subscription_SubscriptionItem_SubscriptionType {
+ if x != nil {
+ return x.Type
+ }
+ return Subscription_SubscriptionItem_ASYNC
+}
+
+type Subscription_Reply struct {
+ state protoimpl.MessageState
+ sizeCache protoimpl.SizeCache
+ unknownFields protoimpl.UnknownFields
+
+ ProducerGroup string `protobuf:"bytes,1,opt,name=producerGroup,proto3" json:"producerGroup,omitempty"`
+ Topic string `protobuf:"bytes,2,opt,name=topic,proto3" json:"topic,omitempty"`
+ Content string `protobuf:"bytes,3,opt,name=content,proto3" json:"content,omitempty"`
+ Ttl string `protobuf:"bytes,4,opt,name=ttl,proto3" json:"ttl,omitempty"`
+ UniqueId string `protobuf:"bytes,5,opt,name=uniqueId,proto3" json:"uniqueId,omitempty"`
+ SeqNum string `protobuf:"bytes,6,opt,name=seqNum,proto3" json:"seqNum,omitempty"`
+ Tag string `protobuf:"bytes,7,opt,name=tag,proto3" json:"tag,omitempty"`
+ Properties map[string]string `protobuf:"bytes,8,rep,name=properties,proto3" json:"properties,omitempty" protobuf_key:"bytes,1,opt,name=key,proto3" protobuf_val:"bytes,2,opt,name=value,proto3"`
+}
+
+func (x *Subscription_Reply) Reset() {
+ *x = Subscription_Reply{}
+ if protoimpl.UnsafeEnabled {
+ mi := &file_eventmesh_client_proto_msgTypes[10]
+ ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
+ ms.StoreMessageInfo(mi)
+ }
+}
+
+func (x *Subscription_Reply) String() string {
+ return protoimpl.X.MessageStringOf(x)
+}
+
+func (*Subscription_Reply) ProtoMessage() {}
+
+func (x *Subscription_Reply) ProtoReflect() protoreflect.Message {
+ mi := &file_eventmesh_client_proto_msgTypes[10]
+ if protoimpl.UnsafeEnabled && x != nil {
+ ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
+ if ms.LoadMessageInfo() == nil {
+ ms.StoreMessageInfo(mi)
+ }
+ return ms
+ }
+ return mi.MessageOf(x)
+}
+
+// Deprecated: Use Subscription_Reply.ProtoReflect.Descriptor instead.
+func (*Subscription_Reply) Descriptor() ([]byte, []int) {
+ return file_eventmesh_client_proto_rawDescGZIP(), []int{4, 1}
+}
+
+func (x *Subscription_Reply) GetProducerGroup() string {
+ if x != nil {
+ return x.ProducerGroup
+ }
+ return ""
+}
+
+func (x *Subscription_Reply) GetTopic() string {
+ if x != nil {
+ return x.Topic
+ }
+ return ""
+}
+
+func (x *Subscription_Reply) GetContent() string {
+ if x != nil {
+ return x.Content
+ }
+ return ""
+}
+
+func (x *Subscription_Reply) GetTtl() string {
+ if x != nil {
+ return x.Ttl
+ }
+ return ""
+}
+
+func (x *Subscription_Reply) GetUniqueId() string {
+ if x != nil {
+ return x.UniqueId
+ }
+ return ""
+}
+
+func (x *Subscription_Reply) GetSeqNum() string {
+ if x != nil {
+ return x.SeqNum
+ }
+ return ""
+}
+
+func (x *Subscription_Reply) GetTag() string {
+ if x != nil {
+ return x.Tag
+ }
+ return ""
+}
+
+func (x *Subscription_Reply) GetProperties() map[string]string {
+ if x != nil {
+ return x.Properties
+ }
+ return nil
+}
+
+type Heartbeat_HeartbeatItem struct {
+ state protoimpl.MessageState
+ sizeCache protoimpl.SizeCache
+ unknownFields protoimpl.UnknownFields
+
+ Topic string `protobuf:"bytes,1,opt,name=topic,proto3" json:"topic,omitempty"`
+ Url string `protobuf:"bytes,2,opt,name=url,proto3" json:"url,omitempty"`
+}
+
+func (x *Heartbeat_HeartbeatItem) Reset() {
+ *x = Heartbeat_HeartbeatItem{}
+ if protoimpl.UnsafeEnabled {
+ mi := &file_eventmesh_client_proto_msgTypes[12]
+ ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
+ ms.StoreMessageInfo(mi)
+ }
+}
+
+func (x *Heartbeat_HeartbeatItem) String() string {
+ return protoimpl.X.MessageStringOf(x)
+}
+
+func (*Heartbeat_HeartbeatItem) ProtoMessage() {}
+
+func (x *Heartbeat_HeartbeatItem) ProtoReflect() protoreflect.Message {
+ mi := &file_eventmesh_client_proto_msgTypes[12]
+ if protoimpl.UnsafeEnabled && x != nil {
+ ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
+ if ms.LoadMessageInfo() == nil {
+ ms.StoreMessageInfo(mi)
+ }
+ return ms
+ }
+ return mi.MessageOf(x)
+}
+
+// Deprecated: Use Heartbeat_HeartbeatItem.ProtoReflect.Descriptor instead.
+func (*Heartbeat_HeartbeatItem) Descriptor() ([]byte, []int) {
+ return file_eventmesh_client_proto_rawDescGZIP(), []int{5, 0}
+}
+
+func (x *Heartbeat_HeartbeatItem) GetTopic() string {
+ if x != nil {
+ return x.Topic
+ }
+ return ""
+}
+
+func (x *Heartbeat_HeartbeatItem) GetUrl() string {
+ if x != nil {
+ return x.Url
+ }
+ return ""
+}
+
+var File_eventmesh_client_proto protoreflect.FileDescriptor
+
+var file_eventmesh_client_proto_rawDesc = []byte{
+ 0x0a, 0x16, 0x65, 0x76, 0x65, 0x6e, 0x74, 0x6d, 0x65, 0x73, 0x68, 0x2d, 0x63, 0x6c, 0x69, 0x65,
+ 0x6e, 0x74, 0x2e, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x12, 0x1e, 0x65, 0x76, 0x65, 0x6e, 0x74, 0x6d,
+ 0x65, 0x73, 0x68, 0x2e, 0x63, 0x6f, 0x6d, 0x6d, 0x6f, 0x6e, 0x2e, 0x70, 0x72, 0x6f, 0x74, 0x6f,
+ 0x63, 0x6f, 0x6c, 0x2e, 0x67, 0x72, 0x70, 0x63, 0x22, 0xc5, 0x02, 0x0a, 0x0d, 0x52, 0x65, 0x71,
+ 0x75, 0x65, 0x73, 0x74, 0x48, 0x65, 0x61, 0x64, 0x65, 0x72, 0x12, 0x10, 0x0a, 0x03, 0x65, 0x6e,
+ 0x76, 0x18, 0x01, 0x20, 0x01, 0x28, 0x09, 0x52, 0x03, 0x65, 0x6e, 0x76, 0x12, 0x16, 0x0a, 0x06,
+ 0x72, 0x65, 0x67, 0x69, 0x6f, 0x6e, 0x18, 0x02, 0x20, 0x01, 0x28, 0x09, 0x52, 0x06, 0x72, 0x65,
+ 0x67, 0x69, 0x6f, 0x6e, 0x12, 0x10, 0x0a, 0x03, 0x69, 0x64, 0x63, 0x18, 0x03, 0x20, 0x01, 0x28,
+ 0x09, 0x52, 0x03, 0x69, 0x64, 0x63, 0x12, 0x0e, 0x0a, 0x02, 0x69, 0x70, 0x18, 0x04, 0x20, 0x01,
+ 0x28, 0x09, 0x52, 0x02, 0x69, 0x70, 0x12, 0x10, 0x0a, 0x03, 0x70, 0x69, 0x64, 0x18, 0x05, 0x20,
+ 0x01, 0x28, 0x09, 0x52, 0x03, 0x70, 0x69, 0x64, 0x12, 0x10, 0x0a, 0x03, 0x73, 0x79, 0x73, 0x18,
+ 0x06, 0x20, 0x01, 0x28, 0x09, 0x52, 0x03, 0x73, 0x79, 0x73, 0x12, 0x1a, 0x0a, 0x08, 0x75, 0x73,
+ 0x65, 0x72, 0x6e, 0x61, 0x6d, 0x65, 0x18, 0x07, 0x20, 0x01, 0x28, 0x09, 0x52, 0x08, 0x75, 0x73,
+ 0x65, 0x72, 0x6e, 0x61, 0x6d, 0x65, 0x12, 0x1a, 0x0a, 0x08, 0x70, 0x61, 0x73, 0x73, 0x77, 0x6f,
+ 0x72, 0x64, 0x18, 0x08, 0x20, 0x01, 0x28, 0x09, 0x52, 0x08, 0x70, 0x61, 0x73, 0x73, 0x77, 0x6f,
+ 0x72, 0x64, 0x12, 0x1a, 0x0a, 0x08, 0x6c, 0x61, 0x6e, 0x67, 0x75, 0x61, 0x67, 0x65, 0x18, 0x09,
+ 0x20, 0x01, 0x28, 0x09, 0x52, 0x08, 0x6c, 0x61, 0x6e, 0x67, 0x75, 0x61, 0x67, 0x65, 0x12, 0x22,
+ 0x0a, 0x0c, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x63, 0x6f, 0x6c, 0x54, 0x79, 0x70, 0x65, 0x18, 0x0a,
+ 0x20, 0x01, 0x28, 0x09, 0x52, 0x0c, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x63, 0x6f, 0x6c, 0x54, 0x79,
+ 0x70, 0x65, 0x12, 0x28, 0x0a, 0x0f, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x63, 0x6f, 0x6c, 0x56, 0x65,
+ 0x72, 0x73, 0x69, 0x6f, 0x6e, 0x18, 0x0b, 0x20, 0x01, 0x28, 0x09, 0x52, 0x0f, 0x70, 0x72, 0x6f,
+ 0x74, 0x6f, 0x63, 0x6f, 0x6c, 0x56, 0x65, 0x72, 0x73, 0x69, 0x6f, 0x6e, 0x12, 0x22, 0x0a, 0x0c,
+ 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x63, 0x6f, 0x6c, 0x44, 0x65, 0x73, 0x63, 0x18, 0x0c, 0x20, 0x01,
+ 0x28, 0x09, 0x52, 0x0c, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x63, 0x6f, 0x6c, 0x44, 0x65, 0x73, 0x63,
+ 0x22, 0xa2, 0x03, 0x0a, 0x0d, 0x53, 0x69, 0x6d, 0x70, 0x6c, 0x65, 0x4d, 0x65, 0x73, 0x73, 0x61,
+ 0x67, 0x65, 0x12, 0x45, 0x0a, 0x06, 0x68, 0x65, 0x61, 0x64, 0x65, 0x72, 0x18, 0x01, 0x20, 0x01,
+ 0x28, 0x0b, 0x32, 0x2d, 0x2e, 0x65, 0x76, 0x65, 0x6e, 0x74, 0x6d, 0x65, 0x73, 0x68, 0x2e, 0x63,
+ 0x6f, 0x6d, 0x6d, 0x6f, 0x6e, 0x2e, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x63, 0x6f, 0x6c, 0x2e, 0x67,
+ 0x72, 0x70, 0x63, 0x2e, 0x52, 0x65, 0x71, 0x75, 0x65, 0x73, 0x74, 0x48, 0x65, 0x61, 0x64, 0x65,
+ 0x72, 0x52, 0x06, 0x68, 0x65, 0x61, 0x64, 0x65, 0x72, 0x12, 0x24, 0x0a, 0x0d, 0x70, 0x72, 0x6f,
+ 0x64, 0x75, 0x63, 0x65, 0x72, 0x47, 0x72, 0x6f, 0x75, 0x70, 0x18, 0x02, 0x20, 0x01, 0x28, 0x09,
+ 0x52, 0x0d, 0x70, 0x72, 0x6f, 0x64, 0x75, 0x63, 0x65, 0x72, 0x47, 0x72, 0x6f, 0x75, 0x70, 0x12,
+ 0x14, 0x0a, 0x05, 0x74, 0x6f, 0x70, 0x69, 0x63, 0x18, 0x03, 0x20, 0x01, 0x28, 0x09, 0x52, 0x05,
+ 0x74, 0x6f, 0x70, 0x69, 0x63, 0x12, 0x18, 0x0a, 0x07, 0x63, 0x6f, 0x6e, 0x74, 0x65, 0x6e, 0x74,
+ 0x18, 0x04, 0x20, 0x01, 0x28, 0x09, 0x52, 0x07, 0x63, 0x6f, 0x6e, 0x74, 0x65, 0x6e, 0x74, 0x12,
+ 0x10, 0x0a, 0x03, 0x74, 0x74, 0x6c, 0x18, 0x05, 0x20, 0x01, 0x28, 0x09, 0x52, 0x03, 0x74, 0x74,
+ 0x6c, 0x12, 0x1a, 0x0a, 0x08, 0x75, 0x6e, 0x69, 0x71, 0x75, 0x65, 0x49, 0x64, 0x18, 0x06, 0x20,
+ 0x01, 0x28, 0x09, 0x52, 0x08, 0x75, 0x6e, 0x69, 0x71, 0x75, 0x65, 0x49, 0x64, 0x12, 0x16, 0x0a,
+ 0x06, 0x73, 0x65, 0x71, 0x4e, 0x75, 0x6d, 0x18, 0x07, 0x20, 0x01, 0x28, 0x09, 0x52, 0x06, 0x73,
+ 0x65, 0x71, 0x4e, 0x75, 0x6d, 0x12, 0x10, 0x0a, 0x03, 0x74, 0x61, 0x67, 0x18, 0x08, 0x20, 0x01,
+ 0x28, 0x09, 0x52, 0x03, 0x74, 0x61, 0x67, 0x12, 0x5d, 0x0a, 0x0a, 0x70, 0x72, 0x6f, 0x70, 0x65,
+ 0x72, 0x74, 0x69, 0x65, 0x73, 0x18, 0x09, 0x20, 0x03, 0x28, 0x0b, 0x32, 0x3d, 0x2e, 0x65, 0x76,
+ 0x65, 0x6e, 0x74, 0x6d, 0x65, 0x73, 0x68, 0x2e, 0x63, 0x6f, 0x6d, 0x6d, 0x6f, 0x6e, 0x2e, 0x70,
+ 0x72, 0x6f, 0x74, 0x6f, 0x63, 0x6f, 0x6c, 0x2e, 0x67, 0x72, 0x70, 0x63, 0x2e, 0x53, 0x69, 0x6d,
+ 0x70, 0x6c, 0x65, 0x4d, 0x65, 0x73, 0x73, 0x61, 0x67, 0x65, 0x2e, 0x50, 0x72, 0x6f, 0x70, 0x65,
+ 0x72, 0x74, 0x69, 0x65, 0x73, 0x45, 0x6e, 0x74, 0x72, 0x79, 0x52, 0x0a, 0x70, 0x72, 0x6f, 0x70,
+ 0x65, 0x72, 0x74, 0x69, 0x65, 0x73, 0x1a, 0x3d, 0x0a, 0x0f, 0x50, 0x72, 0x6f, 0x70, 0x65, 0x72,
+ 0x74, 0x69, 0x65, 0x73, 0x45, 0x6e, 0x74, 0x72, 0x79, 0x12, 0x10, 0x0a, 0x03, 0x6b, 0x65, 0x79,
+ 0x18, 0x01, 0x20, 0x01, 0x28, 0x09, 0x52, 0x03, 0x6b, 0x65, 0x79, 0x12, 0x14, 0x0a, 0x05, 0x76,
+ 0x61, 0x6c, 0x75, 0x65, 0x18, 0x02, 0x20, 0x01, 0x28, 0x09, 0x52, 0x05, 0x76, 0x61, 0x6c, 0x75,
+ 0x65, 0x3a, 0x02, 0x38, 0x01, 0x22, 0x98, 0x04, 0x0a, 0x0c, 0x42, 0x61, 0x74, 0x63, 0x68, 0x4d,
+ 0x65, 0x73, 0x73, 0x61, 0x67, 0x65, 0x12, 0x45, 0x0a, 0x06, 0x68, 0x65, 0x61, 0x64, 0x65, 0x72,
+ 0x18, 0x01, 0x20, 0x01, 0x28, 0x0b, 0x32, 0x2d, 0x2e, 0x65, 0x76, 0x65, 0x6e, 0x74, 0x6d, 0x65,
+ 0x73, 0x68, 0x2e, 0x63, 0x6f, 0x6d, 0x6d, 0x6f, 0x6e, 0x2e, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x63,
+ 0x6f, 0x6c, 0x2e, 0x67, 0x72, 0x70, 0x63, 0x2e, 0x52, 0x65, 0x71, 0x75, 0x65, 0x73, 0x74, 0x48,
+ 0x65, 0x61, 0x64, 0x65, 0x72, 0x52, 0x06, 0x68, 0x65, 0x61, 0x64, 0x65, 0x72, 0x12, 0x24, 0x0a,
+ 0x0d, 0x70, 0x72, 0x6f, 0x64, 0x75, 0x63, 0x65, 0x72, 0x47, 0x72, 0x6f, 0x75, 0x70, 0x18, 0x02,
+ 0x20, 0x01, 0x28, 0x09, 0x52, 0x0d, 0x70, 0x72, 0x6f, 0x64, 0x75, 0x63, 0x65, 0x72, 0x47, 0x72,
+ 0x6f, 0x75, 0x70, 0x12, 0x14, 0x0a, 0x05, 0x74, 0x6f, 0x70, 0x69, 0x63, 0x18, 0x03, 0x20, 0x01,
+ 0x28, 0x09, 0x52, 0x05, 0x74, 0x6f, 0x70, 0x69, 0x63, 0x12, 0x5a, 0x0a, 0x0b, 0x6d, 0x65, 0x73,
+ 0x73, 0x61, 0x67, 0x65, 0x49, 0x74, 0x65, 0x6d, 0x18, 0x04, 0x20, 0x03, 0x28, 0x0b, 0x32, 0x38,
+ 0x2e, 0x65, 0x76, 0x65, 0x6e, 0x74, 0x6d, 0x65, 0x73, 0x68, 0x2e, 0x63, 0x6f, 0x6d, 0x6d, 0x6f,
+ 0x6e, 0x2e, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x63, 0x6f, 0x6c, 0x2e, 0x67, 0x72, 0x70, 0x63, 0x2e,
+ 0x42, 0x61, 0x74, 0x63, 0x68, 0x4d, 0x65, 0x73, 0x73, 0x61, 0x67, 0x65, 0x2e, 0x4d, 0x65, 0x73,
+ 0x73, 0x61, 0x67, 0x65, 0x49, 0x74, 0x65, 0x6d, 0x52, 0x0b, 0x6d, 0x65, 0x73, 0x73, 0x61, 0x67,
+ 0x65, 0x49, 0x74, 0x65, 0x6d, 0x1a, 0xa8, 0x02, 0x0a, 0x0b, 0x4d, 0x65, 0x73, 0x73, 0x61, 0x67,
+ 0x65, 0x49, 0x74, 0x65, 0x6d, 0x12, 0x18, 0x0a, 0x07, 0x63, 0x6f, 0x6e, 0x74, 0x65, 0x6e, 0x74,
+ 0x18, 0x01, 0x20, 0x01, 0x28, 0x09, 0x52, 0x07, 0x63, 0x6f, 0x6e, 0x74, 0x65, 0x6e, 0x74, 0x12,
+ 0x10, 0x0a, 0x03, 0x74, 0x74, 0x6c, 0x18, 0x02, 0x20, 0x01, 0x28, 0x09, 0x52, 0x03, 0x74, 0x74,
+ 0x6c, 0x12, 0x1a, 0x0a, 0x08, 0x75, 0x6e, 0x69, 0x71, 0x75, 0x65, 0x49, 0x64, 0x18, 0x03, 0x20,
+ 0x01, 0x28, 0x09, 0x52, 0x08, 0x75, 0x6e, 0x69, 0x71, 0x75, 0x65, 0x49, 0x64, 0x12, 0x16, 0x0a,
+ 0x06, 0x73, 0x65, 0x71, 0x4e, 0x75, 0x6d, 0x18, 0x04, 0x20, 0x01, 0x28, 0x09, 0x52, 0x06, 0x73,
+ 0x65, 0x71, 0x4e, 0x75, 0x6d, 0x12, 0x10, 0x0a, 0x03, 0x74, 0x61, 0x67, 0x18, 0x05, 0x20, 0x01,
+ 0x28, 0x09, 0x52, 0x03, 0x74, 0x61, 0x67, 0x12, 0x68, 0x0a, 0x0a, 0x70, 0x72, 0x6f, 0x70, 0x65,
+ 0x72, 0x74, 0x69, 0x65, 0x73, 0x18, 0x06, 0x20, 0x03, 0x28, 0x0b, 0x32, 0x48, 0x2e, 0x65, 0x76,
+ 0x65, 0x6e, 0x74, 0x6d, 0x65, 0x73, 0x68, 0x2e, 0x63, 0x6f, 0x6d, 0x6d, 0x6f, 0x6e, 0x2e, 0x70,
+ 0x72, 0x6f, 0x74, 0x6f, 0x63, 0x6f, 0x6c, 0x2e, 0x67, 0x72, 0x70, 0x63, 0x2e, 0x42, 0x61, 0x74,
+ 0x63, 0x68, 0x4d, 0x65, 0x73, 0x73, 0x61, 0x67, 0x65, 0x2e, 0x4d, 0x65, 0x73, 0x73, 0x61, 0x67,
+ 0x65, 0x49, 0x74, 0x65, 0x6d, 0x2e, 0x50, 0x72, 0x6f, 0x70, 0x65, 0x72, 0x74, 0x69, 0x65, 0x73,
+ 0x45, 0x6e, 0x74, 0x72, 0x79, 0x52, 0x0a, 0x70, 0x72, 0x6f, 0x70, 0x65, 0x72, 0x74, 0x69, 0x65,
+ 0x73, 0x1a, 0x3d, 0x0a, 0x0f, 0x50, 0x72, 0x6f, 0x70, 0x65, 0x72, 0x74, 0x69, 0x65, 0x73, 0x45,
+ 0x6e, 0x74, 0x72, 0x79, 0x12, 0x10, 0x0a, 0x03, 0x6b, 0x65, 0x79, 0x18, 0x01, 0x20, 0x01, 0x28,
+ 0x09, 0x52, 0x03, 0x6b, 0x65, 0x79, 0x12, 0x14, 0x0a, 0x05, 0x76, 0x61, 0x6c, 0x75, 0x65, 0x18,
+ 0x02, 0x20, 0x01, 0x28, 0x09, 0x52, 0x05, 0x76, 0x61, 0x6c, 0x75, 0x65, 0x3a, 0x02, 0x38, 0x01,
+ 0x22, 0x5c, 0x0a, 0x08, 0x52, 0x65, 0x73, 0x70, 0x6f, 0x6e, 0x73, 0x65, 0x12, 0x1a, 0x0a, 0x08,
+ 0x72, 0x65, 0x73, 0x70, 0x43, 0x6f, 0x64, 0x65, 0x18, 0x01, 0x20, 0x01, 0x28, 0x09, 0x52, 0x08,
+ 0x72, 0x65, 0x73, 0x70, 0x43, 0x6f, 0x64, 0x65, 0x12, 0x18, 0x0a, 0x07, 0x72, 0x65, 0x73, 0x70,
+ 0x4d, 0x73, 0x67, 0x18, 0x02, 0x20, 0x01, 0x28, 0x09, 0x52, 0x07, 0x72, 0x65, 0x73, 0x70, 0x4d,
+ 0x73, 0x67, 0x12, 0x1a, 0x0a, 0x08, 0x72, 0x65, 0x73, 0x70, 0x54, 0x69, 0x6d, 0x65, 0x18, 0x03,
+ 0x20, 0x01, 0x28, 0x09, 0x52, 0x08, 0x72, 0x65, 0x73, 0x70, 0x54, 0x69, 0x6d, 0x65, 0x22, 0xf1,
+ 0x07, 0x0a, 0x0c, 0x53, 0x75, 0x62, 0x73, 0x63, 0x72, 0x69, 0x70, 0x74, 0x69, 0x6f, 0x6e, 0x12,
+ 0x45, 0x0a, 0x06, 0x68, 0x65, 0x61, 0x64, 0x65, 0x72, 0x18, 0x01, 0x20, 0x01, 0x28, 0x0b, 0x32,
+ 0x2d, 0x2e, 0x65, 0x76, 0x65, 0x6e, 0x74, 0x6d, 0x65, 0x73, 0x68, 0x2e, 0x63, 0x6f, 0x6d, 0x6d,
+ 0x6f, 0x6e, 0x2e, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x63, 0x6f, 0x6c, 0x2e, 0x67, 0x72, 0x70, 0x63,
+ 0x2e, 0x52, 0x65, 0x71, 0x75, 0x65, 0x73, 0x74, 0x48, 0x65, 0x61, 0x64, 0x65, 0x72, 0x52, 0x06,
+ 0x68, 0x65, 0x61, 0x64, 0x65, 0x72, 0x12, 0x24, 0x0a, 0x0d, 0x63, 0x6f, 0x6e, 0x73, 0x75, 0x6d,
+ 0x65, 0x72, 0x47, 0x72, 0x6f, 0x75, 0x70, 0x18, 0x02, 0x20, 0x01, 0x28, 0x09, 0x52, 0x0d, 0x63,
+ 0x6f, 0x6e, 0x73, 0x75, 0x6d, 0x65, 0x72, 0x47, 0x72, 0x6f, 0x75, 0x70, 0x12, 0x6b, 0x0a, 0x11,
+ 0x73, 0x75, 0x62, 0x73, 0x63, 0x72, 0x69, 0x70, 0x74, 0x69, 0x6f, 0x6e, 0x49, 0x74, 0x65, 0x6d,
+ 0x73, 0x18, 0x03, 0x20, 0x03, 0x28, 0x0b, 0x32, 0x3d, 0x2e, 0x65, 0x76, 0x65, 0x6e, 0x74, 0x6d,
+ 0x65, 0x73, 0x68, 0x2e, 0x63, 0x6f, 0x6d, 0x6d, 0x6f, 0x6e, 0x2e, 0x70, 0x72, 0x6f, 0x74, 0x6f,
+ 0x63, 0x6f, 0x6c, 0x2e, 0x67, 0x72, 0x70, 0x63, 0x2e, 0x53, 0x75, 0x62, 0x73, 0x63, 0x72, 0x69,
+ 0x70, 0x74, 0x69, 0x6f, 0x6e, 0x2e, 0x53, 0x75, 0x62, 0x73, 0x63, 0x72, 0x69, 0x70, 0x74, 0x69,
+ 0x6f, 0x6e, 0x49, 0x74, 0x65, 0x6d, 0x52, 0x11, 0x73, 0x75, 0x62, 0x73, 0x63, 0x72, 0x69, 0x70,
+ 0x74, 0x69, 0x6f, 0x6e, 0x49, 0x74, 0x65, 0x6d, 0x73, 0x12, 0x10, 0x0a, 0x03, 0x75, 0x72, 0x6c,
+ 0x18, 0x04, 0x20, 0x01, 0x28, 0x09, 0x52, 0x03, 0x75, 0x72, 0x6c, 0x12, 0x48, 0x0a, 0x05, 0x72,
+ 0x65, 0x70, 0x6c, 0x79, 0x18, 0x05, 0x20, 0x01, 0x28, 0x0b, 0x32, 0x32, 0x2e, 0x65, 0x76, 0x65,
+ 0x6e, 0x74, 0x6d, 0x65, 0x73, 0x68, 0x2e, 0x63, 0x6f, 0x6d, 0x6d, 0x6f, 0x6e, 0x2e, 0x70, 0x72,
+ 0x6f, 0x74, 0x6f, 0x63, 0x6f, 0x6c, 0x2e, 0x67, 0x72, 0x70, 0x63, 0x2e, 0x53, 0x75, 0x62, 0x73,
+ 0x63, 0x72, 0x69, 0x70, 0x74, 0x69, 0x6f, 0x6e, 0x2e, 0x52, 0x65, 0x70, 0x6c, 0x79, 0x52, 0x05,
+ 0x72, 0x65, 0x70, 0x6c, 0x79, 0x1a, 0xcf, 0x02, 0x0a, 0x10, 0x53, 0x75, 0x62, 0x73, 0x63, 0x72,
+ 0x69, 0x70, 0x74, 0x69, 0x6f, 0x6e, 0x49, 0x74, 0x65, 0x6d, 0x12, 0x14, 0x0a, 0x05, 0x74, 0x6f,
+ 0x70, 0x69, 0x63, 0x18, 0x01, 0x20, 0x01, 0x28, 0x09, 0x52, 0x05, 0x74, 0x6f, 0x70, 0x69, 0x63,
+ 0x12, 0x62, 0x0a, 0x04, 0x6d, 0x6f, 0x64, 0x65, 0x18, 0x02, 0x20, 0x01, 0x28, 0x0e, 0x32, 0x4e,
+ 0x2e, 0x65, 0x76, 0x65, 0x6e, 0x74, 0x6d, 0x65, 0x73, 0x68, 0x2e, 0x63, 0x6f, 0x6d, 0x6d, 0x6f,
+ 0x6e, 0x2e, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x63, 0x6f, 0x6c, 0x2e, 0x67, 0x72, 0x70, 0x63, 0x2e,
+ 0x53, 0x75, 0x62, 0x73, 0x63, 0x72, 0x69, 0x70, 0x74, 0x69, 0x6f, 0x6e, 0x2e, 0x53, 0x75, 0x62,
+ 0x73, 0x63, 0x72, 0x69, 0x70, 0x74, 0x69, 0x6f, 0x6e, 0x49, 0x74, 0x65, 0x6d, 0x2e, 0x53, 0x75,
+ 0x62, 0x73, 0x63, 0x72, 0x69, 0x70, 0x74, 0x69, 0x6f, 0x6e, 0x4d, 0x6f, 0x64, 0x65, 0x52, 0x04,
+ 0x6d, 0x6f, 0x64, 0x65, 0x12, 0x62, 0x0a, 0x04, 0x74, 0x79, 0x70, 0x65, 0x18, 0x03, 0x20, 0x01,
+ 0x28, 0x0e, 0x32, 0x4e, 0x2e, 0x65, 0x76, 0x65, 0x6e, 0x74, 0x6d, 0x65, 0x73, 0x68, 0x2e, 0x63,
+ 0x6f, 0x6d, 0x6d, 0x6f, 0x6e, 0x2e, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x63, 0x6f, 0x6c, 0x2e, 0x67,
+ 0x72, 0x70, 0x63, 0x2e, 0x53, 0x75, 0x62, 0x73, 0x63, 0x72, 0x69, 0x70, 0x74, 0x69, 0x6f, 0x6e,
+ 0x2e, 0x53, 0x75, 0x62, 0x73, 0x63, 0x72, 0x69, 0x70, 0x74, 0x69, 0x6f, 0x6e, 0x49, 0x74, 0x65,
+ 0x6d, 0x2e, 0x53, 0x75, 0x62, 0x73, 0x63, 0x72, 0x69, 0x70, 0x74, 0x69, 0x6f, 0x6e, 0x54, 0x79,
+ 0x70, 0x65, 0x52, 0x04, 0x74, 0x79, 0x70, 0x65, 0x22, 0x34, 0x0a, 0x10, 0x53, 0x75, 0x62, 0x73,
+ 0x63, 0x72, 0x69, 0x70, 0x74, 0x69, 0x6f, 0x6e, 0x4d, 0x6f, 0x64, 0x65, 0x12, 0x0e, 0x0a, 0x0a,
+ 0x43, 0x4c, 0x55, 0x53, 0x54, 0x45, 0x52, 0x49, 0x4e, 0x47, 0x10, 0x00, 0x12, 0x10, 0x0a, 0x0c,
+ 0x42, 0x52, 0x4f, 0x41, 0x44, 0x43, 0x41, 0x53, 0x54, 0x49, 0x4e, 0x47, 0x10, 0x01, 0x22, 0x27,
+ 0x0a, 0x10, 0x53, 0x75, 0x62, 0x73, 0x63, 0x72, 0x69, 0x70, 0x74, 0x69, 0x6f, 0x6e, 0x54, 0x79,
+ 0x70, 0x65, 0x12, 0x09, 0x0a, 0x05, 0x41, 0x53, 0x59, 0x4e, 0x43, 0x10, 0x00, 0x12, 0x08, 0x0a,
+ 0x04, 0x53, 0x59, 0x4e, 0x43, 0x10, 0x01, 0x1a, 0xd8, 0x02, 0x0a, 0x05, 0x52, 0x65, 0x70, 0x6c,
+ 0x79, 0x12, 0x24, 0x0a, 0x0d, 0x70, 0x72, 0x6f, 0x64, 0x75, 0x63, 0x65, 0x72, 0x47, 0x72, 0x6f,
+ 0x75, 0x70, 0x18, 0x01, 0x20, 0x01, 0x28, 0x09, 0x52, 0x0d, 0x70, 0x72, 0x6f, 0x64, 0x75, 0x63,
+ 0x65, 0x72, 0x47, 0x72, 0x6f, 0x75, 0x70, 0x12, 0x14, 0x0a, 0x05, 0x74, 0x6f, 0x70, 0x69, 0x63,
+ 0x18, 0x02, 0x20, 0x01, 0x28, 0x09, 0x52, 0x05, 0x74, 0x6f, 0x70, 0x69, 0x63, 0x12, 0x18, 0x0a,
+ 0x07, 0x63, 0x6f, 0x6e, 0x74, 0x65, 0x6e, 0x74, 0x18, 0x03, 0x20, 0x01, 0x28, 0x09, 0x52, 0x07,
+ 0x63, 0x6f, 0x6e, 0x74, 0x65, 0x6e, 0x74, 0x12, 0x10, 0x0a, 0x03, 0x74, 0x74, 0x6c, 0x18, 0x04,
+ 0x20, 0x01, 0x28, 0x09, 0x52, 0x03, 0x74, 0x74, 0x6c, 0x12, 0x1a, 0x0a, 0x08, 0x75, 0x6e, 0x69,
+ 0x71, 0x75, 0x65, 0x49, 0x64, 0x18, 0x05, 0x20, 0x01, 0x28, 0x09, 0x52, 0x08, 0x75, 0x6e, 0x69,
+ 0x71, 0x75, 0x65, 0x49, 0x64, 0x12, 0x16, 0x0a, 0x06, 0x73, 0x65, 0x71, 0x4e, 0x75, 0x6d, 0x18,
+ 0x06, 0x20, 0x01, 0x28, 0x09, 0x52, 0x06, 0x73, 0x65, 0x71, 0x4e, 0x75, 0x6d, 0x12, 0x10, 0x0a,
+ 0x03, 0x74, 0x61, 0x67, 0x18, 0x07, 0x20, 0x01, 0x28, 0x09, 0x52, 0x03, 0x74, 0x61, 0x67, 0x12,
+ 0x62, 0x0a, 0x0a, 0x70, 0x72, 0x6f, 0x70, 0x65, 0x72, 0x74, 0x69, 0x65, 0x73, 0x18, 0x08, 0x20,
+ 0x03, 0x28, 0x0b, 0x32, 0x42, 0x2e, 0x65, 0x76, 0x65, 0x6e, 0x74, 0x6d, 0x65, 0x73, 0x68, 0x2e,
+ 0x63, 0x6f, 0x6d, 0x6d, 0x6f, 0x6e, 0x2e, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x63, 0x6f, 0x6c, 0x2e,
+ 0x67, 0x72, 0x70, 0x63, 0x2e, 0x53, 0x75, 0x62, 0x73, 0x63, 0x72, 0x69, 0x70, 0x74, 0x69, 0x6f,
+ 0x6e, 0x2e, 0x52, 0x65, 0x70, 0x6c, 0x79, 0x2e, 0x50, 0x72, 0x6f, 0x70, 0x65, 0x72, 0x74, 0x69,
+ 0x65, 0x73, 0x45, 0x6e, 0x74, 0x72, 0x79, 0x52, 0x0a, 0x70, 0x72, 0x6f, 0x70, 0x65, 0x72, 0x74,
+ 0x69, 0x65, 0x73, 0x1a, 0x3d, 0x0a, 0x0f, 0x50, 0x72, 0x6f, 0x70, 0x65, 0x72, 0x74, 0x69, 0x65,
+ 0x73, 0x45, 0x6e, 0x74, 0x72, 0x79, 0x12, 0x10, 0x0a, 0x03, 0x6b, 0x65, 0x79, 0x18, 0x01, 0x20,
+ 0x01, 0x28, 0x09, 0x52, 0x03, 0x6b, 0x65, 0x79, 0x12, 0x14, 0x0a, 0x05, 0x76, 0x61, 0x6c, 0x75,
+ 0x65, 0x18, 0x02, 0x20, 0x01, 0x28, 0x09, 0x52, 0x05, 0x76, 0x61, 0x6c, 0x75, 0x65, 0x3a, 0x02,
+ 0x38, 0x01, 0x22, 0xae, 0x03, 0x0a, 0x09, 0x48, 0x65, 0x61, 0x72, 0x74, 0x62, 0x65, 0x61, 0x74,
+ 0x12, 0x45, 0x0a, 0x06, 0x68, 0x65, 0x61, 0x64, 0x65, 0x72, 0x18, 0x01, 0x20, 0x01, 0x28, 0x0b,
+ 0x32, 0x2d, 0x2e, 0x65, 0x76, 0x65, 0x6e, 0x74, 0x6d, 0x65, 0x73, 0x68, 0x2e, 0x63, 0x6f, 0x6d,
+ 0x6d, 0x6f, 0x6e, 0x2e, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x63, 0x6f, 0x6c, 0x2e, 0x67, 0x72, 0x70,
+ 0x63, 0x2e, 0x52, 0x65, 0x71, 0x75, 0x65, 0x73, 0x74, 0x48, 0x65, 0x61, 0x64, 0x65, 0x72, 0x52,
+ 0x06, 0x68, 0x65, 0x61, 0x64, 0x65, 0x72, 0x12, 0x54, 0x0a, 0x0a, 0x63, 0x6c, 0x69, 0x65, 0x6e,
+ 0x74, 0x54, 0x79, 0x70, 0x65, 0x18, 0x02, 0x20, 0x01, 0x28, 0x0e, 0x32, 0x34, 0x2e, 0x65, 0x76,
+ 0x65, 0x6e, 0x74, 0x6d, 0x65, 0x73, 0x68, 0x2e, 0x63, 0x6f, 0x6d, 0x6d, 0x6f, 0x6e, 0x2e, 0x70,
+ 0x72, 0x6f, 0x74, 0x6f, 0x63, 0x6f, 0x6c, 0x2e, 0x67, 0x72, 0x70, 0x63, 0x2e, 0x48, 0x65, 0x61,
+ 0x72, 0x74, 0x62, 0x65, 0x61, 0x74, 0x2e, 0x43, 0x6c, 0x69, 0x65, 0x6e, 0x74, 0x54, 0x79, 0x70,
+ 0x65, 0x52, 0x0a, 0x63, 0x6c, 0x69, 0x65, 0x6e, 0x74, 0x54, 0x79, 0x70, 0x65, 0x12, 0x24, 0x0a,
+ 0x0d, 0x70, 0x72, 0x6f, 0x64, 0x75, 0x63, 0x65, 0x72, 0x47, 0x72, 0x6f, 0x75, 0x70, 0x18, 0x03,
+ 0x20, 0x01, 0x28, 0x09, 0x52, 0x0d, 0x70, 0x72, 0x6f, 0x64, 0x75, 0x63, 0x65, 0x72, 0x47, 0x72,
+ 0x6f, 0x75, 0x70, 0x12, 0x24, 0x0a, 0x0d, 0x63, 0x6f, 0x6e, 0x73, 0x75, 0x6d, 0x65, 0x72, 0x47,
+ 0x72, 0x6f, 0x75, 0x70, 0x18, 0x04, 0x20, 0x01, 0x28, 0x09, 0x52, 0x0d, 0x63, 0x6f, 0x6e, 0x73,
+ 0x75, 0x6d, 0x65, 0x72, 0x47, 0x72, 0x6f, 0x75, 0x70, 0x12, 0x5f, 0x0a, 0x0e, 0x68, 0x65, 0x61,
+ 0x72, 0x74, 0x62, 0x65, 0x61, 0x74, 0x49, 0x74, 0x65, 0x6d, 0x73, 0x18, 0x05, 0x20, 0x03, 0x28,
+ 0x0b, 0x32, 0x37, 0x2e, 0x65, 0x76, 0x65, 0x6e, 0x74, 0x6d, 0x65, 0x73, 0x68, 0x2e, 0x63, 0x6f,
+ 0x6d, 0x6d, 0x6f, 0x6e, 0x2e, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x63, 0x6f, 0x6c, 0x2e, 0x67, 0x72,
+ 0x70, 0x63, 0x2e, 0x48, 0x65, 0x61, 0x72, 0x74, 0x62, 0x65, 0x61, 0x74, 0x2e, 0x48, 0x65, 0x61,
+ 0x72, 0x74, 0x62, 0x65, 0x61, 0x74, 0x49, 0x74, 0x65, 0x6d, 0x52, 0x0e, 0x68, 0x65, 0x61, 0x72,
+ 0x74, 0x62, 0x65, 0x61, 0x74, 0x49, 0x74, 0x65, 0x6d, 0x73, 0x1a, 0x37, 0x0a, 0x0d, 0x48, 0x65,
+ 0x61, 0x72, 0x74, 0x62, 0x65, 0x61, 0x74, 0x49, 0x74, 0x65, 0x6d, 0x12, 0x14, 0x0a, 0x05, 0x74,
+ 0x6f, 0x70, 0x69, 0x63, 0x18, 0x01, 0x20, 0x01, 0x28, 0x09, 0x52, 0x05, 0x74, 0x6f, 0x70, 0x69,
+ 0x63, 0x12, 0x10, 0x0a, 0x03, 0x75, 0x72, 0x6c, 0x18, 0x02, 0x20, 0x01, 0x28, 0x09, 0x52, 0x03,
+ 0x75, 0x72, 0x6c, 0x22, 0x1e, 0x0a, 0x0a, 0x43, 0x6c, 0x69, 0x65, 0x6e, 0x74, 0x54, 0x79, 0x70,
+ 0x65, 0x12, 0x07, 0x0a, 0x03, 0x50, 0x55, 0x42, 0x10, 0x00, 0x12, 0x07, 0x0a, 0x03, 0x53, 0x55,
+ 0x42, 0x10, 0x01, 0x32, 0xcc, 0x02, 0x0a, 0x10, 0x50, 0x75, 0x62, 0x6c, 0x69, 0x73, 0x68, 0x65,
+ 0x72, 0x53, 0x65, 0x72, 0x76, 0x69, 0x63, 0x65, 0x12, 0x62, 0x0a, 0x07, 0x70, 0x75, 0x62, 0x6c,
+ 0x69, 0x73, 0x68, 0x12, 0x2d, 0x2e, 0x65, 0x76, 0x65, 0x6e, 0x74, 0x6d, 0x65, 0x73, 0x68, 0x2e,
+ 0x63, 0x6f, 0x6d, 0x6d, 0x6f, 0x6e, 0x2e, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x63, 0x6f, 0x6c, 0x2e,
+ 0x67, 0x72, 0x70, 0x63, 0x2e, 0x53, 0x69, 0x6d, 0x70, 0x6c, 0x65, 0x4d, 0x65, 0x73, 0x73, 0x61,
+ 0x67, 0x65, 0x1a, 0x28, 0x2e, 0x65, 0x76, 0x65, 0x6e, 0x74, 0x6d, 0x65, 0x73, 0x68, 0x2e, 0x63,
+ 0x6f, 0x6d, 0x6d, 0x6f, 0x6e, 0x2e, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x63, 0x6f, 0x6c, 0x2e, 0x67,
+ 0x72, 0x70, 0x63, 0x2e, 0x52, 0x65, 0x73, 0x70, 0x6f, 0x6e, 0x73, 0x65, 0x12, 0x6c, 0x0a, 0x0c,
+ 0x72, 0x65, 0x71, 0x75, 0x65, 0x73, 0x74, 0x52, 0x65, 0x70, 0x6c, 0x79, 0x12, 0x2d, 0x2e, 0x65,
+ 0x76, 0x65, 0x6e, 0x74, 0x6d, 0x65, 0x73, 0x68, 0x2e, 0x63, 0x6f, 0x6d, 0x6d, 0x6f, 0x6e, 0x2e,
+ 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x63, 0x6f, 0x6c, 0x2e, 0x67, 0x72, 0x70, 0x63, 0x2e, 0x53, 0x69,
+ 0x6d, 0x70, 0x6c, 0x65, 0x4d, 0x65, 0x73, 0x73, 0x61, 0x67, 0x65, 0x1a, 0x2d, 0x2e, 0x65, 0x76,
+ 0x65, 0x6e, 0x74, 0x6d, 0x65, 0x73, 0x68, 0x2e, 0x63, 0x6f, 0x6d, 0x6d, 0x6f, 0x6e, 0x2e, 0x70,
+ 0x72, 0x6f, 0x74, 0x6f, 0x63, 0x6f, 0x6c, 0x2e, 0x67, 0x72, 0x70, 0x63, 0x2e, 0x53, 0x69, 0x6d,
+ 0x70, 0x6c, 0x65, 0x4d, 0x65, 0x73, 0x73, 0x61, 0x67, 0x65, 0x12, 0x66, 0x0a, 0x0c, 0x62, 0x61,
+ 0x74, 0x63, 0x68, 0x50, 0x75, 0x62, 0x6c, 0x69, 0x73, 0x68, 0x12, 0x2c, 0x2e, 0x65, 0x76, 0x65,
+ 0x6e, 0x74, 0x6d, 0x65, 0x73, 0x68, 0x2e, 0x63, 0x6f, 0x6d, 0x6d, 0x6f, 0x6e, 0x2e, 0x70, 0x72,
+ 0x6f, 0x74, 0x6f, 0x63, 0x6f, 0x6c, 0x2e, 0x67, 0x72, 0x70, 0x63, 0x2e, 0x42, 0x61, 0x74, 0x63,
+ 0x68, 0x4d, 0x65, 0x73, 0x73, 0x61, 0x67, 0x65, 0x1a, 0x28, 0x2e, 0x65, 0x76, 0x65, 0x6e, 0x74,
+ 0x6d, 0x65, 0x73, 0x68, 0x2e, 0x63, 0x6f, 0x6d, 0x6d, 0x6f, 0x6e, 0x2e, 0x70, 0x72, 0x6f, 0x74,
+ 0x6f, 0x63, 0x6f, 0x6c, 0x2e, 0x67, 0x72, 0x70, 0x63, 0x2e, 0x52, 0x65, 0x73, 0x70, 0x6f, 0x6e,
+ 0x73, 0x65, 0x32, 0xd1, 0x02, 0x0a, 0x0f, 0x43, 0x6f, 0x6e, 0x73, 0x75, 0x6d, 0x65, 0x72, 0x53,
+ 0x65, 0x72, 0x76, 0x69, 0x63, 0x65, 0x12, 0x63, 0x0a, 0x09, 0x73, 0x75, 0x62, 0x73, 0x63, 0x72,
+ 0x69, 0x62, 0x65, 0x12, 0x2c, 0x2e, 0x65, 0x76, 0x65, 0x6e, 0x74, 0x6d, 0x65, 0x73, 0x68, 0x2e,
+ 0x63, 0x6f, 0x6d, 0x6d, 0x6f, 0x6e, 0x2e, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x63, 0x6f, 0x6c, 0x2e,
+ 0x67, 0x72, 0x70, 0x63, 0x2e, 0x53, 0x75, 0x62, 0x73, 0x63, 0x72, 0x69, 0x70, 0x74, 0x69, 0x6f,
+ 0x6e, 0x1a, 0x28, 0x2e, 0x65, 0x76, 0x65, 0x6e, 0x74, 0x6d, 0x65, 0x73, 0x68, 0x2e, 0x63, 0x6f,
+ 0x6d, 0x6d, 0x6f, 0x6e, 0x2e, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x63, 0x6f, 0x6c, 0x2e, 0x67, 0x72,
+ 0x70, 0x63, 0x2e, 0x52, 0x65, 0x73, 0x70, 0x6f, 0x6e, 0x73, 0x65, 0x12, 0x72, 0x0a, 0x0f, 0x73,
+ 0x75, 0x62, 0x73, 0x63, 0x72, 0x69, 0x62, 0x65, 0x53, 0x74, 0x72, 0x65, 0x61, 0x6d, 0x12, 0x2c,
+ 0x2e, 0x65, 0x76, 0x65, 0x6e, 0x74, 0x6d, 0x65, 0x73, 0x68, 0x2e, 0x63, 0x6f, 0x6d, 0x6d, 0x6f,
+ 0x6e, 0x2e, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x63, 0x6f, 0x6c, 0x2e, 0x67, 0x72, 0x70, 0x63, 0x2e,
+ 0x53, 0x75, 0x62, 0x73, 0x63, 0x72, 0x69, 0x70, 0x74, 0x69, 0x6f, 0x6e, 0x1a, 0x2d, 0x2e, 0x65,
+ 0x76, 0x65, 0x6e, 0x74, 0x6d, 0x65, 0x73, 0x68, 0x2e, 0x63, 0x6f, 0x6d, 0x6d, 0x6f, 0x6e, 0x2e,
+ 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x63, 0x6f, 0x6c, 0x2e, 0x67, 0x72, 0x70, 0x63, 0x2e, 0x53, 0x69,
+ 0x6d, 0x70, 0x6c, 0x65, 0x4d, 0x65, 0x73, 0x73, 0x61, 0x67, 0x65, 0x28, 0x01, 0x30, 0x01, 0x12,
+ 0x65, 0x0a, 0x0b, 0x75, 0x6e, 0x73, 0x75, 0x62, 0x73, 0x63, 0x72, 0x69, 0x62, 0x65, 0x12, 0x2c,
+ 0x2e, 0x65, 0x76, 0x65, 0x6e, 0x74, 0x6d, 0x65, 0x73, 0x68, 0x2e, 0x63, 0x6f, 0x6d, 0x6d, 0x6f,
+ 0x6e, 0x2e, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x63, 0x6f, 0x6c, 0x2e, 0x67, 0x72, 0x70, 0x63, 0x2e,
+ 0x53, 0x75, 0x62, 0x73, 0x63, 0x72, 0x69, 0x70, 0x74, 0x69, 0x6f, 0x6e, 0x1a, 0x28, 0x2e, 0x65,
+ 0x76, 0x65, 0x6e, 0x74, 0x6d, 0x65, 0x73, 0x68, 0x2e, 0x63, 0x6f, 0x6d, 0x6d, 0x6f, 0x6e, 0x2e,
+ 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x63, 0x6f, 0x6c, 0x2e, 0x67, 0x72, 0x70, 0x63, 0x2e, 0x52, 0x65,
+ 0x73, 0x70, 0x6f, 0x6e, 0x73, 0x65, 0x32, 0x74, 0x0a, 0x10, 0x48, 0x65, 0x61, 0x72, 0x74, 0x62,
+ 0x65, 0x61, 0x74, 0x53, 0x65, 0x72, 0x76, 0x69, 0x63, 0x65, 0x12, 0x60, 0x0a, 0x09, 0x68, 0x65,
+ 0x61, 0x72, 0x74, 0x62, 0x65, 0x61, 0x74, 0x12, 0x29, 0x2e, 0x65, 0x76, 0x65, 0x6e, 0x74, 0x6d,
+ 0x65, 0x73, 0x68, 0x2e, 0x63, 0x6f, 0x6d, 0x6d, 0x6f, 0x6e, 0x2e, 0x70, 0x72, 0x6f, 0x74, 0x6f,
+ 0x63, 0x6f, 0x6c, 0x2e, 0x67, 0x72, 0x70, 0x63, 0x2e, 0x48, 0x65, 0x61, 0x72, 0x74, 0x62, 0x65,
+ 0x61, 0x74, 0x1a, 0x28, 0x2e, 0x65, 0x76, 0x65, 0x6e, 0x74, 0x6d, 0x65, 0x73, 0x68, 0x2e, 0x63,
+ 0x6f, 0x6d, 0x6d, 0x6f, 0x6e, 0x2e, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x63, 0x6f, 0x6c, 0x2e, 0x67,
+ 0x72, 0x70, 0x63, 0x2e, 0x52, 0x65, 0x73, 0x70, 0x6f, 0x6e, 0x73, 0x65, 0x42, 0x09, 0x5a, 0x07,
+ 0x2e, 0x3b, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x62, 0x06, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x33,
+}
+
+var (
+ file_eventmesh_client_proto_rawDescOnce sync.Once
+ file_eventmesh_client_proto_rawDescData = file_eventmesh_client_proto_rawDesc
+)
+
+func file_eventmesh_client_proto_rawDescGZIP() []byte {
+ file_eventmesh_client_proto_rawDescOnce.Do(func() {
+ file_eventmesh_client_proto_rawDescData = protoimpl.X.CompressGZIP(file_eventmesh_client_proto_rawDescData)
+ })
+ return file_eventmesh_client_proto_rawDescData
+}
+
+var file_eventmesh_client_proto_enumTypes = make([]protoimpl.EnumInfo, 3)
+var file_eventmesh_client_proto_msgTypes = make([]protoimpl.MessageInfo, 13)
+var file_eventmesh_client_proto_goTypes = []interface{}{
+ (Subscription_SubscriptionItem_SubscriptionMode)(0), // 0: eventmesh.common.protocol.grpc.Subscription.SubscriptionItem.SubscriptionMode
+ (Subscription_SubscriptionItem_SubscriptionType)(0), // 1: eventmesh.common.protocol.grpc.Subscription.SubscriptionItem.SubscriptionType
+ (Heartbeat_ClientType)(0), // 2: eventmesh.common.protocol.grpc.Heartbeat.ClientType
+ (*RequestHeader)(nil), // 3: eventmesh.common.protocol.grpc.RequestHeader
+ (*SimpleMessage)(nil), // 4: eventmesh.common.protocol.grpc.SimpleMessage
+ (*BatchMessage)(nil), // 5: eventmesh.common.protocol.grpc.BatchMessage
+ (*Response)(nil), // 6: eventmesh.common.protocol.grpc.Response
+ (*Subscription)(nil), // 7: eventmesh.common.protocol.grpc.Subscription
+ (*Heartbeat)(nil), // 8: eventmesh.common.protocol.grpc.Heartbeat
+ nil, // 9: eventmesh.common.protocol.grpc.SimpleMessage.PropertiesEntry
+ (*BatchMessage_MessageItem)(nil), // 10: eventmesh.common.protocol.grpc.BatchMessage.MessageItem
+ nil, // 11: eventmesh.common.protocol.grpc.BatchMessage.MessageItem.PropertiesEntry
+ (*Subscription_SubscriptionItem)(nil), // 12: eventmesh.common.protocol.grpc.Subscription.SubscriptionItem
+ (*Subscription_Reply)(nil), // 13: eventmesh.common.protocol.grpc.Subscription.Reply
+ nil, // 14: eventmesh.common.protocol.grpc.Subscription.Reply.PropertiesEntry
+ (*Heartbeat_HeartbeatItem)(nil), // 15: eventmesh.common.protocol.grpc.Heartbeat.HeartbeatItem
+}
+var file_eventmesh_client_proto_depIdxs = []int32{
+ 3, // 0: eventmesh.common.protocol.grpc.SimpleMessage.header:type_name -> eventmesh.common.protocol.grpc.RequestHeader
+ 9, // 1: eventmesh.common.protocol.grpc.SimpleMessage.properties:type_name -> eventmesh.common.protocol.grpc.SimpleMessage.PropertiesEntry
+ 3, // 2: eventmesh.common.protocol.grpc.BatchMessage.header:type_name -> eventmesh.common.protocol.grpc.RequestHeader
+ 10, // 3: eventmesh.common.protocol.grpc.BatchMessage.messageItem:type_name -> eventmesh.common.protocol.grpc.BatchMessage.MessageItem
+ 3, // 4: eventmesh.common.protocol.grpc.Subscription.header:type_name -> eventmesh.common.protocol.grpc.RequestHeader
+ 12, // 5: eventmesh.common.protocol.grpc.Subscription.subscriptionItems:type_name -> eventmesh.common.protocol.grpc.Subscription.SubscriptionItem
+ 13, // 6: eventmesh.common.protocol.grpc.Subscription.reply:type_name -> eventmesh.common.protocol.grpc.Subscription.Reply
+ 3, // 7: eventmesh.common.protocol.grpc.Heartbeat.header:type_name -> eventmesh.common.protocol.grpc.RequestHeader
+ 2, // 8: eventmesh.common.protocol.grpc.Heartbeat.clientType:type_name -> eventmesh.common.protocol.grpc.Heartbeat.ClientType
+ 15, // 9: eventmesh.common.protocol.grpc.Heartbeat.heartbeatItems:type_name -> eventmesh.common.protocol.grpc.Heartbeat.HeartbeatItem
+ 11, // 10: eventmesh.common.protocol.grpc.BatchMessage.MessageItem.properties:type_name -> eventmesh.common.protocol.grpc.BatchMessage.MessageItem.PropertiesEntry
+ 0, // 11: eventmesh.common.protocol.grpc.Subscription.SubscriptionItem.mode:type_name -> eventmesh.common.protocol.grpc.Subscription.SubscriptionItem.SubscriptionMode
+ 1, // 12: eventmesh.common.protocol.grpc.Subscription.SubscriptionItem.type:type_name -> eventmesh.common.protocol.grpc.Subscription.SubscriptionItem.SubscriptionType
+ 14, // 13: eventmesh.common.protocol.grpc.Subscription.Reply.properties:type_name -> eventmesh.common.protocol.grpc.Subscription.Reply.PropertiesEntry
+ 4, // 14: eventmesh.common.protocol.grpc.PublisherService.publish:input_type -> eventmesh.common.protocol.grpc.SimpleMessage
+ 4, // 15: eventmesh.common.protocol.grpc.PublisherService.requestReply:input_type -> eventmesh.common.protocol.grpc.SimpleMessage
+ 5, // 16: eventmesh.common.protocol.grpc.PublisherService.batchPublish:input_type -> eventmesh.common.protocol.grpc.BatchMessage
+ 7, // 17: eventmesh.common.protocol.grpc.ConsumerService.subscribe:input_type -> eventmesh.common.protocol.grpc.Subscription
+ 7, // 18: eventmesh.common.protocol.grpc.ConsumerService.subscribeStream:input_type -> eventmesh.common.protocol.grpc.Subscription
+ 7, // 19: eventmesh.common.protocol.grpc.ConsumerService.unsubscribe:input_type -> eventmesh.common.protocol.grpc.Subscription
+ 8, // 20: eventmesh.common.protocol.grpc.HeartbeatService.heartbeat:input_type -> eventmesh.common.protocol.grpc.Heartbeat
+ 6, // 21: eventmesh.common.protocol.grpc.PublisherService.publish:output_type -> eventmesh.common.protocol.grpc.Response
+ 4, // 22: eventmesh.common.protocol.grpc.PublisherService.requestReply:output_type -> eventmesh.common.protocol.grpc.SimpleMessage
+ 6, // 23: eventmesh.common.protocol.grpc.PublisherService.batchPublish:output_type -> eventmesh.common.protocol.grpc.Response
+ 6, // 24: eventmesh.common.protocol.grpc.ConsumerService.subscribe:output_type -> eventmesh.common.protocol.grpc.Response
+ 4, // 25: eventmesh.common.protocol.grpc.ConsumerService.subscribeStream:output_type -> eventmesh.common.protocol.grpc.SimpleMessage
+ 6, // 26: eventmesh.common.protocol.grpc.ConsumerService.unsubscribe:output_type -> eventmesh.common.protocol.grpc.Response
+ 6, // 27: eventmesh.common.protocol.grpc.HeartbeatService.heartbeat:output_type -> eventmesh.common.protocol.grpc.Response
+ 21, // [21:28] is the sub-list for method output_type
+ 14, // [14:21] is the sub-list for method input_type
+ 14, // [14:14] is the sub-list for extension type_name
+ 14, // [14:14] is the sub-list for extension extendee
+ 0, // [0:14] is the sub-list for field type_name
+}
+
+func init() { file_eventmesh_client_proto_init() }
+func file_eventmesh_client_proto_init() {
+ if File_eventmesh_client_proto != nil {
+ return
+ }
+ if !protoimpl.UnsafeEnabled {
+ file_eventmesh_client_proto_msgTypes[0].Exporter = func(v interface{}, i int) interface{} {
+ switch v := v.(*RequestHeader); i {
+ case 0:
+ return &v.state
+ case 1:
+ return &v.sizeCache
+ case 2:
+ return &v.unknownFields
+ default:
+ return nil
+ }
+ }
+ file_eventmesh_client_proto_msgTypes[1].Exporter = func(v interface{}, i int) interface{} {
+ switch v := v.(*SimpleMessage); i {
+ case 0:
+ return &v.state
+ case 1:
+ return &v.sizeCache
+ case 2:
+ return &v.unknownFields
+ default:
+ return nil
+ }
+ }
+ file_eventmesh_client_proto_msgTypes[2].Exporter = func(v interface{}, i int) interface{} {
+ switch v := v.(*BatchMessage); i {
+ case 0:
+ return &v.state
+ case 1:
+ return &v.sizeCache
+ case 2:
+ return &v.unknownFields
+ default:
+ return nil
+ }
+ }
+ file_eventmesh_client_proto_msgTypes[3].Exporter = func(v interface{}, i int) interface{} {
+ switch v := v.(*Response); i {
+ case 0:
+ return &v.state
+ case 1:
+ return &v.sizeCache
+ case 2:
+ return &v.unknownFields
+ default:
+ return nil
+ }
+ }
+ file_eventmesh_client_proto_msgTypes[4].Exporter = func(v interface{}, i int) interface{} {
+ switch v := v.(*Subscription); i {
+ case 0:
+ return &v.state
+ case 1:
+ return &v.sizeCache
+ case 2:
+ return &v.unknownFields
+ default:
+ return nil
+ }
+ }
+ file_eventmesh_client_proto_msgTypes[5].Exporter = func(v interface{}, i int) interface{} {
+ switch v := v.(*Heartbeat); i {
+ case 0:
+ return &v.state
+ case 1:
+ return &v.sizeCache
+ case 2:
+ return &v.unknownFields
+ default:
+ return nil
+ }
+ }
+ file_eventmesh_client_proto_msgTypes[7].Exporter = func(v interface{}, i int) interface{} {
+ switch v := v.(*BatchMessage_MessageItem); i {
+ case 0:
+ return &v.state
+ case 1:
+ return &v.sizeCache
+ case 2:
+ return &v.unknownFields
+ default:
+ return nil
+ }
+ }
+ file_eventmesh_client_proto_msgTypes[9].Exporter = func(v interface{}, i int) interface{} {
+ switch v := v.(*Subscription_SubscriptionItem); i {
+ case 0:
+ return &v.state
+ case 1:
+ return &v.sizeCache
+ case 2:
+ return &v.unknownFields
+ default:
+ return nil
+ }
+ }
+ file_eventmesh_client_proto_msgTypes[10].Exporter = func(v interface{}, i int) interface{} {
+ switch v := v.(*Subscription_Reply); i {
+ case 0:
+ return &v.state
+ case 1:
+ return &v.sizeCache
+ case 2:
+ return &v.unknownFields
+ default:
+ return nil
+ }
+ }
+ file_eventmesh_client_proto_msgTypes[12].Exporter = func(v interface{}, i int) interface{} {
+ switch v := v.(*Heartbeat_HeartbeatItem); i {
+ case 0:
+ return &v.state
+ case 1:
+ return &v.sizeCache
+ case 2:
+ return &v.unknownFields
+ default:
+ return nil
+ }
+ }
+ }
+ type x struct{}
+ out := protoimpl.TypeBuilder{
+ File: protoimpl.DescBuilder{
+ GoPackagePath: reflect.TypeOf(x{}).PkgPath(),
+ RawDescriptor: file_eventmesh_client_proto_rawDesc,
+ NumEnums: 3,
+ NumMessages: 13,
+ NumExtensions: 0,
+ NumServices: 3,
+ },
+ GoTypes: file_eventmesh_client_proto_goTypes,
+ DependencyIndexes: file_eventmesh_client_proto_depIdxs,
+ EnumInfos: file_eventmesh_client_proto_enumTypes,
+ MessageInfos: file_eventmesh_client_proto_msgTypes,
+ }.Build()
+ File_eventmesh_client_proto = out.File
+ file_eventmesh_client_proto_rawDesc = nil
+ file_eventmesh_client_proto_goTypes = nil
+ file_eventmesh_client_proto_depIdxs = nil
+}
diff --git a/eventmesh-server-go/pkg/runtime/proto/pb/eventmesh-client.proto b/eventmesh-server-go/pkg/runtime/proto/pb/eventmesh-client.proto
new file mode 100644
index 00000000..349e8c1d
--- /dev/null
+++ b/eventmesh-server-go/pkg/runtime/proto/pb/eventmesh-client.proto
@@ -0,0 +1,165 @@
+/*
+ * 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.
+ */
+
+syntax = "proto3";
+
+//package eventmesh.common.protocol.grpc;
+//
+//option java_multiple_files = true;
+//option java_package = "org.apache.eventmesh.common.protocol.grpc.protos";
+//option java_outer_classname = "EventmeshGrpc";
+
+// make sure the protoc and protoc-gen-go is installed on your machine, and has set
+// its directory into path
+// download protoc: https://github.com/protocolbuffers/protobuf/releases
+// install protoc-gen-go: go install google.golang.org/protobuf/cmd/protoc-gen-go
+// generate go code by protoc: protoc --go_out=. eventmesh-client.proto
+
+package eventmesh.common.protocol.grpc;
+
+option go_package= ".;proto";
+
+message RequestHeader {
+ string env = 1;
+ string region = 2;
+ string idc = 3;
+ string ip = 4;
+ string pid = 5;
+ string sys = 6;
+ string username = 7;
+ string password = 8;
+ string language = 9;
+ string protocolType = 10;
+ string protocolVersion = 11;
+ string protocolDesc = 12;
+}
+
+message SimpleMessage {
+ RequestHeader header = 1;
+ string producerGroup = 2;
+ string topic = 3;
+ string content = 4;
+ string ttl = 5;
+ string uniqueId = 6;
+ string seqNum = 7;
+ string tag = 8;
+ map<string, string> properties = 9;
+}
+
+message BatchMessage {
+ RequestHeader header = 1;
+ string producerGroup = 2;
+ string topic = 3;
+
+ message MessageItem {
+ string content = 1;
+ string ttl = 2;
+ string uniqueId = 3;
+ string seqNum = 4;
+ string tag = 5;
+ map<string, string> properties = 6;
+ }
+
+ repeated MessageItem messageItem = 4;
+}
+
+message Response {
+ string respCode = 1;
+ string respMsg = 2;
+ string respTime = 3;
+}
+
+message Subscription {
+ RequestHeader header = 1;
+ string consumerGroup = 2;
+
+ message SubscriptionItem {
+ enum SubscriptionMode {
+ CLUSTERING = 0;
+ BROADCASTING = 1;
+ }
+
+ enum SubscriptionType {
+ ASYNC = 0;
+ SYNC = 1;
+ }
+
+ string topic = 1;
+ SubscriptionMode mode = 2;
+ SubscriptionType type = 3;
+ }
+
+ repeated SubscriptionItem subscriptionItems = 3;
+ string url = 4;
+
+ message Reply {
+ string producerGroup = 1;
+ string topic = 2;
+ string content = 3;
+ string ttl = 4;
+ string uniqueId = 5;
+ string seqNum = 6;
+ string tag = 7;
+ map<string, string> properties = 8;
+ }
+
+ Reply reply = 5;
+}
+
+message Heartbeat {
+ enum ClientType {
+ PUB = 0;
+ SUB = 1;
+ }
+
+ RequestHeader header = 1;
+ ClientType clientType = 2;
+ string producerGroup = 3;
+ string consumerGroup = 4;
+
+ message HeartbeatItem {
+ string topic = 1;
+ string url = 2;
+ }
+
+ repeated HeartbeatItem heartbeatItems = 5;
+}
+
+service PublisherService {
+ // Async event publish
+ rpc publish(SimpleMessage) returns (Response);
+
+ // Sync event publish
+ rpc requestReply(SimpleMessage) returns (SimpleMessage);
+
+ // Async batch event publish
+ rpc batchPublish(BatchMessage) returns (Response);
+}
+
+service ConsumerService {
+ // The subscribed event will be delivered by invoking the webhook url in the Subscription
+ rpc subscribe(Subscription) returns (Response);
+
+ // The subscribed event will be delivered through stream of Message
+ rpc subscribeStream(stream Subscription) returns (stream SimpleMessage);
+
+ rpc unsubscribe(Subscription) returns (Response);
+}
+
+service HeartbeatService {
+ rpc heartbeat(Heartbeat) returns (Response);
+}
\ No newline at end of file
diff --git a/eventmesh-server-go/pkg/runtime/proto/pb/eventmesh-client_grpc.pb.go b/eventmesh-server-go/pkg/runtime/proto/pb/eventmesh-client_grpc.pb.go
new file mode 100644
index 00000000..e96bf325
--- /dev/null
+++ b/eventmesh-server-go/pkg/runtime/proto/pb/eventmesh-client_grpc.pb.go
@@ -0,0 +1,479 @@
+// 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.
+
+// Code generated by protoc-gen-go-grpc. DO NOT EDIT.
+// versions:
+// - protoc-gen-go-grpc v1.2.0
+// - protoc v3.19.4
+// source: eventmesh-client.proto
+
+package pb
+
+import (
+ context "context"
+ grpc "google.golang.org/grpc"
+ codes "google.golang.org/grpc/codes"
+ status "google.golang.org/grpc/status"
+)
+
+// This is a compile-time assertion to ensure that this generated file
+// is compatible with the grpc package it is being compiled against.
+// Requires gRPC-Go v1.32.0 or later.
+const _ = grpc.SupportPackageIsVersion7
+
+// PublisherServiceClient is the client API for PublisherService service.
+//
+// For semantics around ctx use and closing/ending streaming RPCs, please refer to https://pkg.go.dev/google.golang.org/grpc/?tab=doc#ClientConn.NewStream.
+type PublisherServiceClient interface {
+ // Async event publish
+ Publish(ctx context.Context, in *SimpleMessage, opts ...grpc.CallOption) (*Response, error)
+ // Sync event publish
+ RequestReply(ctx context.Context, in *SimpleMessage, opts ...grpc.CallOption) (*SimpleMessage, error)
+ // Async batch event publish
+ BatchPublish(ctx context.Context, in *BatchMessage, opts ...grpc.CallOption) (*Response, error)
+}
+
+type publisherServiceClient struct {
+ cc grpc.ClientConnInterface
+}
+
+func NewPublisherServiceClient(cc grpc.ClientConnInterface) PublisherServiceClient {
+ return &publisherServiceClient{cc}
+}
+
+func (c *publisherServiceClient) Publish(ctx context.Context, in *SimpleMessage, opts ...grpc.CallOption) (*Response, error) {
+ out := new(Response)
+ err := c.cc.Invoke(ctx, "/eventmesh.common.protocol.grpc.PublisherService/publish", in, out, opts...)
+ if err != nil {
+ return nil, err
+ }
+ return out, nil
+}
+
+func (c *publisherServiceClient) RequestReply(ctx context.Context, in *SimpleMessage, opts ...grpc.CallOption) (*SimpleMessage, error) {
+ out := new(SimpleMessage)
+ err := c.cc.Invoke(ctx, "/eventmesh.common.protocol.grpc.PublisherService/requestReply", in, out, opts...)
+ if err != nil {
+ return nil, err
+ }
+ return out, nil
+}
+
+func (c *publisherServiceClient) BatchPublish(ctx context.Context, in *BatchMessage, opts ...grpc.CallOption) (*Response, error) {
+ out := new(Response)
+ err := c.cc.Invoke(ctx, "/eventmesh.common.protocol.grpc.PublisherService/batchPublish", in, out, opts...)
+ if err != nil {
+ return nil, err
+ }
+ return out, nil
+}
+
+// PublisherServiceServer is the server API for PublisherService service.
+// All implementations must embed UnimplementedPublisherServiceServer
+// for forward compatibility
+type PublisherServiceServer interface {
+ // Async event publish
+ Publish(context.Context, *SimpleMessage) (*Response, error)
+ // Sync event publish
+ RequestReply(context.Context, *SimpleMessage) (*SimpleMessage, error)
+ // Async batch event publish
+ BatchPublish(context.Context, *BatchMessage) (*Response, error)
+ mustEmbedUnimplementedPublisherServiceServer()
+}
+
+// UnimplementedPublisherServiceServer must be embedded to have forward compatible implementations.
+type UnimplementedPublisherServiceServer struct {
+}
+
+func (UnimplementedPublisherServiceServer) Publish(context.Context, *SimpleMessage) (*Response, error) {
+ return nil, status.Errorf(codes.Unimplemented, "method Publish not implemented")
+}
+func (UnimplementedPublisherServiceServer) RequestReply(context.Context, *SimpleMessage) (*SimpleMessage, error) {
+ return nil, status.Errorf(codes.Unimplemented, "method RequestReply not implemented")
+}
+func (UnimplementedPublisherServiceServer) BatchPublish(context.Context, *BatchMessage) (*Response, error) {
+ return nil, status.Errorf(codes.Unimplemented, "method BatchPublish not implemented")
+}
+func (UnimplementedPublisherServiceServer) mustEmbedUnimplementedPublisherServiceServer() {}
+
+// UnsafePublisherServiceServer may be embedded to opt out of forward compatibility for this service.
+// Use of this interface is not recommended, as added methods to PublisherServiceServer will
+// result in compilation errors.
+type UnsafePublisherServiceServer interface {
+ mustEmbedUnimplementedPublisherServiceServer()
+}
+
+func RegisterPublisherServiceServer(s grpc.ServiceRegistrar, srv PublisherServiceServer) {
+ s.RegisterService(&PublisherService_ServiceDesc, srv)
+}
+
+func _PublisherService_Publish_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) {
+ in := new(SimpleMessage)
+ if err := dec(in); err != nil {
+ return nil, err
+ }
+ if interceptor == nil {
+ return srv.(PublisherServiceServer).Publish(ctx, in)
+ }
+ info := &grpc.UnaryServerInfo{
+ Server: srv,
+ FullMethod: "/eventmesh.common.protocol.grpc.PublisherService/publish",
+ }
+ handler := func(ctx context.Context, req interface{}) (interface{}, error) {
+ return srv.(PublisherServiceServer).Publish(ctx, req.(*SimpleMessage))
+ }
+ return interceptor(ctx, in, info, handler)
+}
+
+func _PublisherService_RequestReply_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) {
+ in := new(SimpleMessage)
+ if err := dec(in); err != nil {
+ return nil, err
+ }
+ if interceptor == nil {
+ return srv.(PublisherServiceServer).RequestReply(ctx, in)
+ }
+ info := &grpc.UnaryServerInfo{
+ Server: srv,
+ FullMethod: "/eventmesh.common.protocol.grpc.PublisherService/requestReply",
+ }
+ handler := func(ctx context.Context, req interface{}) (interface{}, error) {
+ return srv.(PublisherServiceServer).RequestReply(ctx, req.(*SimpleMessage))
+ }
+ return interceptor(ctx, in, info, handler)
+}
+
+func _PublisherService_BatchPublish_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) {
+ in := new(BatchMessage)
+ if err := dec(in); err != nil {
+ return nil, err
+ }
+ if interceptor == nil {
+ return srv.(PublisherServiceServer).BatchPublish(ctx, in)
+ }
+ info := &grpc.UnaryServerInfo{
+ Server: srv,
+ FullMethod: "/eventmesh.common.protocol.grpc.PublisherService/batchPublish",
+ }
+ handler := func(ctx context.Context, req interface{}) (interface{}, error) {
+ return srv.(PublisherServiceServer).BatchPublish(ctx, req.(*BatchMessage))
+ }
+ return interceptor(ctx, in, info, handler)
+}
+
+// PublisherService_ServiceDesc is the grpc.ServiceDesc for PublisherService service.
+// It's only intended for direct use with grpc.RegisterService,
+// and not to be introspected or modified (even as a copy)
+var PublisherService_ServiceDesc = grpc.ServiceDesc{
+ ServiceName: "eventmesh.common.protocol.grpc.PublisherService",
+ HandlerType: (*PublisherServiceServer)(nil),
+ Methods: []grpc.MethodDesc{
+ {
+ MethodName: "publish",
+ Handler: _PublisherService_Publish_Handler,
+ },
+ {
+ MethodName: "requestReply",
+ Handler: _PublisherService_RequestReply_Handler,
+ },
+ {
+ MethodName: "batchPublish",
+ Handler: _PublisherService_BatchPublish_Handler,
+ },
+ },
+ Streams: []grpc.StreamDesc{},
+ Metadata: "eventmesh-client.proto",
+}
+
+// ConsumerServiceClient is the client API for ConsumerService service.
+//
+// For semantics around ctx use and closing/ending streaming RPCs, please refer to https://pkg.go.dev/google.golang.org/grpc/?tab=doc#ClientConn.NewStream.
+type ConsumerServiceClient interface {
+ // The subscribed event will be delivered by invoking the webhook url in the Subscription
+ Subscribe(ctx context.Context, in *Subscription, opts ...grpc.CallOption) (*Response, error)
+ // The subscribed event will be delivered through stream of Message
+ SubscribeStream(ctx context.Context, opts ...grpc.CallOption) (ConsumerService_SubscribeStreamClient, error)
+ Unsubscribe(ctx context.Context, in *Subscription, opts ...grpc.CallOption) (*Response, error)
+}
+
+type consumerServiceClient struct {
+ cc grpc.ClientConnInterface
+}
+
+func NewConsumerServiceClient(cc grpc.ClientConnInterface) ConsumerServiceClient {
+ return &consumerServiceClient{cc}
+}
+
+func (c *consumerServiceClient) Subscribe(ctx context.Context, in *Subscription, opts ...grpc.CallOption) (*Response, error) {
+ out := new(Response)
+ err := c.cc.Invoke(ctx, "/eventmesh.common.protocol.grpc.ConsumerService/subscribe", in, out, opts...)
+ if err != nil {
+ return nil, err
+ }
+ return out, nil
+}
+
+func (c *consumerServiceClient) SubscribeStream(ctx context.Context, opts ...grpc.CallOption) (ConsumerService_SubscribeStreamClient, error) {
+ stream, err := c.cc.NewStream(ctx, &ConsumerService_ServiceDesc.Streams[0], "/eventmesh.common.protocol.grpc.ConsumerService/subscribeStream", opts...)
+ if err != nil {
+ return nil, err
+ }
+ x := &consumerServiceSubscribeStreamClient{stream}
+ return x, nil
+}
+
+type ConsumerService_SubscribeStreamClient interface {
+ Send(*Subscription) error
+ Recv() (*SimpleMessage, error)
+ grpc.ClientStream
+}
+
+type consumerServiceSubscribeStreamClient struct {
+ grpc.ClientStream
+}
+
+func (x *consumerServiceSubscribeStreamClient) Send(m *Subscription) error {
+ return x.ClientStream.SendMsg(m)
+}
+
+func (x *consumerServiceSubscribeStreamClient) Recv() (*SimpleMessage, error) {
+ m := new(SimpleMessage)
+ if err := x.ClientStream.RecvMsg(m); err != nil {
+ return nil, err
+ }
+ return m, nil
+}
+
+func (c *consumerServiceClient) Unsubscribe(ctx context.Context, in *Subscription, opts ...grpc.CallOption) (*Response, error) {
+ out := new(Response)
+ err := c.cc.Invoke(ctx, "/eventmesh.common.protocol.grpc.ConsumerService/unsubscribe", in, out, opts...)
+ if err != nil {
+ return nil, err
+ }
+ return out, nil
+}
+
+// ConsumerServiceServer is the server API for ConsumerService service.
+// All implementations must embed UnimplementedConsumerServiceServer
+// for forward compatibility
+type ConsumerServiceServer interface {
+ // The subscribed event will be delivered by invoking the webhook url in the Subscription
+ Subscribe(context.Context, *Subscription) (*Response, error)
+ // The subscribed event will be delivered through stream of Message
+ SubscribeStream(ConsumerService_SubscribeStreamServer) error
+ Unsubscribe(context.Context, *Subscription) (*Response, error)
+ mustEmbedUnimplementedConsumerServiceServer()
+}
+
+// UnimplementedConsumerServiceServer must be embedded to have forward compatible implementations.
+type UnimplementedConsumerServiceServer struct {
+}
+
+func (UnimplementedConsumerServiceServer) Subscribe(context.Context, *Subscription) (*Response, error) {
+ return nil, status.Errorf(codes.Unimplemented, "method Subscribe not implemented")
+}
+func (UnimplementedConsumerServiceServer) SubscribeStream(ConsumerService_SubscribeStreamServer) error {
+ return status.Errorf(codes.Unimplemented, "method SubscribeStream not implemented")
+}
+func (UnimplementedConsumerServiceServer) Unsubscribe(context.Context, *Subscription) (*Response, error) {
+ return nil, status.Errorf(codes.Unimplemented, "method Unsubscribe not implemented")
+}
+func (UnimplementedConsumerServiceServer) mustEmbedUnimplementedConsumerServiceServer() {}
+
+// UnsafeConsumerServiceServer may be embedded to opt out of forward compatibility for this service.
+// Use of this interface is not recommended, as added methods to ConsumerServiceServer will
+// result in compilation errors.
+type UnsafeConsumerServiceServer interface {
+ mustEmbedUnimplementedConsumerServiceServer()
+}
+
+func RegisterConsumerServiceServer(s grpc.ServiceRegistrar, srv ConsumerServiceServer) {
+ s.RegisterService(&ConsumerService_ServiceDesc, srv)
+}
+
+func _ConsumerService_Subscribe_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) {
+ in := new(Subscription)
+ if err := dec(in); err != nil {
+ return nil, err
+ }
+ if interceptor == nil {
+ return srv.(ConsumerServiceServer).Subscribe(ctx, in)
+ }
+ info := &grpc.UnaryServerInfo{
+ Server: srv,
+ FullMethod: "/eventmesh.common.protocol.grpc.ConsumerService/subscribe",
+ }
+ handler := func(ctx context.Context, req interface{}) (interface{}, error) {
+ return srv.(ConsumerServiceServer).Subscribe(ctx, req.(*Subscription))
+ }
+ return interceptor(ctx, in, info, handler)
+}
+
+func _ConsumerService_SubscribeStream_Handler(srv interface{}, stream grpc.ServerStream) error {
+ return srv.(ConsumerServiceServer).SubscribeStream(&consumerServiceSubscribeStreamServer{stream})
+}
+
+type ConsumerService_SubscribeStreamServer interface {
+ Send(*SimpleMessage) error
+ Recv() (*Subscription, error)
+ grpc.ServerStream
+}
+
+type consumerServiceSubscribeStreamServer struct {
+ grpc.ServerStream
+}
+
+func (x *consumerServiceSubscribeStreamServer) Send(m *SimpleMessage) error {
+ return x.ServerStream.SendMsg(m)
+}
+
+func (x *consumerServiceSubscribeStreamServer) Recv() (*Subscription, error) {
+ m := new(Subscription)
+ if err := x.ServerStream.RecvMsg(m); err != nil {
+ return nil, err
+ }
+ return m, nil
+}
+
+func _ConsumerService_Unsubscribe_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) {
+ in := new(Subscription)
+ if err := dec(in); err != nil {
+ return nil, err
+ }
+ if interceptor == nil {
+ return srv.(ConsumerServiceServer).Unsubscribe(ctx, in)
+ }
+ info := &grpc.UnaryServerInfo{
+ Server: srv,
+ FullMethod: "/eventmesh.common.protocol.grpc.ConsumerService/unsubscribe",
+ }
+ handler := func(ctx context.Context, req interface{}) (interface{}, error) {
+ return srv.(ConsumerServiceServer).Unsubscribe(ctx, req.(*Subscription))
+ }
+ return interceptor(ctx, in, info, handler)
+}
+
+// ConsumerService_ServiceDesc is the grpc.ServiceDesc for ConsumerService service.
+// It's only intended for direct use with grpc.RegisterService,
+// and not to be introspected or modified (even as a copy)
+var ConsumerService_ServiceDesc = grpc.ServiceDesc{
+ ServiceName: "eventmesh.common.protocol.grpc.ConsumerService",
+ HandlerType: (*ConsumerServiceServer)(nil),
+ Methods: []grpc.MethodDesc{
+ {
+ MethodName: "subscribe",
+ Handler: _ConsumerService_Subscribe_Handler,
+ },
+ {
+ MethodName: "unsubscribe",
+ Handler: _ConsumerService_Unsubscribe_Handler,
+ },
+ },
+ Streams: []grpc.StreamDesc{
+ {
+ StreamName: "subscribeStream",
+ Handler: _ConsumerService_SubscribeStream_Handler,
+ ServerStreams: true,
+ ClientStreams: true,
+ },
+ },
+ Metadata: "eventmesh-client.proto",
+}
+
+// HeartbeatServiceClient is the client API for HeartbeatService service.
+//
+// For semantics around ctx use and closing/ending streaming RPCs, please refer to https://pkg.go.dev/google.golang.org/grpc/?tab=doc#ClientConn.NewStream.
+type HeartbeatServiceClient interface {
+ Heartbeat(ctx context.Context, in *Heartbeat, opts ...grpc.CallOption) (*Response, error)
+}
+
+type heartbeatServiceClient struct {
+ cc grpc.ClientConnInterface
+}
+
+func NewHeartbeatServiceClient(cc grpc.ClientConnInterface) HeartbeatServiceClient {
+ return &heartbeatServiceClient{cc}
+}
+
+func (c *heartbeatServiceClient) Heartbeat(ctx context.Context, in *Heartbeat, opts ...grpc.CallOption) (*Response, error) {
+ out := new(Response)
+ err := c.cc.Invoke(ctx, "/eventmesh.common.protocol.grpc.HeartbeatService/heartbeat", in, out, opts...)
+ if err != nil {
+ return nil, err
+ }
+ return out, nil
+}
+
+// HeartbeatServiceServer is the server API for HeartbeatService service.
+// All implementations must embed UnimplementedHeartbeatServiceServer
+// for forward compatibility
+type HeartbeatServiceServer interface {
+ Heartbeat(context.Context, *Heartbeat) (*Response, error)
+ mustEmbedUnimplementedHeartbeatServiceServer()
+}
+
+// UnimplementedHeartbeatServiceServer must be embedded to have forward compatible implementations.
+type UnimplementedHeartbeatServiceServer struct {
+}
+
+func (UnimplementedHeartbeatServiceServer) Heartbeat(context.Context, *Heartbeat) (*Response, error) {
+ return nil, status.Errorf(codes.Unimplemented, "method Heartbeat not implemented")
+}
+func (UnimplementedHeartbeatServiceServer) mustEmbedUnimplementedHeartbeatServiceServer() {}
+
+// UnsafeHeartbeatServiceServer may be embedded to opt out of forward compatibility for this service.
+// Use of this interface is not recommended, as added methods to HeartbeatServiceServer will
+// result in compilation errors.
+type UnsafeHeartbeatServiceServer interface {
+ mustEmbedUnimplementedHeartbeatServiceServer()
+}
+
+func RegisterHeartbeatServiceServer(s grpc.ServiceRegistrar, srv HeartbeatServiceServer) {
+ s.RegisterService(&HeartbeatService_ServiceDesc, srv)
+}
+
+func _HeartbeatService_Heartbeat_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) {
+ in := new(Heartbeat)
+ if err := dec(in); err != nil {
+ return nil, err
+ }
+ if interceptor == nil {
+ return srv.(HeartbeatServiceServer).Heartbeat(ctx, in)
+ }
+ info := &grpc.UnaryServerInfo{
+ Server: srv,
+ FullMethod: "/eventmesh.common.protocol.grpc.HeartbeatService/heartbeat",
+ }
+ handler := func(ctx context.Context, req interface{}) (interface{}, error) {
+ return srv.(HeartbeatServiceServer).Heartbeat(ctx, req.(*Heartbeat))
+ }
+ return interceptor(ctx, in, info, handler)
+}
+
+// HeartbeatService_ServiceDesc is the grpc.ServiceDesc for HeartbeatService service.
+// It's only intended for direct use with grpc.RegisterService,
+// and not to be introspected or modified (even as a copy)
+var HeartbeatService_ServiceDesc = grpc.ServiceDesc{
+ ServiceName: "eventmesh.common.protocol.grpc.HeartbeatService",
+ HandlerType: (*HeartbeatServiceServer)(nil),
+ Methods: []grpc.MethodDesc{
+ {
+ MethodName: "heartbeat",
+ Handler: _HeartbeatService_Heartbeat_Handler,
+ },
+ },
+ Streams: []grpc.StreamDesc{},
+ Metadata: "eventmesh-client.proto",
+}
---------------------------------------------------------------------
To unsubscribe, e-mail: commits-unsubscribe@eventmesh.apache.org
For additional commands, e-mail: commits-help@eventmesh.apache.org