Skip to content

Commit 280c5b7

Browse files
committed
Fix RequestReply duplicate reply deadlock
1 parent 5411419 commit 280c5b7

2 files changed

Lines changed: 75 additions & 8 deletions

File tree

pkg/requestreply/ingress_handler.go

Lines changed: 17 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -232,6 +232,14 @@ func (h *IngressHandler) deleteEvent(event *cloudevents.Event, rr *v1alpha1.Requ
232232
delete(h.entries[rr.GetNamespacedName()], event.ID())
233233
}
234234

235+
func (h *IngressHandler) getEvent(id string, rr *v1alpha1.RequestReply) (*proxiedRequest, bool) {
236+
h.requestLock.RLock()
237+
defer h.requestLock.RUnlock()
238+
239+
pr, ok := h.entries[rr.GetNamespacedName()][id]
240+
return pr, ok
241+
}
242+
235243
func (h *IngressHandler) handleNewEvent(ctx context.Context, responseWriter http.ResponseWriter, event *cloudevents.Event, rr *v1alpha1.RequestReply, headers http.Header) {
236244
h.logger.Debug("handling new event")
237245
pr := h.addEvent(responseWriter, event, rr)
@@ -289,9 +297,6 @@ func (h *IngressHandler) handleNewEvent(ctx context.Context, responseWriter http
289297
}
290298

291299
func (h *IngressHandler) handleReplyEvent(responseWriter http.ResponseWriter, event *cloudevents.Event, rr *v1alpha1.RequestReply) {
292-
h.requestLock.RLock()
293-
defer h.requestLock.RUnlock()
294-
295300
h.logger.Debug("handling a response event")
296301

297302
// TODO: with OIDC enabled, we can skip validation of the key if we validate the identity of the trigger making the request
@@ -336,15 +341,19 @@ func (h *IngressHandler) handleReplyEvent(responseWriter http.ResponseWriter, ev
336341
return
337342
}
338343

339-
responseWriter.WriteHeader(http.StatusAccepted)
340-
341344
id := strings.Split(replyIdString, ":")[0]
342-
pr, ok := h.entries[rr.GetNamespacedName()][id]
345+
pr, ok := h.getEvent(id, rr)
343346
if !ok {
344347
h.logger.Warn("no event found matching the reply id, discarding event", zap.String("reply id", id))
348+
responseWriter.WriteHeader(http.StatusAccepted)
345349
return
346350
}
347351

348-
// send the reply event back to the original response writer
349-
pr.replyEvent <- event
352+
select {
353+
case pr.replyEvent <- event:
354+
responseWriter.WriteHeader(http.StatusAccepted)
355+
default:
356+
h.logger.Warn("reply event already delivered or duplicate reply received, discarding event", zap.String("reply id", id))
357+
responseWriter.WriteHeader(http.StatusAccepted)
358+
}
350359
}

pkg/requestreply/ingress_handler_test.go

Lines changed: 58 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -24,6 +24,7 @@ import (
2424
"net/http/httptest"
2525
"strings"
2626
"testing"
27+
"time"
2728

2829
cloudevents "github.com/cloudevents/sdk-go/v2"
2930
cehttp "github.com/cloudevents/sdk-go/v2/protocol/http"
@@ -210,6 +211,63 @@ func TestHandlerServeHttp(t *testing.T) {
210211
}
211212
}
212213

214+
func TestHandleReplyEventIgnoresDuplicateReplyWithoutBlocking(t *testing.T) {
215+
t.Parallel()
216+
217+
rr := makeRequestReply("my-request-reply", "default")
218+
inflight := cloudevents.NewEvent()
219+
inflight.SetID("1234567890")
220+
221+
firstReply := inflight.Clone()
222+
if err := SetCorrelationId(&firstReply, "replyid", exampleKey, 0); err != nil {
223+
t.Fatalf("failed to create first reply id: %v", err)
224+
}
225+
226+
duplicateReply := inflight.Clone()
227+
if err := SetCorrelationId(&duplicateReply, "replyid", exampleKey, 0); err != nil {
228+
t.Fatalf("failed to create duplicate reply id: %v", err)
229+
}
230+
231+
keyStore := &AESKeyStore{}
232+
keyStore.addAesKey(rr.GetNamespacedName(), "key", exampleKey)
233+
234+
handler := &IngressHandler{
235+
logger: zap.NewNop(),
236+
keyStore: keyStore,
237+
entries: map[types.NamespacedName]map[string]*proxiedRequest{
238+
rr.GetNamespacedName(): {
239+
inflight.ID(): {
240+
responseWriter: httptest.NewRecorder(),
241+
replyEvent: make(chan *cloudevents.Event, 1),
242+
},
243+
},
244+
},
245+
}
246+
247+
// Fill the reply channel once to simulate an already delivered reply.
248+
pr, ok := handler.getEvent(inflight.ID(), rr)
249+
if !ok {
250+
t.Fatal("expected inflight request to exist")
251+
}
252+
pr.replyEvent <- &firstReply
253+
254+
recorder := httptest.NewRecorder()
255+
done := make(chan struct{})
256+
257+
go func() {
258+
handler.handleReplyEvent(recorder, &duplicateReply, rr)
259+
close(done)
260+
}()
261+
262+
select {
263+
case <-done:
264+
case <-time.After(2 * time.Second):
265+
t.Fatal("handleReplyEvent blocked on duplicate reply")
266+
}
267+
268+
assert.Equal(t, http.StatusAccepted, recorder.Result().StatusCode)
269+
}
270+
213271
type testServerHandler struct {
214272
makeReplyEvent func(e *cloudevents.Event) *cloudevents.Event
215273
callbackHandler http.Handler

0 commit comments

Comments
 (0)