Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions containers/container.go
Original file line number Diff line number Diff line change
Expand Up @@ -1387,6 +1387,11 @@ func (c *Container) onL7RequestWithResult(pid uint32, fd uint64, timestamp uint6
if r.PayloadSize > uint64(len(r.Payload)) {
missing = r.PayloadSize - uint64(len(r.Payload))
}
// Events of this connection the kernel could not deliver (its ring
// buffer was full) leave a gap the parser cannot see in the bytes.
if r.LostBefore != 0 {
parser.Lost(r.LostBefore)
}
requests := parser.Parse(r.Method, r.Payload, r.KernelTime, missing)

// HTTP/2 has the weakest detection heuristic of any protocol here — it
Expand Down
2 changes: 2 additions & 0 deletions containers/l7_self_metrics.go
Original file line number Diff line number Diff line change
Expand Up @@ -133,6 +133,8 @@ 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
// events_lost events of the connection were lost before this one
// (full ring buffer); that direction was resynchronized
// 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
Expand Down
20 changes: 10 additions & 10 deletions ebpftracer/ebpf.go

Large diffs are not rendered by default.

26 changes: 25 additions & 1 deletion ebpftracer/ebpf/l7/l7.c
Original file line number Diff line number Diff line change
Expand Up @@ -98,7 +98,7 @@ struct l7_event {
__u8 protocol;
__u8 method;
__u8 is_tls; // payload came from a TLS library hook, i.e. it is plaintext
__u8 padding;
__u8 lost_before; // connection.l7_lost when this event was sent
__u32 statement_id;
__u64 payload_size;
__u64 response_size;
Expand Down Expand Up @@ -294,12 +294,28 @@ void send_event(void *ctx, struct l7_event *e, struct connection_id cid, struct
if (len > sizeof(struct l7_event)) {
len = sizeof(struct l7_event);
}
// A record with a response carries the whole payload buffer, so the bytes
// between payload_size and MAX_PAYLOAD_SIZE are left over from an earlier
// event on this CPU. They are plaintext the agent already received, the
// ring is readable only by the agent, and the decoder copies exactly
// payload_size bytes; zeroing 4 KB per event would cost more than it buys.
//
// A lost event leaves a gap in its connection's stream that a stateful
// parser cannot see on its own: HTTP/2 frames are cut mid-way and the
// HPACK table misses insertions, so later headers decode to wrong values.
// Record the loss on the connection and hand it to the next event that is
// delivered. The read-modify-write is not atomic (atomic OR needs 5.12);
// two CPUs racing on one connection can lose a flag, rarely.
e->lost_before = conn->l7_lost;
Comment thread
mayankpande88 marked this conversation as resolved.
if (bpf_ringbuf_output(&l7_events, e, len, 0)) {
__u32 zero = 0;
__u64 *drops = bpf_map_lookup_elem(&l7_ringbuf_drops, &zero);
if (drops) {
*drops += 1;
}
conn->l7_lost |= e->method == METHOD_HTTP2_SERVER_FRAMES ? 2 : 1;
} else if (e->lost_before) {
conn->l7_lost = 0;
}
}

Expand Down Expand Up @@ -565,6 +581,14 @@ int trace_enter_write(void *ctx, __u64 fd, __u16 is_tls, char *buf, __u64 size,
// here on.
conn_on_stack.timestamp = bpf_ktime_get_ns();
bpf_map_update_elem(&active_connections, &cid, &conn_on_stack, BPF_NOEXIST);
// Continue on the map entry (another CPU's, if it won the race):
// the protocol detected below, and a loss recorded by
// send_event, would otherwise be written to this stack copy and
// discarded, and the next write would have to detect afresh.
struct connection *tracked = bpf_map_lookup_elem(&active_connections, &cid);
if (tracked) {
conn = tracked;
}
}
} else if (conn->tls) {
if (is_tls_clienthello(payload, size)) {
Expand Down
5 changes: 5 additions & 0 deletions ebpftracer/ebpf/tcp/state.c
Original file line number Diff line number Diff line change
Expand Up @@ -80,6 +80,11 @@ struct connection {
// ciphertext for it; see mark_tls in l7/l7.c.
__u8 tls;
__u16 dport; // Destination port for protocol detection
// l7_lost records L7 events of this connection that l7_events had no room
// for since its last delivered one: bit 0 written frames, bit 1 read
// frames. The next delivered event carries it (l7_event.lost_before) and
// clears it. Stored in what was tail padding; the struct stays 32 bytes.
__u8 l7_lost;
};

struct {
Expand Down
66 changes: 40 additions & 26 deletions ebpftracer/l7/hpack.go
Original file line number Diff line number Diff line change
Expand Up @@ -26,9 +26,12 @@ import (
// 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
// ring holds the n entries the decoder knows, oldest at head. A ring
// makes eviction O(1): shifting a slice per eviction cost ~34 ms on one
// 64 KB block of minimal insertions after a size update to 64 KiB.
ring []hpackEntry
head, n int
size uint32 // RFC 7541 4.1 size of the known entries
maxSize uint32 // set by dynamic table size updates
}

Expand Down Expand Up @@ -61,12 +64,21 @@ 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.empty()
d.maxSize = hpackDefaultTableSize
}

// empty drops every entry, releasing their strings.
func (d *hpackDecoder) empty() {
clear(d.ring)
d.head, d.n, d.size = 0, 0, 0
}

// entry returns the i-th known entry, 0 being the oldest.
func (d *hpackDecoder) entry(i int) *hpackEntry {
return &d.ring[(d.head+i)%len(d.ring)]
}

// 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
Expand Down Expand Up @@ -94,10 +106,11 @@ func (d *hpackDecoder) decode(block []byte, emit func(name, value string)) (unkn
if size, block, err = hpackInt(block, 5); err != nil {
return unknown, err
}
if size > hpackMaxTableSize {
return unknown, errHpackInvalid
}
d.maxSize = uint32(size)
// The decoder only mirrors the table; clamping an outsized
// update keeps decoding the rest of the block, and at worst
// evicts entries the encoder still has, which then read as
// unknown rather than wrong.
d.maxSize = uint32(min(size, hpackMaxTableSize))
d.evict(0)
default: // 6.2 Literal Header Field
prefix, index := uint8(4), false
Expand Down Expand Up @@ -147,39 +160,40 @@ func (d *hpackDecoder) field(idx uint64) (hpackEntry, bool) {
return hpackEntry{HeaderField: hpackStaticTable[idx-1]}, true
}
i := idx - uint64(len(hpackStaticTable)) // 1 = newest
if i > uint64(len(d.dynamic)) {
if i > uint64(d.n) {
return hpackEntry{}, false
}
return d.dynamic[len(d.dynamic)-int(i)], true
return *d.entry(d.n - 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
d.empty()
return
}
d.evict(size)
d.dynamic = append(d.dynamic, f)
if d.n == len(d.ring) {
grown := make([]hpackEntry, max(2*len(d.ring), 16))
for i := 0; i < d.n; i++ {
grown[i] = *d.entry(i)
}
d.ring, d.head = grown, 0
}
*d.entry(d.n) = f
d.n++
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]
for d.size+room > d.maxSize && d.n > 0 {
f := d.entry(0)
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]
*f = hpackEntry{} // release its strings
d.head = (d.head + 1) % len(d.ring)
d.n--
}
}

Expand Down
41 changes: 40 additions & 1 deletion ebpftracer/l7/hpack_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -245,7 +245,6 @@ func TestHpackDecoderRejectsMalformedBlocks(t *testing.T) {
"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 {
Expand Down Expand Up @@ -290,3 +289,43 @@ func BenchmarkHpackDecoder(b *testing.B) {
}
})
}

// A size update above what the decoder will hold is clamped, not an error
// that would lose the rest of the block.
func TestHpackDecoderClampsOversizedTableUpdate(t *testing.T) {
block := []byte{0x3f, 0xe1, 0xff, 0x7f} // size update to ~2 MiB
var buf bytes.Buffer
_ = hpack.NewEncoder(&buf).WriteField(hpack.HeaderField{Name: "x-a", Value: "1"})
block = append(block, buf.Bytes()...)
d := newHpackDecoder()
got, unknown := decodeAll(t, d, block)
if unknown != 0 || len(got) != 1 || got[0].Value != "1" {
t.Fatalf("got %v, %d unknown", got, unknown)
}
if d.maxSize != hpackMaxTableSize {
t.Errorf("maxSize = %d, want %d", d.maxSize, hpackMaxTableSize)
}
}

// A 64 KB block of minimal insertions after a size update to 64 KiB evicts
// on every insertion once the table is full. With a slice that shifted per
// eviction this took ~34 ms; the ring makes it linear.
func manySmallInsertions() []byte {
var b bytes.Buffer
b.Write([]byte{0x3f, 0xe1, 0xff, 0x03}) // size update to 64 KiB
for b.Len() < 64*1024 {
b.Write([]byte{0x40, 0x01, 'a', 0x00}) // literal with indexing, name "a", empty value
}
return b.Bytes()
}

func BenchmarkHpackDecoderManySmallInsertions(b *testing.B) {
block := manySmallInsertions()
emit := func(string, string) {}
b.ReportAllocs()
for i := 0; i < b.N; i++ {
if _, err := newHpackDecoder().decode(block, emit); err != nil {
b.Fatal(err)
}
}
}
Loading
Loading