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..6140afa --- /dev/null +++ b/ebpftracer/l7/hpack.go @@ -0,0 +1,299 @@ +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() { + clear(d.dynamic) // release the strings the backing array still holds + 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[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 + } + 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 { + 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] + } +} + +// 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[:int(n)] + b = b[int(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..9243d58 --- /dev/null +++ b/ebpftracer/l7/http2_hpack_test.go @@ -0,0 +1,190 @@ +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") + } +} + +// 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"]) + } +}