Skip to content

Commit 87b0ebd

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

4 files changed

Lines changed: 286 additions & 49 deletions

File tree

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

Lines changed: 10 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -120,7 +120,16 @@ func (t *grpcTransport) Subscribe(ctx context.Context) error {
120120

121121
values := header.Get(constants.GRPCSubscriptionIDKey)
122122
if len(values) != 1 {
123-
return fmt.Errorf("expected exactly one subscription-id header, got %d", len(values))
123+
// Header() succeeded but no subscription-id was sent (header is nil or empty).
124+
// This typically means the server rejected the subscription before sending headers
125+
// (e.g., authorization failure). The actual error is only available via Recv().
126+
// Call Recv() to get the real error from the server.
127+
_, recvErr := subClient.Recv()
128+
if recvErr != nil {
129+
return recvErr
130+
}
131+
// If Recv() didn't return an error, this is a server-side configuration issue
132+
return fmt.Errorf("no subscription-id in header (%v)", values)
124133
}
125134
t.subID = values[0]
126135
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
}

pkg/cloudevents/server/grpc/broker_test.go

Lines changed: 224 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -133,3 +133,227 @@ func TestServer(t *testing.T) {
133133
t.Error("received event is different")
134134
}
135135
}
136+
137+
// TestSubscriptionHeaderImmediateSend verifies that the subscription ID header
138+
// is sent immediately upon subscription, preventing the "got 0 headers" error.
139+
func TestSubscriptionHeaderImmediateSend(t *testing.T) {
140+
grpcServerOptions := []grpc.ServerOption{}
141+
grpcServer := grpc.NewServer(grpcServerOptions...)
142+
defer grpcServer.Stop()
143+
144+
grpcEventServer := NewGRPCBroker(NewBrokerOptions())
145+
pbv1.RegisterCloudEventServiceServer(grpcServer, grpcEventServer)
146+
147+
svc := &testService{evts: make(map[string]*cloudevents.Event)}
148+
grpcEventServer.RegisterService(context.Background(), dataType, svc)
149+
150+
ctx, cancel := context.WithCancel(context.Background())
151+
defer cancel()
152+
153+
lis, err := net.Listen("tcp", "127.0.0.1:0")
154+
if err != nil {
155+
t.Fatalf("failed to listen: %v", err)
156+
}
157+
t.Cleanup(func() {
158+
grpcServer.GracefulStop()
159+
_ = lis.Close()
160+
})
161+
162+
go func() {
163+
if err := grpcServer.Serve(lis); err != nil {
164+
t.Errorf("failed to serve: %v", err)
165+
}
166+
}()
167+
168+
grpcClientOptions := grpccli.NewGRPCOptions()
169+
grpcClientOptions.Dialer = &grpccli.GRPCDialer{URL: lis.Addr().String()}
170+
agentOption := grpcv2.NewAgentOptions(grpcClientOptions, "cluster1", "agent1", dataType)
171+
172+
// Test that connection and subscription work without "got 0 headers" error
173+
if err := agentOption.CloudEventsTransport.Connect(ctx); err != nil {
174+
t.Fatalf("failed to connect: %v", err)
175+
}
176+
177+
if err := agentOption.CloudEventsTransport.Subscribe(ctx); err != nil {
178+
t.Fatalf("failed to subscribe, expected header to be sent immediately: %v", err)
179+
}
180+
}
181+
182+
// TestReconnectionScenario simulates a client restart (disconnect and reconnect)
183+
// to verify that subscription headers are properly sent on reconnection.
184+
func TestReconnectionScenario(t *testing.T) {
185+
grpcServerOptions := []grpc.ServerOption{}
186+
grpcServer := grpc.NewServer(grpcServerOptions...)
187+
defer grpcServer.Stop()
188+
189+
grpcEventServer := NewGRPCBroker(NewBrokerOptions())
190+
pbv1.RegisterCloudEventServiceServer(grpcServer, grpcEventServer)
191+
192+
svc := &testService{evts: make(map[string]*cloudevents.Event)}
193+
grpcEventServer.RegisterService(context.Background(), dataType, svc)
194+
195+
ctx, cancel := context.WithCancel(context.Background())
196+
defer cancel()
197+
198+
lis, err := net.Listen("tcp", "127.0.0.1:0")
199+
if err != nil {
200+
t.Fatalf("failed to listen: %v", err)
201+
}
202+
t.Cleanup(func() {
203+
grpcServer.GracefulStop()
204+
_ = lis.Close()
205+
})
206+
207+
go func() {
208+
if err := grpcServer.Serve(lis); err != nil {
209+
t.Errorf("failed to serve: %v", err)
210+
}
211+
}()
212+
213+
grpcClientOptions := grpccli.NewGRPCOptions()
214+
grpcClientOptions.Dialer = &grpccli.GRPCDialer{URL: lis.Addr().String()}
215+
216+
// First connection and subscription
217+
agentOption := grpcv2.NewAgentOptions(grpcClientOptions, "cluster1", "agent1", dataType)
218+
if err := agentOption.CloudEventsTransport.Connect(ctx); err != nil {
219+
t.Fatalf("failed to connect: %v", err)
220+
}
221+
222+
if err := agentOption.CloudEventsTransport.Subscribe(ctx); err != nil {
223+
t.Fatalf("failed to subscribe on first connection: %v", err)
224+
}
225+
226+
// Simulate client restart by closing and reconnecting
227+
if err := agentOption.CloudEventsTransport.Close(ctx); err != nil {
228+
t.Fatalf("failed to close: %v", err)
229+
}
230+
231+
// Create a new transport for reconnection
232+
grpcClientOptions2 := grpccli.NewGRPCOptions()
233+
grpcClientOptions2.Dialer = &grpccli.GRPCDialer{URL: lis.Addr().String()}
234+
agentOption2 := grpcv2.NewAgentOptions(grpcClientOptions2, "cluster1", "agent1", dataType)
235+
236+
// Reconnect
237+
if err := agentOption2.CloudEventsTransport.Connect(ctx); err != nil {
238+
t.Fatalf("failed to reconnect: %v", err)
239+
}
240+
241+
// This should not fail with "got 0 headers" error
242+
if err := agentOption2.CloudEventsTransport.Subscribe(ctx); err != nil {
243+
t.Fatalf("failed to subscribe after reconnection: %v", err)
244+
}
245+
}
246+
247+
// TestConcurrentSubscriptions verifies that multiple clients can subscribe
248+
// simultaneously without header race conditions.
249+
func TestConcurrentSubscriptions(t *testing.T) {
250+
grpcServerOptions := []grpc.ServerOption{}
251+
grpcServer := grpc.NewServer(grpcServerOptions...)
252+
defer grpcServer.Stop()
253+
254+
grpcEventServer := NewGRPCBroker(NewBrokerOptions())
255+
pbv1.RegisterCloudEventServiceServer(grpcServer, grpcEventServer)
256+
257+
svc := &testService{evts: make(map[string]*cloudevents.Event)}
258+
grpcEventServer.RegisterService(context.Background(), dataType, svc)
259+
260+
ctx, cancel := context.WithCancel(context.Background())
261+
defer cancel()
262+
263+
lis, err := net.Listen("tcp", "127.0.0.1:0")
264+
if err != nil {
265+
t.Fatalf("failed to listen: %v", err)
266+
}
267+
t.Cleanup(func() {
268+
grpcServer.GracefulStop()
269+
_ = lis.Close()
270+
})
271+
272+
go func() {
273+
if err := grpcServer.Serve(lis); err != nil {
274+
t.Errorf("failed to serve: %v", err)
275+
}
276+
}()
277+
278+
// Create and subscribe multiple clients concurrently
279+
numClients := 10
280+
errCh := make(chan error, numClients)
281+
282+
for i := 0; i < numClients; i++ {
283+
go func(clientID int) {
284+
grpcClientOptions := grpccli.NewGRPCOptions()
285+
grpcClientOptions.Dialer = &grpccli.GRPCDialer{URL: lis.Addr().String()}
286+
agentOption := grpcv2.NewAgentOptions(grpcClientOptions, "cluster1", "agent1", dataType)
287+
288+
if err := agentOption.CloudEventsTransport.Connect(ctx); err != nil {
289+
errCh <- err
290+
return
291+
}
292+
293+
if err := agentOption.CloudEventsTransport.Subscribe(ctx); err != nil {
294+
errCh <- err
295+
return
296+
}
297+
298+
errCh <- nil
299+
}(i)
300+
}
301+
302+
// Wait for all clients to complete
303+
for i := 0; i < numClients; i++ {
304+
if err := <-errCh; err != nil {
305+
t.Errorf("client %d failed: %v", i, err)
306+
}
307+
}
308+
}
309+
310+
// TestMultipleRapidReconnections simulates rapid reconnection scenarios
311+
// that could trigger the race condition.
312+
func TestMultipleRapidReconnections(t *testing.T) {
313+
grpcServerOptions := []grpc.ServerOption{}
314+
grpcServer := grpc.NewServer(grpcServerOptions...)
315+
defer grpcServer.Stop()
316+
317+
grpcEventServer := NewGRPCBroker(NewBrokerOptions())
318+
pbv1.RegisterCloudEventServiceServer(grpcServer, grpcEventServer)
319+
320+
svc := &testService{evts: make(map[string]*cloudevents.Event)}
321+
grpcEventServer.RegisterService(context.Background(), dataType, svc)
322+
323+
ctx, cancel := context.WithCancel(context.Background())
324+
defer cancel()
325+
326+
lis, err := net.Listen("tcp", "127.0.0.1:0")
327+
if err != nil {
328+
t.Fatalf("failed to listen: %v", err)
329+
}
330+
t.Cleanup(func() {
331+
grpcServer.GracefulStop()
332+
_ = lis.Close()
333+
})
334+
335+
go func() {
336+
if err := grpcServer.Serve(lis); err != nil {
337+
t.Errorf("failed to serve: %v", err)
338+
}
339+
}()
340+
341+
// Perform multiple rapid reconnections
342+
for i := 0; i < 5; i++ {
343+
grpcClientOptions := grpccli.NewGRPCOptions()
344+
grpcClientOptions.Dialer = &grpccli.GRPCDialer{URL: lis.Addr().String()}
345+
agentOption := grpcv2.NewAgentOptions(grpcClientOptions, "cluster1", "agent1", dataType)
346+
347+
if err := agentOption.CloudEventsTransport.Connect(ctx); err != nil {
348+
t.Fatalf("reconnection %d: failed to connect: %v", i, err)
349+
}
350+
351+
if err := agentOption.CloudEventsTransport.Subscribe(ctx); err != nil {
352+
t.Fatalf("reconnection %d: failed to subscribe: %v", i, err)
353+
}
354+
355+
if err := agentOption.CloudEventsTransport.Close(ctx); err != nil {
356+
t.Fatalf("reconnection %d: failed to close: %v", i, err)
357+
}
358+
}
359+
}

0 commit comments

Comments
 (0)