Skip to content

Commit 2118b9d

Browse files
committed
send the header immediately
Signed-off-by: Wei Liu <liuweixa@redhat.com>
1 parent bfb55c4 commit 2118b9d

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
@@ -124,34 +124,28 @@ func (bkr *GRPCBroker) Publish(ctx context.Context, pubReq *pbv1.PublishRequest)
124124
return &emptypb.Empty{}, nil
125125
}
126126

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

135136
bkr.mu.Lock()
136137
defer bkr.mu.Unlock()
137138

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

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

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

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

187-
logger := klog.FromContext(subCtx).WithValues("clusterName", subReq.ClusterName)
188+
logger := klog.FromContext(subCtx).WithValues("clusterName", subReq.ClusterName, "subID", subID)
188189

189190
// TODO make the channel size configurable
190191
eventCh := make(chan *pbv1.CloudEvent, 100)
@@ -195,6 +196,35 @@ func (bkr *GRPCBroker) Subscribe(subReq *pbv1.SubscriptionRequest, subServer pbv
195196
}
196197
sendErrCh := make(chan error, 1)
197198

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

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

0 commit comments

Comments
 (0)