Skip to content

Commit 9e9f97a

Browse files
authored
send the header immediately (#194)
Signed-off-by: Wei Liu <liuweixa@redhat.com>
1 parent 4f1b863 commit 9e9f97a

5 files changed

Lines changed: 316 additions & 53 deletions

File tree

pkg/cloudevents/generic/clients/baseclient.go

Lines changed: 17 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -210,10 +210,23 @@ func (c *baseClient) subscribe(ctx context.Context, receive receiveFn) {
210210
case <-ctx.Done():
211211
return
212212
case <-c.subscribeChan:
213-
if err := c.transport.Subscribe(ctx); err != nil {
214-
// Failed to send subscribe request, it should be connection failed, will retry on next reconnection
215-
runtime.HandleErrorWithContext(ctx, err, "failed to subscribe after connection")
216-
continue
213+
// Retry subscribe with backoff until success or context cancellation
214+
for {
215+
if err := c.transport.Subscribe(ctx); err != nil {
216+
runtime.HandleErrorWithContext(ctx, err, "failed to subscribe after connection")
217+
218+
// Wait with backoff before retrying
219+
select {
220+
case <-ctx.Done():
221+
return
222+
case <-wait.RealTimer(DelayFn()).C():
223+
// Continue to retry
224+
}
225+
continue
226+
}
227+
228+
// Subscribe succeeded, break out of retry loop
229+
break
217230
}
218231

219232
// Send startReceiverSignal to start/restart the receiver after successful subscription.

pkg/cloudevents/generic/options/v2/grpc/transport.go

Lines changed: 23 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,7 @@ import (
44
"context"
55
"fmt"
66
"sync"
7+
"time"
78

89
cloudevents "github.com/cloudevents/sdk-go/v2"
910
"github.com/cloudevents/sdk-go/v2/binding"
@@ -120,7 +121,28 @@ func (t *grpcTransport) Subscribe(ctx context.Context) error {
120121

121122
values := header.Get(constants.GRPCSubscriptionIDKey)
122123
if len(values) != 1 {
123-
return fmt.Errorf("expected exactly one subscription-id header, got %d", len(values))
124+
// Header() succeeded but no subscription-id was sent (header is nil or empty).
125+
// This typically means the server rejected the subscription before sending headers
126+
// (e.g., authorization failure). The actual error is only available via Recv().
127+
// Call Recv() to get the real error from the server.
128+
recvErrCh := make(chan error, 1)
129+
go func() {
130+
_, err := subClient.Recv()
131+
recvErrCh <- err
132+
}()
133+
select {
134+
case recvErr := <-recvErrCh:
135+
if recvErr != nil {
136+
return recvErr
137+
}
138+
case <-ctx.Done():
139+
return ctx.Err()
140+
case <-time.After(5 * time.Second):
141+
_ = subClient.CloseSend()
142+
return fmt.Errorf("no subscription-id in header (%v): recv timeout", header)
143+
}
144+
// If Recv() didn't return an error, this is a server-side configuration issue
145+
return fmt.Errorf("no subscription-id in header (%v)", header)
124146
}
125147
t.subID = values[0]
126148
t.subClient = subClient

pkg/cloudevents/server/grpc/broker.go

Lines changed: 45 additions & 43 deletions
Original file line numberDiff line numberDiff line change
@@ -122,34 +122,28 @@ func (bkr *GRPCBroker) Publish(ctx context.Context, pubReq *pbv1.PublishRequest)
122122
return &emptypb.Empty{}, nil
123123
}
124124

125-
// register registers a subscriber and return client id and error channel.
126-
func (bkr *GRPCBroker) register(ctx context.Context,
125+
// registerSubscriber registers a subscriber with a pre-generated ID.
126+
// The subscription header must already be sent before calling this function.
127+
func (bkr *GRPCBroker) registerSubscriber(ctx context.Context,
128+
id string,
127129
dataType types.CloudEventsDataType,
128130
subReq *pbv1.SubscriptionRequest,
129-
subServer pbv1.CloudEventService_SubscribeServer,
130-
handler resourceHandler) (string, error) {
131+
handler resourceHandler) error {
131132
logger := klog.FromContext(ctx)
132133

133134
bkr.mu.Lock()
134135
defer bkr.mu.Unlock()
135136

136-
id := uuid.NewString()
137+
logger.Info("registering subscriber", "id", id, "clusterName", subReq.ClusterName, "dataType", dataType)
138+
137139
bkr.subscribers[id] = &subscriber{
138140
clusterName: subReq.ClusterName,
139141
dataType: dataType,
140142
handler: handler,
141143
}
142144

143-
// Signal subscriber is registered
144-
if err := subServer.SendHeader(metadata.Pairs(constants.GRPCSubscriptionIDKey, id)); err != nil {
145-
logger.Error(err, "failed to send subscription header, unregister subscriber", "subID", id)
146-
delete(bkr.subscribers, id)
147-
return "", err
148-
}
149-
logger.V(4).Info("register a subscriber", "id", id, "clusterName", subReq.ClusterName, "dataType", dataType)
150145
metrics.IncGRPCCESubscribersMetric(subReq.ClusterName, dataType.String())
151-
152-
return id, nil
146+
return nil
153147
}
154148

155149
// unregister a subscriber by id
@@ -179,10 +173,17 @@ func (bkr *GRPCBroker) Subscribe(subReq *pbv1.SubscriptionRequest, subServer pbv
179173
return fmt.Errorf("invalid subscription request: invalid data type %v", err)
180174
}
181175

176+
// Generate subscription ID and send header IMMEDIATELY, before any other operations
177+
// This ensures the client receives the header as soon as possible after the stream is established
178+
subID := uuid.NewString()
179+
if err := subServer.SendHeader(metadata.Pairs(constants.GRPCSubscriptionIDKey, subID)); err != nil {
180+
return fmt.Errorf("failed to send subscription header for subID %s: %w", subID, err)
181+
}
182+
182183
subCtx, cancel := context.WithCancel(subServer.Context())
183184
defer cancel()
184185

185-
logger := klog.FromContext(subCtx).WithValues("clusterName", subReq.ClusterName)
186+
logger := klog.FromContext(subCtx).WithValues("clusterName", subReq.ClusterName, "subID", subID)
186187

187188
// TODO make the channel size configurable
188189
eventCh := make(chan *pbv1.CloudEvent, 100)
@@ -193,6 +194,35 @@ func (bkr *GRPCBroker) Subscribe(subReq *pbv1.SubscriptionRequest, subServer pbv
193194
}
194195
sendErrCh := make(chan error, 1)
195196

197+
// Register the subscriber with the ID we already created and sent in the header
198+
err = bkr.registerSubscriber(klog.NewContext(subCtx, logger), subID, *dataType, subReq, func(handlerCtx context.Context, subID string, evt *cloudevents.Event) error {
199+
// convert the cloudevents.Event to pbv1.CloudEvent
200+
// WARNING: don't use "pbEvt, err := pb.ToProto(evt)" to convert cloudevent to protobuf
201+
pbEvt := &pbv1.CloudEvent{}
202+
if err := grpcprotocol.WritePBMessage(handlerCtx, binding.ToMessage(evt), pbEvt); err != nil {
203+
return fmt.Errorf("failed to convert cloudevent to protobuf for resource(%s): %v", evt.ID(), err)
204+
}
205+
206+
// send the cloudevent to the subscriber
207+
klog.FromContext(handlerCtx).V(4).Info("sending the event to spec subscribers",
208+
"subID", subID, "eventType", evt.Type(), "extensions", evt.Extensions())
209+
select {
210+
case eventCh <- pbEvt:
211+
case <-subCtx.Done():
212+
// The context of the stream has been canceled or completed.
213+
// This could happen if:
214+
// - The client closed the connection or canceled the stream.
215+
// - The server closed the stream, potentially due to a shutdown.
216+
// No error is returned here because the stream closure is expected.
217+
return nil
218+
}
219+
220+
return nil
221+
})
222+
if err != nil {
223+
return err
224+
}
225+
196226
// send events
197227
// The grpc send is not concurrency safe and non-blocking, see: https://github.com/grpc/grpc-go/blob/v1.75.1/stream.go#L1571
198228
// Return the error without wrapping, as it includes the gRPC error code and message for further handling.
@@ -236,34 +266,6 @@ func (bkr *GRPCBroker) Subscribe(subReq *pbv1.SubscriptionRequest, subServer pbv
236266
}
237267
}()
238268

239-
subID, err := bkr.register(klog.NewContext(subCtx, logger), *dataType, subReq, subServer, func(handlerCtx context.Context, subID string, evt *cloudevents.Event) error {
240-
// convert the cloudevents.Event to pbv1.CloudEvent
241-
// WARNING: don't use "pbEvt, err := pb.ToProto(evt)" to convert cloudevent to protobuf
242-
pbEvt := &pbv1.CloudEvent{}
243-
if err := grpcprotocol.WritePBMessage(handlerCtx, binding.ToMessage(evt), pbEvt); err != nil {
244-
return fmt.Errorf("failed to convert cloudevent to protobuf for resource(%s): %v", evt.ID(), err)
245-
}
246-
247-
// send the cloudevent to the subscriber
248-
logger.V(4).Info("sending the event to spec subscribers",
249-
"subID", subID, "eventType", evt.Type(), "extensions", evt.Extensions())
250-
select {
251-
case eventCh <- pbEvt:
252-
case <-subCtx.Done():
253-
// The context of the stream has been canceled or completed.
254-
// This could happen if:
255-
// - The client closed the connection or canceled the stream.
256-
// - The server closed the stream, potentially due to a shutdown.
257-
// No error is returned here because the stream closure is expected.
258-
return nil
259-
}
260-
261-
return nil
262-
})
263-
if err != nil {
264-
return err
265-
}
266-
267269
if heartbeater != nil {
268270
go heartbeater.Start(subCtx)
269271
}

0 commit comments

Comments
 (0)