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"]) + } +}