Skip to content

Commit 4f1b863

Browse files
authored
refactor grpc broker interface (#193)
Signed-off-by: Wei Liu <liuweixa@redhat.com>
1 parent bfb55c4 commit 4f1b863

7 files changed

Lines changed: 85 additions & 165 deletions

File tree

pkg/cloudevents/clients/serviceaccount/integration_test.go

Lines changed: 16 additions & 39 deletions
Original file line numberDiff line numberDiff line change
@@ -3,7 +3,6 @@ package serviceaccount
33
import (
44
"context"
55
"net"
6-
"sync"
76
"testing"
87
"time"
98

@@ -21,36 +20,11 @@ import (
2120

2221
// tokenRequestService implements a service that handles token requests
2322
type tokenRequestService struct {
24-
mu sync.RWMutex
25-
tokenRequests map[string]*authenticationv1.TokenRequest
2623
handler server.EventHandler
2724
handleStatusUpdateFn func(context.Context, *cloudevents.Event) error
2825
}
2926

30-
func (s *tokenRequestService) Get(ctx context.Context, resourceID string) (*cloudevents.Event, error) {
31-
s.mu.RLock()
32-
tokenReq, ok := s.tokenRequests[resourceID]
33-
s.mu.RUnlock()
34-
if !ok {
35-
return nil, nil
36-
}
37-
38-
codec := NewTokenRequestCodec()
39-
eventType := cetypes.CloudEventsType{
40-
CloudEventsDataType: TokenRequestDataType,
41-
SubResource: cetypes.SubResourceStatus,
42-
Action: cetypes.UpdateRequestAction,
43-
}
44-
45-
evt, err := codec.Encode("test-source", eventType, tokenReq)
46-
if err != nil {
47-
return nil, err
48-
}
49-
50-
return evt, nil
51-
}
52-
53-
func (s *tokenRequestService) List(listOpts cetypes.ListOptions) ([]*cloudevents.Event, error) {
27+
func (s *tokenRequestService) List(_ context.Context, listOpts cetypes.ListOptions) ([]*cloudevents.Event, error) {
5428
return nil, nil
5529
}
5630

@@ -71,15 +45,23 @@ func (s *tokenRequestService) HandleStatusUpdate(ctx context.Context, evt *cloud
7145
ExpirationTimestamp: metav1.Time{Time: time.Now().Add(1 * time.Hour)},
7246
}
7347

74-
s.mu.Lock()
75-
s.tokenRequests[string(tokenReq.UID)] = tokenReq
76-
s.mu.Unlock()
77-
7848
// Notify the handler that we have a response ready
7949
if s.handler == nil {
8050
return nil
8151
}
82-
return s.handler.OnCreate(ctx, TokenRequestDataType, string(tokenReq.UID))
52+
53+
eventType := cetypes.CloudEventsType{
54+
CloudEventsDataType: TokenRequestDataType,
55+
SubResource: cetypes.SubResourceSpec,
56+
Action: cetypes.UpdateRequestAction,
57+
}
58+
59+
createEvt, err := codec.Encode("test-source", eventType, tokenReq)
60+
if err != nil {
61+
return err
62+
}
63+
64+
return s.handler.HandleEvent(ctx, createEvt)
8365
}
8466

8567
func (s *tokenRequestService) RegisterHandler(_ context.Context, handler server.EventHandler) {
@@ -95,9 +77,7 @@ func TestCreateToken_Integration(t *testing.T) {
9577
pbv1.RegisterCloudEventServiceServer(grpcServer, grpcEventServer)
9678

9779
// Create and register the token request service
98-
svc := &tokenRequestService{
99-
tokenRequests: make(map[string]*authenticationv1.TokenRequest),
100-
}
80+
svc := &tokenRequestService{}
10181
grpcEventServer.RegisterService(context.Background(), TokenRequestDataType, svc)
10282

10383
// Start listening
@@ -172,7 +152,6 @@ func TestCreateToken_Timeout(t *testing.T) {
172152

173153
// Create a service that receives but never responds
174154
svc := &tokenRequestService{
175-
tokenRequests: make(map[string]*authenticationv1.TokenRequest),
176155
handleStatusUpdateFn: func(ctx context.Context, evt *cloudevents.Event) error {
177156
// Just receive but don't respond
178157
return nil
@@ -231,9 +210,7 @@ func TestCreateToken_MultipleRequests(t *testing.T) {
231210
grpcEventServer := grpcserver.NewGRPCBroker(grpcserver.NewBrokerOptions())
232211
pbv1.RegisterCloudEventServiceServer(grpcServer, grpcEventServer)
233212

234-
svc := &tokenRequestService{
235-
tokenRequests: make(map[string]*authenticationv1.TokenRequest),
236-
}
213+
svc := &tokenRequestService{}
237214
grpcEventServer.RegisterService(context.Background(), TokenRequestDataType, svc)
238215

239216
lis, err := net.Listen("tcp", "127.0.0.1:0")

pkg/cloudevents/server/grpc/broker.go

Lines changed: 34 additions & 82 deletions
Original file line numberDiff line numberDiff line change
@@ -11,8 +11,6 @@ import (
1111
"open-cluster-management.io/sdk-go/pkg/cloudevents/server/grpc/heartbeat"
1212
"open-cluster-management.io/sdk-go/pkg/cloudevents/server/grpc/metrics"
1313

14-
"k8s.io/apimachinery/pkg/api/errors"
15-
1614
cloudevents "github.com/cloudevents/sdk-go/v2"
1715
"github.com/cloudevents/sdk-go/v2/binding"
1816
cloudeventstypes "github.com/cloudevents/sdk-go/v2/types"
@@ -315,30 +313,36 @@ func (bkr *GRPCBroker) respondResyncSpecRequest(ctx context.Context, eventDataTy
315313
return fmt.Errorf("failed to find service for event type %s", eventDataType)
316314
}
317315

318-
objs, err := service.List(types.ListOptions{ClusterName: clusterName, CloudEventsDataType: eventDataType})
316+
evts, err := service.List(ctx, types.ListOptions{ClusterName: clusterName, CloudEventsDataType: eventDataType})
319317
if err != nil {
320318
return err
321319
}
322320

323-
if len(objs) == 0 {
321+
if len(evts) == 0 {
324322
log.V(4).Info("no objs from the lister, do nothing")
325323
return nil
326324
}
327325

328-
for _, obj := range objs {
326+
for _, evt := range evts {
329327
// respond with the deleting resource regardless of the resource version
330-
objLogger := log.WithValues("eventType", obj.Type(), "extensions", obj.Extensions())
331-
if _, ok := obj.Extensions()[types.ExtensionDeletionTimestamp]; ok {
328+
objLogger := log.WithValues("eventType", evt.Type(), "extensions", evt.Extensions())
329+
if _, ok := evt.Extensions()[types.ExtensionDeletionTimestamp]; ok {
332330
objLogger.V(4).Info("respond spec resync request")
333-
err = bkr.handleRes(ctx, obj, eventDataType, "delete_request")
331+
deleteEventTypes := types.CloudEventsType{
332+
CloudEventsDataType: eventDataType,
333+
SubResource: types.SubResourceSpec,
334+
Action: types.DeleteRequestAction,
335+
}
336+
evt.SetType(deleteEventTypes.String())
337+
err = bkr.HandleEvent(ctx, evt)
334338
if err != nil {
335339
objLogger.Error(err, "failed to handle resync spec request")
336340
}
337341
continue
338342
}
339343

340-
lastResourceVersion := findResourceVersion(obj.ID(), resourceVersions.Versions)
341-
currentResourceVersion, err := cloudeventstypes.ToInteger(obj.Extensions()[types.ExtensionResourceVersion])
344+
lastResourceVersion := findResourceVersion(evt.ID(), resourceVersions.Versions)
345+
currentResourceVersion, err := cloudeventstypes.ToInteger(evt.Extensions()[types.ExtensionResourceVersion])
342346
if err != nil {
343347
objLogger.V(4).Info("ignore the event since it has a invalid resourceVersion", "error", err)
344348
continue
@@ -348,7 +352,13 @@ func (bkr *GRPCBroker) respondResyncSpecRequest(ctx context.Context, eventDataTy
348352
// the newer work to agent
349353
if currentResourceVersion == 0 || int64(currentResourceVersion) > lastResourceVersion {
350354
objLogger.V(4).Info("respond spec resync request")
351-
err := bkr.handleRes(ctx, obj, eventDataType, "update_request")
355+
updateEventTypes := types.CloudEventsType{
356+
CloudEventsDataType: eventDataType,
357+
SubResource: types.SubResourceSpec,
358+
Action: types.UpdateRequestAction,
359+
}
360+
evt.SetType(updateEventTypes.String())
361+
err := bkr.HandleEvent(ctx, evt)
352362
if err != nil {
353363
objLogger.Error(err, "failed to handle resync spec request")
354364
}
@@ -357,14 +367,15 @@ func (bkr *GRPCBroker) respondResyncSpecRequest(ctx context.Context, eventDataTy
357367

358368
// the resources do not exist on the source, but exist on the agent, delete them
359369
for _, rv := range resourceVersions.Versions {
360-
_, exists := getObj(rv.ResourceID, objs)
370+
_, exists := getObj(rv.ResourceID, evts)
361371
if exists {
362372
continue
363373
}
364374

365375
deleteEventTypes := types.CloudEventsType{
366376
CloudEventsDataType: eventDataType,
367377
SubResource: types.SubResourceSpec,
378+
Action: types.DeleteRequestAction,
368379
}
369380
obj := types.NewEventBuilder("source", deleteEventTypes).
370381
WithResourceID(rv.ResourceID).
@@ -375,7 +386,7 @@ func (bkr *GRPCBroker) respondResyncSpecRequest(ctx context.Context, eventDataTy
375386

376387
// send a delete event for the current resource
377388
log.V(4).Info("respond spec resync request")
378-
err := bkr.handleRes(ctx, &obj, eventDataType, "delete_request")
389+
err := bkr.HandleEvent(ctx, &obj)
379390
if err != nil {
380391
log.Error(err, "failed to handle delete request")
381392
}
@@ -384,22 +395,20 @@ func (bkr *GRPCBroker) respondResyncSpecRequest(ctx context.Context, eventDataTy
384395
return nil
385396
}
386397

387-
// handleRes publish the resource to the correct subscriber.
388-
func (bkr *GRPCBroker) handleRes(
389-
ctx context.Context,
390-
evt *cloudevents.Event,
391-
t types.CloudEventsDataType,
392-
action types.EventAction) error {
398+
// HandleEvent publish the event to the correct subscriber.
399+
func (bkr *GRPCBroker) HandleEvent(ctx context.Context, evt *cloudevents.Event) error {
400+
if evt == nil {
401+
return fmt.Errorf("event is nil")
402+
}
393403

394404
bkr.mu.RLock()
395405
defer bkr.mu.RUnlock()
396406

397-
eventType := types.CloudEventsType{
398-
CloudEventsDataType: t,
399-
SubResource: types.SubResourceSpec,
400-
Action: action,
407+
eventType, err := types.ParseCloudEventsType(evt.Type())
408+
if err != nil {
409+
return err
401410
}
402-
evt.SetType(eventType.String())
411+
evtDataType := eventType.CloudEventsDataType
403412

404413
clusterNameValue, err := evt.Context.GetExtension(types.ExtensionClusterName)
405414
if err != nil {
@@ -412,7 +421,7 @@ func (bkr *GRPCBroker) handleRes(
412421
// the resource consumer name and its data type is in the subscriber list, ensuring
413422
// the event will be only processed when the consumer is subscribed to the current
414423
// broker.
415-
if subscriber.clusterName == clusterName && subscriber.dataType == t {
424+
if subscriber.clusterName == clusterName && subscriber.dataType == evtDataType {
416425
if err := subscriber.handler(ctx, subID, evt); err != nil {
417426
return err
418427
}
@@ -421,63 +430,6 @@ func (bkr *GRPCBroker) handleRes(
421430
return nil
422431
}
423432

424-
// OnCreate is called by the controller when a resource is created on the maestro server.
425-
func (bkr *GRPCBroker) OnCreate(ctx context.Context, t types.CloudEventsDataType, id string) error {
426-
service, ok := bkr.services[t]
427-
if !ok {
428-
return fmt.Errorf("failed to find service for event type %s", t)
429-
}
430-
431-
resource, err := service.Get(ctx, id)
432-
// if the resource is not found, it indicates the resource has been processed.
433-
if errors.IsNotFound(err) {
434-
return nil
435-
}
436-
if err != nil {
437-
return err
438-
}
439-
440-
return bkr.handleRes(ctx, resource, t, "create_request")
441-
}
442-
443-
// OnUpdate is called by the controller when a resource is updated on the maestro server.
444-
func (bkr *GRPCBroker) OnUpdate(ctx context.Context, t types.CloudEventsDataType, id string) error {
445-
service, ok := bkr.services[t]
446-
if !ok {
447-
return fmt.Errorf("failed to find service for event type %s", t)
448-
}
449-
450-
resource, err := service.Get(ctx, id)
451-
// if the resource is not found, it indicates the resource has been processed.
452-
if errors.IsNotFound(err) {
453-
return nil
454-
}
455-
if err != nil {
456-
return err
457-
}
458-
459-
return bkr.handleRes(ctx, resource, t, "update_request")
460-
}
461-
462-
// OnDelete is called by the controller when a resource is deleted from the maestro server.
463-
func (bkr *GRPCBroker) OnDelete(ctx context.Context, t types.CloudEventsDataType, id string) error {
464-
service, ok := bkr.services[t]
465-
if !ok {
466-
return fmt.Errorf("failed to find service for event type %s", t)
467-
}
468-
469-
resource, err := service.Get(ctx, id)
470-
// if the resource is not found, it indicates the resource has been processed.
471-
if errors.IsNotFound(err) {
472-
return nil
473-
}
474-
if err != nil {
475-
return err
476-
}
477-
478-
return bkr.handleRes(ctx, resource, t, "delete_request")
479-
}
480-
481433
// IsConsumerSubscribed returns true if the consumer is subscribed to the broker for resource spec.
482434
func (bkr *GRPCBroker) IsConsumerSubscribed(consumerName string) bool {
483435
bkr.mu.RLock()

pkg/cloudevents/server/grpc/broker_test.go

Lines changed: 2 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -2,7 +2,6 @@ package grpc
22

33
import (
44
"context"
5-
"errors"
65
"net"
76
"testing"
87

@@ -26,16 +25,8 @@ type testService struct {
2625
handler server.EventHandler
2726
}
2827

29-
func (s *testService) Get(ctx context.Context, resourceID string) (*cloudevents.Event, error) {
30-
evt, ok := s.evts[resourceID]
31-
if !ok {
32-
return nil, errors.New("not found")
33-
}
34-
return evt, nil
35-
}
36-
3728
// List the cloudEvent from the service
38-
func (s *testService) List(listOpts cetypes.ListOptions) ([]*cloudevents.Event, error) {
29+
func (s *testService) List(_ context.Context, listOpts cetypes.ListOptions) ([]*cloudevents.Event, error) {
3930
evts := make([]*cloudevents.Event, 0, len(s.evts))
4031
for _, evt := range s.evts {
4132
evts = append(evts, evt)
@@ -56,7 +47,7 @@ func (s *testService) RegisterHandler(_ context.Context, handler server.EventHan
5647

5748
func (s *testService) create(evt *cloudevents.Event) error {
5849
s.evts[evt.ID()] = evt
59-
return s.handler.OnCreate(context.TODO(), dataType, evt.ID())
50+
return s.handler.HandleEvent(context.TODO(), evt)
6051
}
6152

6253
func TestServer(t *testing.T) {

pkg/cloudevents/server/interface.go

Lines changed: 3 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,7 @@ package server
33
import (
44
"context"
55

6+
cloudevents "github.com/cloudevents/sdk-go/v2"
67
"k8s.io/apimachinery/pkg/util/sets"
78

89
"open-cluster-management.io/sdk-go/pkg/cloudevents/generic/types"
@@ -22,14 +23,8 @@ type AgentEventServer interface {
2223
}
2324

2425
type EventHandler interface {
25-
// OnCreate is the callback when resource is created in the service.
26-
OnCreate(ctx context.Context, t types.CloudEventsDataType, resourceID string) error
27-
28-
// OnUpdate is the callback when resource is updated in the service.
29-
OnUpdate(ctx context.Context, t types.CloudEventsDataType, resourceID string) error
30-
31-
// OnDelete is the callback when resource is deleted from the service.
32-
OnDelete(ctx context.Context, t types.CloudEventsDataType, resourceID string) error
26+
// HandleEvent publish the event to the correct subscriber.
27+
HandleEvent(ctx context.Context, evt *cloudevents.Event) error
3328
}
3429

3530
// TODO SourceEventServer to handle the grpc conversation between consumers and grpcserver.

pkg/cloudevents/server/store.go

Lines changed: 2 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,7 @@ package server
22

33
import (
44
"context"
5+
56
cloudevents "github.com/cloudevents/sdk-go/v2"
67
cetypes "open-cluster-management.io/sdk-go/pkg/cloudevents/generic/types"
78
)
@@ -11,11 +12,8 @@ import (
1112

1213
// TODO need a method to check if an event has been processed already.
1314
type Service interface {
14-
// Get the cloudEvent based on resourceID from the service
15-
Get(ctx context.Context, resourceID string) (*cloudevents.Event, error)
16-
1715
// List the cloudEvent from the service
18-
List(listOpts cetypes.ListOptions) ([]*cloudevents.Event, error)
16+
List(ctx context.Context, listOpts cetypes.ListOptions) ([]*cloudevents.Event, error)
1917

2018
// HandleStatusUpdate processes the resource status update from the agent.
2119
HandleStatusUpdate(ctx context.Context, evt *cloudevents.Event) error

pkg/server/grpc/metrics/metrics_test.go

Lines changed: 1 addition & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -207,11 +207,7 @@ func newMockWorkService() *mockWorkService {
207207
return &mockWorkService{}
208208
}
209209

210-
func (s *mockWorkService) Get(ctx context.Context, resourceID string) (*cloudevents.Event, error) {
211-
return nil, nil
212-
}
213-
214-
func (s *mockWorkService) List(listOpts cetypes.ListOptions) ([]*cloudevents.Event, error) {
210+
func (s *mockWorkService) List(_ context.Context, listOpts cetypes.ListOptions) ([]*cloudevents.Event, error) {
215211
return nil, nil
216212
}
217213

0 commit comments

Comments
 (0)