Skip to content
Open
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
223 changes: 151 additions & 72 deletions pkg/clickhouse/ingester.go
Original file line number Diff line number Diff line change
Expand Up @@ -55,6 +55,11 @@ func (i *Ingester) Ingest(ctx context.Context, record arrow.RecordBatch) error {

schema := record.Schema()

// Counts locations whose encoded bytes could not be fully decoded, so a
// systematic encoder/decoder drift is visible instead of just producing
// profiles with unnamed frames.
malformedLocations := 0

// Find column indices
nameIdx := findColumnIndex(schema, profile.ColumnName)
sampleTypeIdx := findColumnIndex(schema, profile.ColumnSampleType)
Expand Down Expand Up @@ -102,7 +107,8 @@ func (i *Ingester) Ingest(ctx context.Context, record arrow.RecordBatch) error {
}

// Extract stacktrace data
stacktraceData := extractStacktraceData(record, stacktraceIdx, row)
stacktraceData, malformed := extractStacktraceData(record, stacktraceIdx, row)
malformedLocations += malformed

// Append to batch
err := batch.Append(
Expand Down Expand Up @@ -135,6 +141,16 @@ func (i *Ingester) Ingest(ctx context.Context, record arrow.RecordBatch) error {
}
}

// One line per batch, not per location: a drift affects every location in
// the batch, and logging each one would bury the signal it is meant to be.
if malformedLocations > 0 {
level.Warn(i.logger).Log(
"msg", "could not fully decode some encoded locations; their frames are stored without function info",
"locations", malformedLocations,
"rows", record.NumRows(),
)
}

if err := batch.Send(); err != nil {
return fmt.Errorf("failed to send batch: %w", err)
}
Expand Down Expand Up @@ -167,84 +183,144 @@ type LineInfo struct {
}

// decodeLineInfo decodes line and function information from the encoded location data.
// It returns the first line's info (most profiles have one line per location).
func decodeLineInfo(data []byte) LineInfo {
var n int
//
// Every read is bounds-checked, and a short or malformed record yields whatever
// had been decoded before the record ran out rather than panicking. The bytes
// are produced by this server's own encoders and so are self-consistent in
// normal operation, which is exactly why the failure mode matters: the realistic
// way to get a malformed record is encoder/decoder drift, and there is no
// recovery interceptor on the ingest path, so an unchecked index there takes the
// process down instead of failing one request.
func decodeLineInfo(data []byte) (LineInfo, bool) {
info := LineInfo{}
offset := 0

// Skip addr
_, offset := varint.Uvarint(data)

// Read number of lines
numLines, n := varint.Uvarint(data[offset:])
offset += n

// Check if has mapping
hasMapping := data[offset] == 0x1
offset++

if hasMapping {
// Skip buildID
length, n := varint.Uvarint(data[offset:])
offset += n + int(length)

// Skip filename
length, n = varint.Uvarint(data[offset:])
offset += n + int(length)

// Skip memoryStart
_, n = varint.Uvarint(data[offset:])
offset += n

// Skip memoryLength
_, n = varint.Uvarint(data[offset:])
// uvarint reports false when the record has run out or the varint is
// malformed; Uvarint returns n <= 0 for both.
uvarint := func() (uint64, bool) {
if offset >= len(data) {
return 0, false
}
v, n := varint.Uvarint(data[offset:])
if n <= 0 {
return 0, false
}
offset += n
return v, true
}
// str reads a length-prefixed string.
str := func() (string, bool) {
length, ok := uvarint()
// The length is unsigned, so one larger than MaxInt converts to a
// negative int and offset+int(length) lands BELOW offset -- which slips
// past a naive offset+int(length) > len(data) check straight into a
// panicking slice. Compare in the space the length was read in, against
// the bytes that actually remain.
if !ok || length > uint64(len(data)-offset) {
return "", false
}
v := string(data[offset : offset+int(length)])
offset += int(length)
return v, true
}
flag := func() (bool, bool) {
if offset >= len(data) {
return false, false
}
v := data[offset] == 0x1
offset++
return v, true
}
skip := func(n int) bool {
for range n {
if _, ok := uvarint(); !ok {
return false
}
}
return true
}

// Skip mappingOffset
_, n = varint.Uvarint(data[offset:])
offset += n
if _, ok := uvarint(); !ok { // address
return info, false
}
numLines, ok := uvarint()
if !ok {
return info, false
}
hasMapping, ok := flag()
if !ok {
return info, false
}
if hasMapping {
if _, ok := str(); !ok { // buildID
return info, false
}
if _, ok := str(); !ok { // filename
return info, false
}
// memoryStart, memoryLength, mappingOffset
if !skip(3) {
return info, false
}
}

if numLines > 0 {
// Read first line info (we only store one line per location)
lineNum, n := varint.Uvarint(data[offset:])
offset += n
info.LineNumber = int64(lineNum)
// A location with no lines is a valid shape, not a malformed record.
if numLines == 0 {
return info, true
}

hasFunction := data[offset] == 0x1
offset++
// Only the first line is kept: the schema stores one function per location.
// Location.Line[0] is the innermost inlined function, which is the right one
// to keep, but every inlined caller above it is dropped here.
lineNum, ok := uvarint()
if !ok {
return info, false
}
info.LineNumber = int64(lineNum)

// The column. pprof carries no column information, so EncodePprofLocation
// writes a uvarint zero here -- a single 0x00 byte. Leaving it unread makes
// the hasFunction read below land on the column instead of the flag, where
// it is always false, which silently discards the function name, system
// name, filename and start line of every already-symbolized location.
if _, ok := uvarint(); !ok {
return info, false
}

if hasFunction {
// Read startLine
startLine, n := varint.Uvarint(data[offset:])
offset += n
info.FunctionStartLine = int64(startLine)

// Read function name
length, n := varint.Uvarint(data[offset:])
offset += n
info.FunctionName = string(data[offset : offset+int(length)])
offset += int(length)

// Read system name
length, n = varint.Uvarint(data[offset:])
offset += n
info.FunctionSystemName = string(data[offset : offset+int(length)])
offset += int(length)

// Read filename
length, n = varint.Uvarint(data[offset:])
offset += n
info.FunctionFilename = string(data[offset : offset+int(length)])
}
hasFunction, ok := flag()
if !ok {
return info, false
}
// A line with no function is a valid shape too.
if !hasFunction {
return info, true
}

return info
startLine, ok := uvarint()
if !ok {
return info, false
}
info.FunctionStartLine = int64(startLine)
if info.FunctionName, ok = str(); !ok {
return info, false
}
if info.FunctionSystemName, ok = str(); !ok {
return info, false
}
info.FunctionFilename, ok = str()
return info, ok
}

// extractStacktraceData extracts stacktrace information from the encoded binary column.
// The stacktrace column contains encoded location data that needs to be decoded.
func extractStacktraceData(record arrow.RecordBatch, colIdx, row int) StacktraceData {
// The second return value counts locations whose encoded bytes could not be
// fully decoded. They are still written, with whatever was recovered, so one
// bad location does not discard the rest of the batch -- but the count is
// reported by the caller, because a decoder that silently degrades every
// profile to unnamed frames is indistinguishable from legitimately
// unsymbolized ones.
func extractStacktraceData(record arrow.RecordBatch, colIdx, row int) (StacktraceData, int) {
malformed := 0
data := StacktraceData{
Addresses: []uint64{},
MappingStarts: []uint64{},
Expand All @@ -260,30 +336,30 @@ func extractStacktraceData(record arrow.RecordBatch, colIdx, row int) Stacktrace
}

if colIdx < 0 {
return data
return data, 0
}

col := record.Column(colIdx)
listCol, ok := col.(*array.List)
if !ok {
return data
return data, 0
}

if listCol.IsNull(row) {
return data
return data, 0
}

start, end := listCol.ValueOffsets(row)
values := listCol.ListValues()

dictCol, ok := values.(*array.Dictionary)
if !ok {
return data
return data, 0
}

binaryDict, ok := dictCol.Dictionary().(*array.Binary)
if !ok {
return data
return data, 0
}

for idx := int(start); idx < int(end); idx++ {
Expand All @@ -305,15 +381,18 @@ func extractStacktraceData(record arrow.RecordBatch, colIdx, row int) Stacktrace
data.MappingBuildIDs = append(data.MappingBuildIDs, string(symInfo.BuildID))

// Decode line/function info
lineInfo := decodeLineInfo(encodedLocation)
lineInfo, ok := decodeLineInfo(encodedLocation)
if !ok {
malformed++
}
data.LineNumbers = append(data.LineNumbers, lineInfo.LineNumber)
data.FunctionNames = append(data.FunctionNames, lineInfo.FunctionName)
data.FunctionSystemNames = append(data.FunctionSystemNames, lineInfo.FunctionSystemName)
data.FunctionFilenames = append(data.FunctionFilenames, lineInfo.FunctionFilename)
data.FunctionStartLines = append(data.FunctionStartLines, lineInfo.FunctionStartLine)
}

return data
return data, malformed
}

func findColumnIndex(schema *arrow.Schema, name string) int {
Expand Down
Loading