Skip to content

Commit 502f5f8

Browse files
committed
Add grpc service register.
Signed-off-by: xuezhaojun <zxue@redhat.com>
1 parent 5ef72d1 commit 502f5f8

3 files changed

Lines changed: 36 additions & 35 deletions

File tree

pkg/cloudevents/server/grpc/options/server.go

Lines changed: 36 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@ import (
55
"crypto/tls"
66
"crypto/x509"
77
"fmt"
8+
"net"
89
"os"
910
"sync"
1011

@@ -29,15 +30,16 @@ type PreStartHook interface {
2930
}
3031

3132
type Server struct {
32-
options *GRPCServerOptions
33-
authenticators []authn.Authenticator
34-
authorizers []authz.Authorizer
35-
services map[types.CloudEventsDataType]server.Service
36-
hooks []PreStartHook
33+
options *GRPCServerOptions
34+
authenticators []authn.Authenticator
35+
authorizers []authz.Authorizer
36+
cloudeventServices map[types.CloudEventsDataType]server.Service
37+
grpcService []func(*grpc.Server)
38+
hooks []PreStartHook
3739
}
3840

3941
func NewServer(opt *GRPCServerOptions) *Server {
40-
return &Server{options: opt, services: make(map[types.CloudEventsDataType]server.Service)}
42+
return &Server{options: opt, cloudeventServices: make(map[types.CloudEventsDataType]server.Service)}
4143
}
4244

4345
func (s *Server) WithAuthenticator(authenticator authn.Authenticator) *Server {
@@ -51,7 +53,7 @@ func (s *Server) WithAuthorizer(authorizer authz.Authorizer) *Server {
5153
}
5254

5355
func (s *Server) WithService(t types.CloudEventsDataType, service server.Service) *Server {
54-
s.services[t] = service
56+
s.cloudeventServices[t] = service
5557
return s
5658
}
5759

@@ -60,6 +62,11 @@ func (s *Server) WithPreStartHooks(hooks ...PreStartHook) *Server {
6062
return s
6163
}
6264

65+
func (s *Server) WithGRPCService(fn func(*grpc.Server)) *Server {
66+
s.grpcService = append(s.grpcService, fn)
67+
return s
68+
}
69+
6370
func (s *Server) Run(ctx context.Context) error {
6471
var grpcServerOptions []grpc.ServerOption
6572
grpcServerOptions = append(grpcServerOptions, grpc.MaxRecvMsgSize(s.options.MaxReceiveMessageSize))
@@ -116,19 +123,38 @@ func (s *Server) Run(ctx context.Context) error {
116123
newAuthzStreamInterceptor(s.authorizers...)))
117124

118125
grpcServer := grpc.NewServer(grpcServerOptions...)
119-
grpcEventServer := grpcserver.NewGRPCBroker(grpcServer)
120126

121-
for t, service := range s.services {
127+
// Register the CloudEventServiceServer
128+
grpcEventServer := grpcserver.NewGRPCBroker(grpcServer)
129+
for t, service := range s.cloudeventServices {
122130
grpcEventServer.RegisterService(t, service)
123131
}
124132

133+
// Register other grpc services
134+
for _, fn := range s.grpcService {
135+
fn(grpcServer)
136+
}
137+
125138
// start hook
126139
for _, hook := range s.hooks {
127140
hook.Run(ctx)
128141
}
129142

130-
go grpcEventServer.Start(ctx, ":"+s.options.ServerBindPort)
143+
// start grpc server
144+
lis, err := net.Listen("tcp", ":"+s.options.ServerBindPort)
145+
if err != nil {
146+
return fmt.Errorf("failed to listen: %v", err)
147+
}
148+
go func() {
149+
if err := grpcServer.Serve(lis); err != nil {
150+
klog.Errorf("failed to serve gRPC broker: %v", err)
151+
}
152+
}()
153+
131154
<-ctx.Done()
155+
156+
klog.Infof("Shutting down gRPC server")
157+
grpcServer.GracefulStop()
132158
return nil
133159
}
134160

pkg/cloudevents/server/grpc/server.go

Lines changed: 0 additions & 22 deletions
Original file line numberDiff line numberDiff line change
@@ -4,12 +4,10 @@ import (
44
"context"
55
"encoding/json"
66
"fmt"
7-
"net"
87
"sync"
98
"time"
109

1110
"k8s.io/apimachinery/pkg/api/errors"
12-
utilruntime "k8s.io/apimachinery/pkg/util/runtime"
1311

1412
cloudevents "github.com/cloudevents/sdk-go/v2"
1513
"github.com/cloudevents/sdk-go/v2/binding"
@@ -78,26 +76,6 @@ func (bkr *GRPCBroker) Subscribers() sets.Set[string] {
7876
return subscribers
7977
}
8078

81-
// Start starts the gRPC broker at the given address
82-
func (bkr *GRPCBroker) Start(ctx context.Context, addr string) {
83-
logger := klog.FromContext(ctx)
84-
logger.Info("Starting gRPC broker at addr", "addr", addr)
85-
lis, err := net.Listen("tcp", addr)
86-
if err != nil {
87-
utilruntime.Must(fmt.Errorf("failed to listen: %v", err))
88-
}
89-
go func() {
90-
if err := bkr.grpcServer.Serve(lis); err != nil {
91-
utilruntime.Must(fmt.Errorf("failed to serve gRPC broker: %v", err))
92-
}
93-
}()
94-
95-
// wait until context is canceled
96-
<-ctx.Done()
97-
klog.Infof("Shutting down gRPC broker")
98-
bkr.grpcServer.GracefulStop()
99-
}
100-
10179
// Publish in stub implementation for agent publish resource status.
10280
func (bkr *GRPCBroker) Publish(ctx context.Context, pubReq *pbv1.PublishRequest) (*emptypb.Empty, error) {
10381
logger := klog.FromContext(ctx)

pkg/cloudevents/server/interface.go

Lines changed: 0 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -17,9 +17,6 @@ type AgentEventServer interface {
1717
// RegisterService registers a backend service with a certain data type.
1818
RegisterService(t types.CloudEventsDataType, service Service)
1919

20-
// Start initiates the EventServer to listen to agents.
21-
Start(ctx context.Context, addr string)
22-
2320
// Subscribers returns all current subscribers who subscribe to this server.
2421
Subscribers() sets.Set[string]
2522
}

0 commit comments

Comments
 (0)