Skip to content

Commit 194bf57

Browse files
committed
Fix RequestReply duplicate reply deadlock
Signed-off-by: rootp1 <arnav.iitr@gmail.com>
1 parent 5411419 commit 194bf57

2 files changed

Lines changed: 65 additions & 10 deletions

File tree

pkg/requestreply/ingress_handler.go

Lines changed: 19 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -215,7 +215,9 @@ func (h *IngressHandler) addEvent(responseWriter http.ResponseWriter, event *clo
215215
pr := &proxiedRequest{
216216
received: time.Now(),
217217
responseWriter: responseWriter,
218-
replyEvent: make(chan *cloudevents.Event, 1),
218+
// capacity must stay 1: handleReplyEvent relies on a full buffer to
219+
// detect and drop duplicate replies without blocking.
220+
replyEvent: make(chan *cloudevents.Event, 1),
219221
}
220222
if h.entries[rr.GetNamespacedName()] == nil {
221223
h.entries[rr.GetNamespacedName()] = make(map[string]*proxiedRequest)
@@ -232,6 +234,14 @@ func (h *IngressHandler) deleteEvent(event *cloudevents.Event, rr *v1alpha1.Requ
232234
delete(h.entries[rr.GetNamespacedName()], event.ID())
233235
}
234236

237+
func (h *IngressHandler) getEvent(id string, rr *v1alpha1.RequestReply) (*proxiedRequest, bool) {
238+
h.requestLock.RLock()
239+
defer h.requestLock.RUnlock()
240+
241+
pr, ok := h.entries[rr.GetNamespacedName()][id]
242+
return pr, ok
243+
}
244+
235245
func (h *IngressHandler) handleNewEvent(ctx context.Context, responseWriter http.ResponseWriter, event *cloudevents.Event, rr *v1alpha1.RequestReply, headers http.Header) {
236246
h.logger.Debug("handling new event")
237247
pr := h.addEvent(responseWriter, event, rr)
@@ -289,9 +299,6 @@ func (h *IngressHandler) handleNewEvent(ctx context.Context, responseWriter http
289299
}
290300

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

297304
// TODO: with OIDC enabled, we can skip validation of the key if we validate the identity of the trigger making the request
@@ -336,15 +343,17 @@ func (h *IngressHandler) handleReplyEvent(responseWriter http.ResponseWriter, ev
336343
return
337344
}
338345

339-
responseWriter.WriteHeader(http.StatusAccepted)
340-
341346
id := strings.Split(replyIdString, ":")[0]
342-
pr, ok := h.entries[rr.GetNamespacedName()][id]
347+
pr, ok := h.getEvent(id, rr)
343348
if !ok {
344349
h.logger.Warn("no event found matching the reply id, discarding event", zap.String("reply id", id))
345-
return
350+
} else {
351+
select {
352+
case pr.replyEvent <- event:
353+
default:
354+
h.logger.Warn("reply event already delivered or duplicate reply received, discarding event", zap.String("reply id", id))
355+
}
346356
}
347357

348-
// send the reply event back to the original response writer
349-
pr.replyEvent <- event
358+
responseWriter.WriteHeader(http.StatusAccepted)
350359
}

pkg/requestreply/ingress_handler_test.go

Lines changed: 46 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,51 @@ func TestHandlerServeHttp(t *testing.T) {
210211
}
211212
}
212213

214+
func TestHandleReplyEventIgnoresDuplicateReplyWithoutBlocking(t *testing.T) {
215+
t.Parallel()
216+
217+
ctx, _ := reconcilertesting.SetupFakeContext(t, setupInformerSelector)
218+
219+
rr := makeRequestReply("my-request-reply", "default")
220+
inflight := cloudevents.NewEvent()
221+
inflight.SetID("1234567890")
222+
223+
firstReply := inflight.Clone()
224+
if err := SetCorrelationId(&firstReply, "replyid", exampleKey, 0); err != nil {
225+
t.Fatalf("failed to create first reply id: %v", err)
226+
}
227+
228+
duplicateReply := inflight.Clone()
229+
if err := SetCorrelationId(&duplicateReply, "replyid", exampleKey, 0); err != nil {
230+
t.Fatalf("failed to create duplicate reply id: %v", err)
231+
}
232+
233+
keyStore := &AESKeyStore{}
234+
keyStore.addAesKey(rr.GetNamespacedName(), "key", exampleKey)
235+
236+
handler := NewHandler(zap.NewNop(), requestreplyinformerfake.Get(ctx), configmapinformerfake.Get(ctx).Lister().ConfigMaps("ns"), keyStore, 0)
237+
238+
// Fill the reply channel once to simulate an already delivered reply.
239+
pr := handler.addEvent(httptest.NewRecorder(), &inflight, rr)
240+
pr.replyEvent <- &firstReply
241+
242+
recorder := httptest.NewRecorder()
243+
done := make(chan struct{})
244+
245+
go func() {
246+
handler.handleReplyEvent(recorder, &duplicateReply, rr)
247+
close(done)
248+
}()
249+
250+
select {
251+
case <-done:
252+
case <-time.After(2 * time.Second):
253+
t.Fatal("handleReplyEvent blocked on duplicate reply")
254+
}
255+
256+
assert.Equal(t, http.StatusAccepted, recorder.Result().StatusCode)
257+
}
258+
213259
type testServerHandler struct {
214260
makeReplyEvent func(e *cloudevents.Event) *cloudevents.Event
215261
callbackHandler http.Handler

0 commit comments

Comments
 (0)