From 84b3782e53cb11b585a4908ad00f86a5f5d43282 Mon Sep 17 00:00:00 2001 From: mayankpande88 Date: Tue, 6 Oct 2026 00:43:51 +0530 Subject: [PATCH 1/3] fix(l7): keep decoding HTTP/2 headers after the HPACK table falls out of sync The parser decoded header blocks with hpack.Decoder and replaced it with a fresh one after any error. When the agent joins a connection after it opened, or loses a header block, the peer's encoder keeps referencing dynamic table entries the decoder never saw. hpack.Decoder fails at the first such reference, so the block's later insertions are never applied and the next block fails the same way. Resetting changes nothing: on a connection whose encoder references an old entry in every request, every block failed for good. Whatever the table supplies was lost with it: :path and :authority of requests, and grpc-status and any :status outside the static table (201, 401, 503, ...) of responses. Header blocks are now decoded by a decoder that skips references to entries it does not hold and applies every insertion. Indices count back from the newest entry, so everything inserted since it lost track sits at the index the encoder uses, and it converges on the encoder's table as old entries are evicted. In a test that joins a connection after 20 requests, static-named headers (:method, :path, :authority, content-type) decode in full from request 63; the reset-on-error decoder failed all 380 blocks. Missing a block's insertions would shift older indices, so the table is now reset where a header block is known to be lost: a HEADERS or CONTINUATION frame cut by truncation, a pending header block dropped, a CONTINUATION without its HEADERS. A block that decodes to pseudo-headers impossible for its direction (:method in a response, :status in a request) also resets it. References the decoder cannot resolve are counted as the hpack_partial stage, not as decode errors. Larger dynamic table size updates, which hpack.Decoder(4096) rejected, are accepted. --- containers/l7_self_metrics.go | 15 +- ebpftracer/l7/hpack.go | 294 ++++++++++++++++++++++++++++++ ebpftracer/l7/hpack_test.go | 292 +++++++++++++++++++++++++++++ ebpftracer/l7/http2.go | 159 +++++++++------- ebpftracer/l7/http2_hpack_test.go | 154 ++++++++++++++++ 5 files changed, 844 insertions(+), 70 deletions(-) create mode 100644 ebpftracer/l7/hpack.go create mode 100644 ebpftracer/l7/hpack_test.go create mode 100644 ebpftracer/l7/http2_hpack_test.go diff --git a/containers/l7_self_metrics.go b/containers/l7_self_metrics.go index e08bc2f..c4c89cc 100644 --- a/containers/l7_self_metrics.go +++ b/containers/l7_self_metrics.go @@ -12,14 +12,16 @@ import ( ) var ( - // HPACKDecodeErrorsTotal counts HPACK decode failures in the HTTP/2 - // parser: the decoder's dynamic table no longer matches the peer's, as - // when the agent joined a long-lived connection mid-stream or missed a - // HEADERS frame. + // HPACKDecodeErrorsTotal counts HTTP/2 header blocks that were not valid + // HPACK, or decoded to pseudo-headers that cannot be right for their + // direction: the decoder's dynamic table had drifted from the peer's, as + // when a HEADERS frame was lost unnoticed. References to entries inserted + // before the agent joined the connection are not errors; they are counted + // as the hpack_partial stage of node_agent_http2_stage_total. HPACKDecodeErrorsTotal = prometheus.NewCounter( prometheus.CounterOpts{ Name: "node_agent_hpack_decode_errors_total", - Help: "Total HPACK decode errors in HTTP/2 parser (mid-stream join indicator)", + Help: "HTTP/2 header blocks that failed to decode or decoded to implausible headers; the decoder's table is reset", }, ) @@ -128,6 +130,9 @@ var ( // 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 + // hpack_partial a header block referenced table entries the decoder + // does not hold (inserted before it joined or was + // reset); the other headers in it were decoded // hpack_error HPACK block failed to decode; decoder reset // // A request is only emitted with BOTH response_status and end_stream, so diff --git a/ebpftracer/l7/hpack.go b/ebpftracer/l7/hpack.go new file mode 100644 index 0000000..9a3fceb --- /dev/null +++ b/ebpftracer/l7/hpack.go @@ -0,0 +1,294 @@ +package l7 + +import ( + "errors" + + "golang.org/x/net/http2/hpack" +) + +// hpackDecoder decodes HPACK header blocks (RFC 7541) like hpack.Decoder, +// except that a reference to a dynamic table entry it does not hold is +// skipped instead of failing the block. +// +// The agent routinely decodes a connection without the peers' full table +// history: it joined the connection after it opened, or a header block was +// lost to truncation. hpack.Decoder fails at the first reference to an entry +// it never saw, so the insertions later in that block are never applied, the +// next block fails the same way, and resetting it changes nothing: once the +// encoder references an old entry on every request, every block of that +// connection fails for good. +// +// Indices count back from the newest entry. Every entry inserted since the +// decoder lost track is therefore at the index the encoder uses, provided no +// insertion is missed. Applying every insertion and skipping only the +// references beyond what the decoder holds converges on the encoder's table +// as its older entries are evicted. A block that was lost (and so its +// insertions) breaks that alignment silently; callers reset the decoder when +// they know a block was lost. +type hpackDecoder struct { + // dynamic holds the entries the decoder knows, newest last. + dynamic []hpackEntry + size uint32 // RFC 7541 4.1 size of dynamic + maxSize uint32 // set by dynamic table size updates +} + +// hpackEntry is a dynamic table entry. noName marks one inserted by a literal +// whose indexed name the decoder did not hold: its value is known, its name +// is not. +type hpackEntry struct { + hpack.HeaderField + noName bool +} + +const ( + hpackDefaultTableSize = 4096 + // hpackMaxTableSize bounds what a size update may set. Peers that agreed + // on a larger SETTINGS_HEADER_TABLE_SIZE use more than 4096; hpack.Decoder + // rejected their updates. + hpackMaxTableSize = 64 * 1024 + hpackEntryOverhead = 32 +) + +var ( + errHpackTruncated = errors.New("hpack: header block truncated") + errHpackInvalid = errors.New("hpack: invalid representation") +) + +func newHpackDecoder() *hpackDecoder { + return &hpackDecoder{maxSize: hpackDefaultTableSize} +} + +// reset forgets the dynamic table, for when the caller knows a header block +// was lost and the table no longer matches the encoder's. +func (d *hpackDecoder) reset() { + d.dynamic = d.dynamic[:0] + d.size = 0 + d.maxSize = hpackDefaultTableSize +} + +// decode decodes one complete header block, calling emit for every field it +// can resolve. unknown counts the fields it could not: references to dynamic +// entries it does not hold, and literals whose indexed name it does not hold +// (their values are still inserted, to keep later indices aligned). An error +// means the block is not valid HPACK, and the table may no longer match. +func (d *hpackDecoder) decode(block []byte, emit func(name, value string)) (unknown int, err error) { + for len(block) > 0 { + b := block[0] + switch { + case b&0x80 != 0: // 6.1 Indexed Header Field + var idx uint64 + if idx, block, err = hpackInt(block, 7); err != nil { + return unknown, err + } + if idx == 0 { + return unknown, errHpackInvalid + } + if f, ok := d.field(idx); ok && !f.noName { + emit(f.Name, f.Value) + } else { + unknown++ + } + case b&0xe0 == 0x20: // 6.3 Dynamic Table Size Update + var size uint64 + if size, block, err = hpackInt(block, 5); err != nil { + return unknown, err + } + if size > hpackMaxTableSize { + return unknown, errHpackInvalid + } + d.maxSize = uint32(size) + d.evict(0) + default: // 6.2 Literal Header Field + prefix, index := uint8(4), false + if b&0xc0 == 0x40 { // with incremental indexing + prefix, index = 6, true + } + var nameIdx uint64 + if nameIdx, block, err = hpackInt(block, prefix); err != nil { + return unknown, err + } + var name, value string + nameKnown := true + if nameIdx == 0 { + if name, block, err = hpackString(block); err != nil { + return unknown, err + } + } else if f, ok := d.field(nameIdx); ok && !f.noName { + name = f.Name + } else { + nameKnown = false + } + if value, block, err = hpackString(block); err != nil { + return unknown, err + } + if nameKnown { + emit(name, value) + } else { + unknown++ + } + if index { + // An entry whose name is unknown still takes its place in the + // table. Its size is underestimated by the name's length, so it + // is evicted later than the encoder evicts it: an entry the + // decoder keeps too long sits past every index the encoder + // still uses, while one evicted too early would shift them. + d.insert(hpackEntry{HeaderField: hpack.HeaderField{Name: name, Value: value}, noName: !nameKnown}) + } + } + } + return unknown, nil +} + +// field resolves an index into the static table or the known part of the +// dynamic table. +func (d *hpackDecoder) field(idx uint64) (hpackEntry, bool) { + if idx <= uint64(len(hpackStaticTable)) { + return hpackEntry{HeaderField: hpackStaticTable[idx-1]}, true + } + i := idx - uint64(len(hpackStaticTable)) // 1 = newest + if i > uint64(len(d.dynamic)) { + return hpackEntry{}, false + } + return d.dynamic[uint64(len(d.dynamic))-i], true +} + +func (d *hpackDecoder) insert(f hpackEntry) { + size := uint32(len(f.Name)+len(f.Value)) + hpackEntryOverhead + if size > d.maxSize { + // RFC 7541 4.4: an entry larger than the table empties it. + d.dynamic = d.dynamic[:0] + d.size = 0 + return + } + d.evict(size) + d.dynamic = append(d.dynamic, f) + d.size += size +} + +// evict drops the oldest entries until room more bytes fit. +func (d *hpackDecoder) evict(room uint32) { + n := 0 + for d.size+room > d.maxSize && n < len(d.dynamic) { + f := d.dynamic[n] + d.size -= uint32(len(f.Name)+len(f.Value)) + hpackEntryOverhead + n++ + } + if n > 0 { + d.dynamic = append(d.dynamic[:0], d.dynamic[n:]...) + } +} + +// hpackInt decodes an RFC 7541 5.1 integer with an n-bit prefix. +func hpackInt(b []byte, n uint8) (uint64, []byte, error) { + if len(b) == 0 { + return 0, b, errHpackTruncated + } + mask := byte(1< 28 { // nothing in a header block needs more than 32 bits + return 0, nil, errHpackInvalid + } + } + return 0, nil, errHpackTruncated +} + +// hpackString decodes an RFC 7541 5.2 string literal. +func hpackString(b []byte) (string, []byte, error) { + if len(b) == 0 { + return "", b, errHpackTruncated + } + huffman := b[0]&0x80 != 0 + n, b, err := hpackInt(b, 7) + if err != nil { + return "", nil, err + } + if n > uint64(len(b)) { + return "", nil, errHpackTruncated + } + raw := b[:n] + b = b[n:] + if !huffman { + return string(raw), b, nil + } + s, err := hpack.HuffmanDecodeToString(raw) + if err != nil { + return "", nil, errHpackInvalid + } + return s, b, nil +} + +// hpackStaticTable is RFC 7541 Appendix A. +var hpackStaticTable = [...]hpack.HeaderField{ + {Name: ":authority"}, + {Name: ":method", Value: "GET"}, + {Name: ":method", Value: "POST"}, + {Name: ":path", Value: "/"}, + {Name: ":path", Value: "/index.html"}, + {Name: ":scheme", Value: "http"}, + {Name: ":scheme", Value: "https"}, + {Name: ":status", Value: "200"}, + {Name: ":status", Value: "204"}, + {Name: ":status", Value: "206"}, + {Name: ":status", Value: "304"}, + {Name: ":status", Value: "400"}, + {Name: ":status", Value: "404"}, + {Name: ":status", Value: "500"}, + {Name: "accept-charset"}, + {Name: "accept-encoding", Value: "gzip, deflate"}, + {Name: "accept-language"}, + {Name: "accept-ranges"}, + {Name: "accept"}, + {Name: "access-control-allow-origin"}, + {Name: "age"}, + {Name: "allow"}, + {Name: "authorization"}, + {Name: "cache-control"}, + {Name: "content-disposition"}, + {Name: "content-encoding"}, + {Name: "content-language"}, + {Name: "content-length"}, + {Name: "content-location"}, + {Name: "content-range"}, + {Name: "content-type"}, + {Name: "cookie"}, + {Name: "date"}, + {Name: "etag"}, + {Name: "expect"}, + {Name: "expires"}, + {Name: "from"}, + {Name: "host"}, + {Name: "if-match"}, + {Name: "if-modified-since"}, + {Name: "if-none-match"}, + {Name: "if-range"}, + {Name: "if-unmodified-since"}, + {Name: "last-modified"}, + {Name: "link"}, + {Name: "location"}, + {Name: "max-forwards"}, + {Name: "proxy-authenticate"}, + {Name: "proxy-authorization"}, + {Name: "range"}, + {Name: "referer"}, + {Name: "refresh"}, + {Name: "retry-after"}, + {Name: "server"}, + {Name: "set-cookie"}, + {Name: "strict-transport-security"}, + {Name: "transfer-encoding"}, + {Name: "user-agent"}, + {Name: "vary"}, + {Name: "via"}, + {Name: "www-authenticate"}, +} diff --git a/ebpftracer/l7/hpack_test.go b/ebpftracer/l7/hpack_test.go new file mode 100644 index 0000000..f60b463 --- /dev/null +++ b/ebpftracer/l7/hpack_test.go @@ -0,0 +1,292 @@ +package l7 + +import ( + "bytes" + "fmt" + "math/rand" + "testing" + + "golang.org/x/net/http2/hpack" +) + +// requestHeaders is a header list shaped like real traffic: a few stable +// headers the encoder indexes once and references from then on, a path from a +// small set, and a per-request ID that is inserted every time and so turns the +// dynamic table over. +func requestHeaders(i int) []hpack.HeaderField { + return []hpack.HeaderField{ + {Name: ":method", Value: "GET"}, + {Name: ":scheme", Value: "https"}, + {Name: ":authority", Value: "api.example.com"}, + {Name: ":path", Value: fmt.Sprintf("/v1/items/%d", i%7)}, + {Name: "user-agent", Value: "client/1.2.3"}, + {Name: "x-request-id", Value: fmt.Sprintf("req-%08d-%08d", i, i*7919)}, + } +} + +func encodeBlocks(t testing.TB, n int, headers func(int) []hpack.HeaderField) [][]byte { + var buf bytes.Buffer + enc := hpack.NewEncoder(&buf) + blocks := make([][]byte, n) + for i := range blocks { + buf.Reset() + for _, f := range headers(i) { + if err := enc.WriteField(f); err != nil { + t.Fatal(err) + } + } + blocks[i] = append([]byte(nil), buf.Bytes()...) + } + return blocks +} + +func decodeAll(t testing.TB, d *hpackDecoder, block []byte) ([]hpack.HeaderField, int) { + var got []hpack.HeaderField + unknown, err := d.decode(block, func(name, value string) { + got = append(got, hpack.HeaderField{Name: name, Value: value}) + }) + if err != nil { + t.Fatalf("decode: %v", err) + } + return got, unknown +} + +func sameFields(a, b []hpack.HeaderField) bool { + if len(a) != len(b) { + return false + } + for i := range a { + if a[i].Name != b[i].Name || a[i].Value != b[i].Value { + return false + } + } + return true +} + +// With the whole connection seen, the decoder must decode exactly what +// hpack.Decoder decodes, across indexing modes, Huffman coding, evictions and +// table size updates. +func TestHpackDecoderMatchesReference(t *testing.T) { + rng := rand.New(rand.NewSource(1)) + var buf bytes.Buffer + enc := hpack.NewEncoder(&buf) + ref := hpack.NewDecoder(hpackDefaultTableSize, nil) + d := newHpackDecoder() + for i := 0; i < 2000; i++ { + buf.Reset() + if i%300 == 150 { + enc.SetMaxDynamicTableSize(uint32(256 + rng.Intn(3840))) + } + var want []hpack.HeaderField + for j := 0; j < 1+rng.Intn(12); j++ { + f := hpack.HeaderField{ + Name: fmt.Sprintf("x-h%d", rng.Intn(20)), + Value: fmt.Sprintf("%x", rng.Int63n(1<= 0 && !complete { + t.Fatalf("block %d: static-named headers missing after recovering at block %d: got %v, want %v", i, recoveredAt, got, want) + } + } + if recoveredAt < 0 { + t.Fatal("never recovered") + } + t.Logf("tolerant decoder complete from block %d; reset-on-error decoder failed %d of %d blocks", recoveredAt, refErrors, n-join) + if refErrors != n-join { + t.Fatalf("reference decoder failed %d of %d blocks; this test no longer shows the cascade", refErrors, n-join) + } +} + +// A literal whose indexed name the decoder does not hold still inserts an +// entry; skipping the insertion would shift every older index by one and +// decode later references to the wrong header. +func TestHpackDecoderKeepsAlignmentThroughUnknownNames(t *testing.T) { + var buf bytes.Buffer + enc := hpack.NewEncoder(&buf) + write := func(fs ...hpack.HeaderField) []byte { + buf.Reset() + for _, f := range fs { + if err := enc.WriteField(f); err != nil { + t.Fatal(err) + } + } + return append([]byte(nil), buf.Bytes()...) + } + write(hpack.HeaderField{Name: "x-custom", Value: "a"}) // before the join + d := newHpackDecoder() + // Name from the unknown entry, new value: inserted with an unknown name. + b1 := write(hpack.HeaderField{Name: "x-custom", Value: "b"}) + b2 := write(hpack.HeaderField{Name: "x-other", Value: "c"}) + b3 := write(hpack.HeaderField{Name: "x-other", Value: "c"}, hpack.HeaderField{Name: "x-custom", Value: "b"}) + + if got, unknown := decodeAll(t, d, b1); len(got) != 0 || unknown != 1 { + t.Fatalf("b1: got %v, %d unknown", got, unknown) + } + decodeAll(t, d, b2) + got, unknown := decodeAll(t, d, b3) + if unknown != 1 || len(got) != 1 || got[0].Name != "x-other" || got[0].Value != "c" { + t.Fatalf("b3: got %v, %d unknown; want x-other: c and one unknown", got, unknown) + } +} + +// Peers that agreed on a larger SETTINGS_HEADER_TABLE_SIZE send larger size +// updates; hpack.Decoder(4096) rejected them. +func TestHpackDecoderAcceptsLargerTable(t *testing.T) { + var buf bytes.Buffer + enc := hpack.NewEncoder(&buf) + enc.SetMaxDynamicTableSizeLimit(16384) + enc.SetMaxDynamicTableSize(16384) + if err := enc.WriteField(hpack.HeaderField{Name: "x-a", Value: "1"}); err != nil { + t.Fatal(err) + } + got, unknown := decodeAll(t, newHpackDecoder(), buf.Bytes()) + if unknown != 0 || len(got) != 1 || got[0].Value != "1" { + t.Fatalf("got %v, %d unknown", got, unknown) + } +} + +func TestHpackDecoderRejectsMalformedBlocks(t *testing.T) { + for name, block := range map[string][]byte{ + "index 0": {0x80}, + "truncated integer": {0xff, 0x80}, + "truncated string": {0x40, 0x05, 'a'}, + "integer overflow": {0xff, 0xff, 0xff, 0xff, 0xff, 0xff, 0x01}, + "huge size update": {0x3f, 0xe1, 0xff, 0x7f}, + "bad huffman": {0x40, 0x81, 0xff, 0x00}, + } { + if _, err := newHpackDecoder().decode(block, func(string, string) {}); err == nil { + t.Errorf("%s: no error", name) + } + } + // Random bytes must never panic. + rng := rand.New(rand.NewSource(2)) + d := newHpackDecoder() + for i := 0; i < 20000; i++ { + b := make([]byte, rng.Intn(64)) + rng.Read(b) + if _, err := d.decode(b, func(string, string) {}); err != nil { + d.reset() + } + } +} + +func BenchmarkHpackDecoder(b *testing.B) { + blocks := encodeBlocks(b, 1000, requestHeaders) + emit := func(string, string) {} + b.Run("tolerant", func(b *testing.B) { + b.ReportAllocs() + for i := 0; i < b.N; i++ { + d := newHpackDecoder() + for _, block := range blocks { + if _, err := d.decode(block, emit); err != nil { + b.Fatal(err) + } + } + } + }) + b.Run("x/net", func(b *testing.B) { + b.ReportAllocs() + for i := 0; i < b.N; i++ { + d := hpack.NewDecoder(hpackDefaultTableSize, func(hpack.HeaderField) {}) + for _, block := range blocks { + if _, err := d.Write(block); err != nil { + b.Fatal(err) + } + } + } + }) +} diff --git a/ebpftracer/l7/http2.go b/ebpftracer/l7/http2.go index b276d23..d7ba22d 100644 --- a/ebpftracer/l7/http2.go +++ b/ebpftracer/l7/http2.go @@ -8,7 +8,6 @@ import ( "time" "golang.org/x/net/http2" - "golang.org/x/net/http2/hpack" "k8s.io/klog/v2" ) @@ -163,8 +162,8 @@ type Http2Parser struct { // frames, indefinitely, because the protocol is cached per connection. sawValidFrame bool - clientDecoder *hpack.Decoder - serverDecoder *hpack.Decoder + clientDecoder *hpackDecoder + serverDecoder *hpackDecoder activeRequests map[uint32]*Http2Request lastGcTime uint64 @@ -195,29 +194,39 @@ type Http2Parser struct { func NewHttp2Parser() *Http2Parser { return &Http2Parser{ - clientDecoder: hpack.NewDecoder(4096, nil), - serverDecoder: hpack.NewDecoder(4096, nil), + clientDecoder: newHpackDecoder(), + serverDecoder: newHpackDecoder(), activeRequests: make(map[uint32]*Http2Request), statuses: make(map[uint32]Status), grpcStatuses: make(map[uint32]Status), } } -// resetDecoder creates a fresh HPACK decoder after an unrecoverable decode error. -// This discards the dynamic table but preserves static table (indices 1-61) functionality. -// Headers like :method (2,3), :path (4,5), :scheme (6,7), :status (8-14) remain decodable. -// The dynamic table gradually rebuilds as the encoder sends new literal-with-indexing headers. +// resetDecoder forgets a direction's HPACK dynamic table, when a header block +// was lost or decoded to garbage and the table can no longer be trusted to +// match the encoder's. The static table (indices 1-61) keeps working, and the +// decoder converges on the encoder's table again as new entries are inserted: +// see hpackDecoder. func (p *Http2Parser) resetDecoder(method Method) { switch method { case MethodHttp2ClientFrames: - p.clientDecoder = hpack.NewDecoder(4096, nil) + p.clientDecoder.reset() p.clientDecoderDegraded = true case MethodHttp2ServerFrames: - p.serverDecoder = hpack.NewDecoder(4096, nil) + p.serverDecoder.reset() p.serverDecoderDegraded = true } } +// 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) { + if *pending != nil { + *pending = nil + p.resetDecoder(method) + } +} + // SawValidFrame reports whether the last Parse call decoded at least one // structurally valid frame header. func (p *Http2Parser) SawValidFrame() bool { @@ -270,58 +279,63 @@ func extractHeaderBlockFragment(flags http2.Flags, framePayload []byte) []byte { } // decodeHeaderBlock processes a complete HPACK-encoded header block for a stream. -// It sets up the emit function, writes the HPACK data to the decoder, and handles errors. func (p *Http2Parser) decodeHeaderBlock( method Method, streamId uint32, endStream bool, hpackData []byte, - decoder *hpack.Decoder, + decoder *hpackDecoder, statuses map[uint32]Status, grpcStatuses map[uint32]Status, kernelTime uint64, ) { + // implausible is set when the block decodes to pseudo-headers that cannot + // be right for its direction: a sign the dynamic table has drifted from + // the encoder's because a block was lost unnoticed. + implausible := false + var emit func(name, value string) switch method { case MethodHttp2ClientFrames: req := p.activeRequests[streamId] - if req == nil { - if len(p.activeRequests) >= maxActiveRequests { - // Too many active streams; set no-op emit so decoder.Write still - // processes the HPACK block (keeps dynamic table in sync) without - // dereferencing a nil request. - decoder.SetEmitFunc(func(hf hpack.HeaderField) {}) - break - } + if req == nil && len(p.activeRequests) < maxActiveRequests { req = &Http2Request{ kernelTime: kernelTime, } p.activeRequests[streamId] = req p.stage("stream_created") } - decoder.SetEmitFunc(func(hf hpack.HeaderField) { - switch hf.Name { + // 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": - if req.Method == "" && isHttpMethod(hf.Value) { - req.Method = hf.Value + if !isHttpMethod(value) { + implausible = true + } else if req != nil && req.Method == "" { + req.Method = value } case ":path": - if req.Path == "" && isHttpPath(hf.Value) { - req.Path = hf.Value + if !isHttpPath(value) { + implausible = true + } else if req != nil && req.Path == "" { + req.Path = value } case ":scheme": - if req.Scheme == "" && isHttpScheme(hf.Value) { - req.Scheme = hf.Value + if req != nil && req.Scheme == "" && isHttpScheme(value) { + req.Scheme = value } case ":authority": - if req.Authority == "" && hf.Value != "" { - req.Authority = hf.Value + if req != nil && req.Authority == "" && value != "" { + req.Authority = value } case "content-type": - if req.ContentType == "" && hf.Value != "" { - req.ContentType = hf.Value + if req != nil && req.ContentType == "" && value != "" { + req.ContentType = value } + case ":status": + implausible = true } - }) + } case MethodHttp2ServerFrames: req := p.activeRequests[streamId] @@ -331,10 +345,14 @@ func (p *Http2Parser) decodeHeaderBlock( statuses[streamId] = 0 } } - decoder.SetEmitFunc(func(hf hpack.HeaderField) { - switch hf.Name { + emit = func(name, value string) { + switch name { case ":status": - s, _ := strconv.Atoi(hf.Value) + s, err := strconv.Atoi(value) + if err != nil || s < 100 || s > 999 { + implausible = true + return + } if req != nil { req.Status = Status(s) if !req.hasResponseStatus { @@ -344,13 +362,15 @@ func (p *Http2Parser) decodeHeaderBlock( } statuses[streamId] = Status(s) case "grpc-status": - s, _ := strconv.Atoi(hf.Value) + s, _ := strconv.Atoi(value) if req != nil { req.GrpcStatus = Status(s) } grpcStatuses[streamId] = Status(s) + case ":method", ":path", ":scheme", ":authority": + implausible = true } - }) + } // Check for END_STREAM flag on HEADERS (no body response) if req != nil && endStream { if !req.responseEndStream { @@ -358,32 +378,34 @@ func (p *Http2Parser) decodeHeaderBlock( } req.responseEndStream = true } + default: + return } - // Decode the complete HPACK header block. - // The emit function (set above) fires per-header, so headers decoded before any - // error are already stored in the request struct (e.g., :method from static index 3, - // :status from static index 8). We preserve these partial results on error. - if _, err := decoder.Write(hpackData); err != nil { - // HPACK decode error - commonly happens during mid-stream join when the agent - // starts monitoring after HTTP/2 connection was established. The remote encoder's - // dynamic table has entries our decoder doesn't have. - klog.V(3).Infof("http2: HPACK decode error on stream %d: %v (partial headers preserved)", streamId, err) - if OnHPACKDecodeError != nil { - OnHPACKDecodeError() - } - p.stage("hpack_error") - + // Fields are emitted as they decode, so on an error the ones before it + // (often :method or :status, from the static table) are kept. + unknown, err := decoder.decode(hpackData, emit) + if err != nil || unknown > 0 || implausible { // Mark the request as having partial headers so downstream can apply fallbacks if req := p.activeRequests[streamId]; req != nil { req.PartialHeaders = true } - - // Reset the decoder to prevent cascading failures. After a decode error, the - // decoder's internal buffer position and dynamic table are desynchronized. - // A fresh decoder starts with an empty dynamic table but static table (indices 1-61) - // always works. The dynamic table rebuilds from new literal-with-indexing headers. + } + switch { + case err != nil || implausible: + klog.V(3).Infof("http2: HPACK decode error on stream %d: %v, implausible=%v (partial headers preserved)", streamId, err, implausible) + if OnHPACKDecodeError != nil { + OnHPACKDecodeError() + } + p.stage("hpack_error") + // The block is not valid HPACK, or decoded to headers it cannot + // contain: either way the table no longer matches the encoder's. p.resetDecoder(method) + case unknown > 0: + // References to entries inserted before the decoder joined, or before + // it was reset. Expected, and not an error: the decoder catches up as + // the encoder inserts new entries. + p.stage("hpack_partial") } } @@ -410,7 +432,7 @@ func (p *Http2Parser) Parse(method Method, payload []byte, kernelTime uint64, mi return nil } - var decoder *hpack.Decoder + var decoder *hpackDecoder clear(p.statuses) clear(p.grpcStatuses) statuses := p.statuses @@ -597,8 +619,10 @@ frameLoop: // Validate: CONTINUATION must follow a HEADERS on the same stream if *pendingHeaders == nil || (*pendingHeaders).streamId != h.StreamId { - // Protocol error or we missed the HEADERS frame -- discard - *pendingHeaders = nil + // We missed the HEADERS frame (or this is a protocol error): + // the block it starts is lost. + p.dropPendingHeaders(method, pendingHeaders) + p.resetDecoder(method) continue } @@ -606,7 +630,7 @@ frameLoop: pending := *pendingHeaders if len(pending.fragments)+len(continuationPayload) > maxPendingHeaderBlockSize { // Too large, discard the pending header block - *pendingHeaders = nil + p.dropPendingHeaders(method, pendingHeaders) continue } pending.fragments = append(pending.fragments, continuationPayload...) @@ -654,12 +678,17 @@ frameLoop: if length >= captured+missing { *skip = length - captured - missing } + // A header block cut short never reaches the decoder, and the + // entries it inserted are missing from the table. + if t := http2.FrameType(payload[offset+3]); t == http2.FrameHeaders || t == http2.FrameContinuation { + p.resetDecoder(method) + } } *partialFrame = nil // A header block interrupted by truncation can never be completed by a // CONTINUATION frame, and feeding its fragments to the decoder later // would desync the dynamic table just as badly. - *pendingHeaders = nil + p.dropPendingHeaders(method, pendingHeaders) } else if offset < len(payload) { remaining := payload[offset:] // Only save if it looks like start of a valid frame (has at least some bytes) @@ -739,8 +768,8 @@ frameLoop: } } // Clear stale pending headers - p.clientPendingHeaders = nil - p.serverPendingHeaders = nil + p.dropPendingHeaders(MethodHttp2ClientFrames, &p.clientPendingHeaders) + p.dropPendingHeaders(MethodHttp2ServerFrames, &p.serverPendingHeaders) } p.lastGcTime = kernelTime } diff --git a/ebpftracer/l7/http2_hpack_test.go b/ebpftracer/l7/http2_hpack_test.go new file mode 100644 index 0000000..b633f6b --- /dev/null +++ b/ebpftracer/l7/http2_hpack_test.go @@ -0,0 +1,154 @@ +package l7 + +import ( + "bytes" + "fmt" + "testing" + + "golang.org/x/net/http2" + "golang.org/x/net/http2/hpack" +) + +// countStages records parser stages for the duration of a test. +func countStages(t *testing.T) map[string]int { + stages := map[string]int{} + prev := OnHttp2Stage + OnHttp2Stage = func(stage, dest string) { stages[stage]++ } + t.Cleanup(func() { OnHttp2Stage = prev }) + return stages +} + +func requestPath(i int) string { return fmt.Sprintf("/v1/items/%d", i%7) } + +// encodeRequests encodes n request header blocks on one connection. +func encodeRequests(n int) [][]byte { + var buf bytes.Buffer + enc := hpack.NewEncoder(&buf) + out := make([][]byte, n) + for i := range out { + buf.Reset() + for _, f := range []hpack.HeaderField{ + {Name: ":method", Value: "GET"}, + {Name: ":scheme", Value: "https"}, + {Name: ":authority", Value: "api.example.com"}, + {Name: ":path", Value: requestPath(i)}, + {Name: "user-agent", Value: "client/1.2.3"}, + {Name: "x-request-id", Value: fmt.Sprintf("req-%08d-%08d", i, i*7919)}, + } { + _ = enc.WriteField(f) + } + out[i] = append([]byte(nil), buf.Bytes()...) + } + return out +} + +func streamID(i int) uint32 { return uint32(2*i + 1) } + +// A parser that joins a connection mid-stream used to fail every header block +// for good (see TestHpackDecoderRecoversAfterJoiningMidStream); it now +// decodes method, path and authority again once the peer's table turns over, +// and never counts the references it cannot resolve as errors. +func TestHttp2ParserRecoversAfterJoiningMidStream(t *testing.T) { + stages := countStages(t) + const join, n = 20, 300 + blocks := encodeRequests(n) + p := NewHttp2Parser() + for i := join; i < n; i++ { + p.Parse(MethodHttp2ClientFrames, frame(http2.FrameHeaders, http2FlagEndHeaders|http2FlagEndStream, streamID(i), blocks[i]), uint64(i), 0) + req := p.activeRequests[streamID(i)] + if req == nil { + t.Fatalf("request %d: no stream", i) + } + if req.Path != "" && req.Path != requestPath(i) { + t.Fatalf("request %d: path %q, want %q", i, req.Path, requestPath(i)) + } + if i >= n-50 && (req.Method != "GET" || req.Path != requestPath(i) || req.Authority != "api.example.com") { + t.Fatalf("request %d not fully decoded: %+v", i, *req) + } + delete(p.activeRequests, streamID(i)) // as a response would + } + if stages["hpack_error"] != 0 { + t.Errorf("hpack_error = %d, want 0", stages["hpack_error"]) + } + if stages["hpack_partial"] == 0 { + t.Error("no hpack_partial: the join was not exercised") + } +} + +// A HEADERS frame cut short by truncation never reaches the decoder, so the +// entries it inserted are missing and every older index is off by their +// count. Decoding on would resolve later references to the wrong entries; +// the parser resets the table instead, and those references come back as +// unknown, never as wrong values. +func TestHttp2ParserResetsTableWhenAHeaderBlockIsLost(t *testing.T) { + const lost, n = 10, 200 + blocks := encodeRequests(n) + + // Without the reset, a decoder that misses the block decodes later + // references to the wrong entries. Check that first, or this test proves + // nothing. + d := newHpackDecoder() + wrong := 0 + for i := 0; i < n; i++ { + if i == lost { + continue + } + d.decode(blocks[i], func(name, value string) { + if name == ":path" && value != requestPath(i) { + wrong++ + } + }) + } + if wrong == 0 { + t.Fatal("losing a block did not misdecode anything; the scenario is too weak") + } + + stages := countStages(t) + p := NewHttp2Parser() + for i := 0; i < n; i++ { + f := frame(http2.FrameHeaders, http2FlagEndHeaders|http2FlagEndStream, streamID(i), blocks[i]) + if i == lost { + // The kernel captured only the first bytes of this write. + p.Parse(MethodHttp2ClientFrames, f[:http2FrameHeaderLength+2], uint64(i), uint64(len(f)-http2FrameHeaderLength-2)) + continue + } + p.Parse(MethodHttp2ClientFrames, f, uint64(i), 0) + req := p.activeRequests[streamID(i)] + if req == nil { + t.Fatalf("request %d: no stream", i) + } + if req.Path != "" && req.Path != requestPath(i) { + t.Fatalf("request %d: path %q, want %q", i, req.Path, requestPath(i)) + } + if req.Authority != "" && req.Authority != "api.example.com" { + t.Fatalf("request %d: authority %q", i, req.Authority) + } + if i >= n-50 && req.Path != requestPath(i) { + t.Fatalf("request %d: not recovered: %+v", i, *req) + } + delete(p.activeRequests, streamID(i)) + } + if stages["hpack_error"] != 0 { + t.Errorf("hpack_error = %d, want 0", stages["hpack_error"]) + } +} + +// Pseudo-headers that cannot appear in a block's direction mean the table +// has drifted (a block was lost unnoticed): the block counts as an error and +// the table is reset rather than trusted. +func TestHttp2ParserResetsTableOnImplausibleHeaders(t *testing.T) { + stages := countStages(t) + p := NewHttp2Parser() + p.serverDecoder.insert(hpackEntry{HeaderField: hpack.HeaderField{Name: "x-a", Value: "1"}}) + + var buf bytes.Buffer + _ = hpack.NewEncoder(&buf).WriteField(hpack.HeaderField{Name: ":method", Value: "GET"}) + p.Parse(MethodHttp2ServerFrames, frame(http2.FrameHeaders, http2FlagEndHeaders, 1, buf.Bytes()), 1, 0) + + if stages["hpack_error"] != 1 { + t.Errorf("hpack_error = %d, want 1", stages["hpack_error"]) + } + if len(p.serverDecoder.dynamic) != 0 { + t.Error("server table not reset") + } +} From 05641fc71b5f0a07b0c79dfa75729f450416d289 Mon Sep 17 00:00:00 2001 From: mayankpande88 Date: Tue, 6 Oct 2026 01:32:48 +0530 Subject: [PATCH 2/3] fix(l7): release evicted HPACK entries' strings Evicting, resetting or emptying the dynamic table resliced it, so the backing array kept the dropped entries' strings alive until overwritten. Clear them. Index with int after the bounds checks. --- ebpftracer/l7/hpack.go | 13 +++++++++---- 1 file changed, 9 insertions(+), 4 deletions(-) diff --git a/ebpftracer/l7/hpack.go b/ebpftracer/l7/hpack.go index 9a3fceb..6140afa 100644 --- a/ebpftracer/l7/hpack.go +++ b/ebpftracer/l7/hpack.go @@ -61,6 +61,7 @@ func newHpackDecoder() *hpackDecoder { // reset forgets the dynamic table, for when the caller knows a header block // was lost and the table no longer matches the encoder's. func (d *hpackDecoder) reset() { + clear(d.dynamic) // release the strings the backing array still holds d.dynamic = d.dynamic[:0] d.size = 0 d.maxSize = hpackDefaultTableSize @@ -149,13 +150,14 @@ func (d *hpackDecoder) field(idx uint64) (hpackEntry, bool) { if i > uint64(len(d.dynamic)) { return hpackEntry{}, false } - return d.dynamic[uint64(len(d.dynamic))-i], true + return d.dynamic[len(d.dynamic)-int(i)], true } func (d *hpackDecoder) insert(f hpackEntry) { size := uint32(len(f.Name)+len(f.Value)) + hpackEntryOverhead if size > d.maxSize { // RFC 7541 4.4: an entry larger than the table empties it. + clear(d.dynamic) d.dynamic = d.dynamic[:0] d.size = 0 return @@ -174,7 +176,10 @@ func (d *hpackDecoder) evict(room uint32) { n++ } if n > 0 { - d.dynamic = append(d.dynamic[:0], d.dynamic[n:]...) + copy(d.dynamic, d.dynamic[n:]) + // Release the evicted strings: the backing array keeps the tail. + clear(d.dynamic[len(d.dynamic)-n:]) + d.dynamic = d.dynamic[:len(d.dynamic)-n] } } @@ -216,8 +221,8 @@ func hpackString(b []byte) (string, []byte, error) { if n > uint64(len(b)) { return "", nil, errHpackTruncated } - raw := b[:n] - b = b[n:] + raw := b[:int(n)] + b = b[int(n):] if !huffman { return string(raw), b, nil } From 41057520cfe6aefcc25f2c12958082aa513d6774 Mon Sep 17 00:00:00 2001 From: mayankpande88 Date: Tue, 6 Oct 2026 02:09:04 +0530 Subject: [PATCH 3/3] test(l7): a table size update in a later header block decodes Go's HTTP/2 client sends one a few requests into every connection, once it has the server's SETTINGS. hpack.Decoder, fed with Write and never Close as the parser must, rejected it as not at the beginning of a header block, which put every Go client connection into the reset-on-error cascade. --- ebpftracer/l7/http2_hpack_test.go | 36 +++++++++++++++++++++++++++++++ 1 file changed, 36 insertions(+) diff --git a/ebpftracer/l7/http2_hpack_test.go b/ebpftracer/l7/http2_hpack_test.go index b633f6b..9243d58 100644 --- a/ebpftracer/l7/http2_hpack_test.go +++ b/ebpftracer/l7/http2_hpack_test.go @@ -152,3 +152,39 @@ func TestHttp2ParserResetsTableOnImplausibleHeaders(t *testing.T) { t.Error("server table not reset") } } + +// Go's HTTP/2 client starts a header block with a dynamic table size update +// once it has the server's SETTINGS, a few requests into every connection. +// hpack.Decoder, fed with Write and never Close as the parser must, does not +// know where a block starts and rejected that update as "not at the beginning +// of a header block": every Go client connection fell into the +// reset-on-error cascade within its first requests. +func TestHttp2ParserAcceptsTableSizeUpdateInLaterBlock(t *testing.T) { + stages := countStages(t) + var buf bytes.Buffer + enc := hpack.NewEncoder(&buf) + p := NewHttp2Parser() + for i := 0; i < 20; i++ { + buf.Reset() + if i == 3 { + enc.SetMaxDynamicTableSize(4096) // what the client does on SETTINGS + } + for _, f := range []hpack.HeaderField{ + {Name: ":authority", Value: "api.example.com"}, + {Name: ":method", Value: "GET"}, + {Name: ":path", Value: "/v1/items"}, + {Name: ":scheme", Value: "https"}, + {Name: "user-agent", Value: "Go-http-client/2.0"}, + } { + _ = enc.WriteField(f) + } + p.Parse(MethodHttp2ClientFrames, frame(http2.FrameHeaders, http2FlagEndHeaders|http2FlagEndStream, streamID(i), buf.Bytes()), uint64(i), 0) + req := p.activeRequests[streamID(i)] + if req == nil || req.Path != "/v1/items" || req.Authority != "api.example.com" { + t.Fatalf("request %d: %+v", i, req) + } + } + if stages["hpack_error"] != 0 || stages["hpack_partial"] != 0 { + t.Errorf("hpack_error = %d, hpack_partial = %d, want 0", stages["hpack_error"], stages["hpack_partial"]) + } +}