From 927c96c16433faba3364c05951658740c476fec1 Mon Sep 17 00:00:00 2001 From: mayankpande88 Date: Tue, 6 Oct 2026 01:13:05 +0530 Subject: [PATCH] fix(l7): make room for new HTTP/2 streams by dropping the oldest waiting one The parser tracks at most 100 requests per connection that are waiting for their response. A request whose response it never sees (the read was cut short by truncation, or the event was lost) waits until the stream GC, two minutes or more. Once 100 of those piled up, new streams were refused, silently: every request on the connection was dropped until the GC ran. At the limit the request that has waited longest is now dropped instead, and counted as the stream_evicted stage. --- containers/l7_self_metrics.go | 3 +++ ebpftracer/l7/http2.go | 33 +++++++++++++++++++++--- ebpftracer/l7/http2_hpack_test.go | 42 +++++++++++++++++++++++++++++++ 3 files changed, 74 insertions(+), 4 deletions(-) diff --git a/containers/l7_self_metrics.go b/containers/l7_self_metrics.go index c4c89cc..9c57af7 100644 --- a/containers/l7_self_metrics.go +++ b/containers/l7_self_metrics.go @@ -127,6 +127,9 @@ var ( // inferred. Stages, in order: // // stream_created client HEADERS decoded, request object created + // stream_evicted a request still waiting for its response dropped to + // make room: the connection had too many such requests, + // nearly always ones whose response was lost // response_status :status seen on the response // end_stream END_STREAM flag seen (a frame flag, not HPACK) // completed both of the above -> request emitted diff --git a/ebpftracer/l7/http2.go b/ebpftracer/l7/http2.go index d7ba22d..a81e478 100644 --- a/ebpftracer/l7/http2.go +++ b/ebpftracer/l7/http2.go @@ -105,7 +105,8 @@ const ( maxPendingHeaderBlockSize = 64 * 1024 // Max concurrent HTTP/2 streams tracked per connection. - // Prevents unbounded memory growth when responses never complete (orphan streams). + // Prevents unbounded memory growth when responses never complete (orphan + // streams); at the limit the oldest stream makes way (evictOldestRequest). maxActiveRequests = 100 ) @@ -218,6 +219,29 @@ func (p *Http2Parser) resetDecoder(method Method) { } } +// evictOldestRequest makes room for a new stream by dropping the request +// that has waited longest for its response. +// +// The streams that fill the table are mostly ones whose response the parser +// will never see: it was in a read cut short by truncation, or in an event +// lost before it. They are only collected after http2DecoderGcInterval. +// Refusing new streams until then dropped every request on a busy +// connection for minutes, silently; the oldest stream is the one least +// likely to still complete. +func (p *Http2Parser) evictOldestRequest() { + var oldestId uint32 + var oldest *Http2Request + for id, r := range p.activeRequests { + if oldest == nil || r.kernelTime < oldest.kernelTime { + oldestId, oldest = id, r + } + } + if oldest != nil { + delete(p.activeRequests, oldestId) + p.stage("stream_evicted") + } +} + // dropPendingHeaders discards a header block still waiting for CONTINUATION // frames. Its insertions never reach the table, so the table is reset too. func (p *Http2Parser) dropPendingHeaders(method Method, pending **pendingHeaderBlock) { @@ -297,15 +321,16 @@ func (p *Http2Parser) decodeHeaderBlock( switch method { case MethodHttp2ClientFrames: req := p.activeRequests[streamId] - if req == nil && len(p.activeRequests) < maxActiveRequests { + if req == nil { + if len(p.activeRequests) >= maxActiveRequests { + p.evictOldestRequest() + } req = &Http2Request{ kernelTime: kernelTime, } p.activeRequests[streamId] = req p.stage("stream_created") } - // With too many active streams req stays nil: the block is still - // decoded, to keep the dynamic table in sync. emit = func(name, value string) { switch name { case ":method": diff --git a/ebpftracer/l7/http2_hpack_test.go b/ebpftracer/l7/http2_hpack_test.go index 9243d58..8d83a13 100644 --- a/ebpftracer/l7/http2_hpack_test.go +++ b/ebpftracer/l7/http2_hpack_test.go @@ -188,3 +188,45 @@ func TestHttp2ParserAcceptsTableSizeUpdateInLaterBlock(t *testing.T) { t.Errorf("hpack_error = %d, hpack_partial = %d, want 0", stages["hpack_error"], stages["hpack_partial"]) } } + +func responseFrame(streamID uint32) []byte { + var buf bytes.Buffer + _ = hpack.NewEncoder(&buf).WriteField(hpack.HeaderField{Name: ":status", Value: "200"}) + return frame(http2.FrameHeaders, http2FlagEndHeaders|http2FlagEndStream, streamID, buf.Bytes()) +} + +// Requests whose responses are lost (cut from a truncated read) stay active +// until the stream GC, minutes later. Once maxActiveRequests of them pile up, +// new streams used to be refused, so every request on the connection was +// dropped until the GC ran. The oldest waiting request now makes room. +func TestHttp2ParserEvictsOldestStreamAtCap(t *testing.T) { + stages := countStages(t) + p := NewHttp2Parser() + orphans := maxActiveRequests + 20 + for i := 0; i < orphans; i++ { + p.Parse(MethodHttp2ClientFrames, headersFrame(streamID(i), "/orphan"), uint64(i), 0) + } + completed := 0 + for i := orphans; i < orphans+50; i++ { + p.Parse(MethodHttp2ClientFrames, headersFrame(streamID(i), "/live"), uint64(i), 0) + for _, r := range p.Parse(MethodHttp2ServerFrames, responseFrame(streamID(i)), uint64(i), 0) { + if r.Path == "/live" { + completed++ + } + } + } + if completed != 50 { + t.Errorf("completed %d of 50 requests made after the orphans filled the table", completed) + } + if len(p.activeRequests) > maxActiveRequests { + t.Errorf("%d active requests, cap is %d", len(p.activeRequests), maxActiveRequests) + } + if p.activeRequests[streamID(0)] != nil || p.activeRequests[streamID(orphans-1)] == nil { + t.Error("evicted the wrong streams: the oldest must go first") + } + // 20 orphans past the cap, then one for the first live request; each + // live request completes and frees its own slot. + if stages["stream_evicted"] != 21 { + t.Errorf("stream_evicted = %d, want 21", stages["stream_evicted"]) + } +}