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

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
35 changes: 33 additions & 2 deletions pkg/engine/compat.go
Original file line number Diff line number Diff line change
Expand Up @@ -57,6 +57,7 @@ type rowBuilder struct {
lbsBuilder *labels.Builder
metadataBuilder *labels.Builder
parsedBuilder *labels.Builder
parsedEmptyKeys []string
}

func (b *streamsResultBuilder) CollectRecord(rec arrow.Record) {
Expand Down Expand Up @@ -140,11 +141,16 @@ func (b *streamsResultBuilder) CollectRecord(rec arrow.Record) {
if b.rowBuilders[rowIdx].parsedBuilder.Get(shortName) != "" {
return
}

b.rowBuilders[rowIdx].parsedBuilder.Set(shortName, parsedVal)
b.rowBuilders[rowIdx].lbsBuilder.Set(shortName, parsedVal)
if b.rowBuilders[rowIdx].metadataBuilder.Get(shortName) != "" {
b.rowBuilders[rowIdx].metadataBuilder.Del(shortName)
}
// If the parsed value is empty, the builder won't accept it as it's not a valid Prometheus-style label. We must add it later for LogQL compatibility.
if parsedVal == "" {
b.rowBuilders[rowIdx].parsedEmptyKeys = append(b.rowBuilders[rowIdx].parsedEmptyKeys, shortName)
}
})
}
}
Expand All @@ -160,16 +166,39 @@ func (b *streamsResultBuilder) CollectRecord(rec arrow.Record) {
continue
}

// For compatibility with LogQL, empty parsed labels need to be added to the stream labels & parsed label sets.
// The Prometheus label builder does not allow empty strings for label values, so we must work around it by creating a new builder and adding the empty labels to it.
var lbsString string
parsedLbs := logproto.FromLabelsToLabelAdapters(b.rowBuilders[rowIdx].parsedBuilder.Labels())
if len(b.rowBuilders[rowIdx].parsedEmptyKeys) > 0 {
newLbsBuilder := labels.NewScratchBuilder(lbs.Len())
lbs.Range(func(label labels.Label) {
newLbsBuilder.Add(label.Name, label.Value)
})

for _, key := range b.rowBuilders[rowIdx].parsedEmptyKeys {
newLbsBuilder.Add(key, "")
parsedLbs = append(parsedLbs, logproto.LabelAdapter{Name: key, Value: ""})
}
newLbsBuilder.Sort()
lbsString = newLbsBuilder.Labels().String()
sort.Slice(parsedLbs, func(i, j int) bool {
return parsedLbs[i].Name < parsedLbs[j].Name
})
} else {
lbsString = lbs.String()
}

entry := logproto.Entry{
Timestamp: ts,
Line: line,
StructuredMetadata: logproto.FromLabelsToLabelAdapters(b.rowBuilders[rowIdx].metadataBuilder.Labels()),
Parsed: logproto.FromLabelsToLabelAdapters(b.rowBuilders[rowIdx].parsedBuilder.Labels()),
Parsed: parsedLbs,
}
b.resetRowBuilder(rowIdx)

// Add entry to appropriate stream
key := lbs.String()
key := lbsString
idx, ok := b.streams[key]
if !ok {
idx = len(b.data)
Expand Down Expand Up @@ -203,6 +232,7 @@ func (b *streamsResultBuilder) ensureRowBuilders(newLen int) {
lbsBuilder: labels.NewBuilder(labels.EmptyLabels()),
metadataBuilder: labels.NewBuilder(labels.EmptyLabels()),
parsedBuilder: labels.NewBuilder(labels.EmptyLabels()),
parsedEmptyKeys: make([]string, 0),
}
}
}
Expand All @@ -213,6 +243,7 @@ func (b *streamsResultBuilder) resetRowBuilder(i int) {
b.rowBuilders[i].lbsBuilder.Reset(labels.EmptyLabels())
b.rowBuilders[i].metadataBuilder.Reset(labels.EmptyLabels())
b.rowBuilders[i].parsedBuilder.Reset(labels.EmptyLabels())
b.rowBuilders[i].parsedEmptyKeys = b.rowBuilders[i].parsedEmptyKeys[:0]
}

func forEachNotNullRowColValue(numRows int, col arrow.Array, f func(rowIdx int)) {
Expand Down
53 changes: 53 additions & 0 deletions pkg/engine/compat_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -352,6 +352,7 @@ func TestStreamsResultBuilder(t *testing.T) {
require.NotNil(t, builder.rowBuilders[i].lbsBuilder, "lbsBuilder should be initialized")
require.NotNil(t, builder.rowBuilders[i].metadataBuilder, "metadataBuilder should be initialized")
require.NotNil(t, builder.rowBuilders[i].parsedBuilder, "parsedBuilder should be initialized")
require.Equal(t, 0, len(builder.rowBuilders[i].parsedEmptyKeys), "parsedEmptyKeys should be empty")
}
})

Expand Down Expand Up @@ -420,6 +421,58 @@ func TestStreamsResultBuilder(t *testing.T) {
}
require.Equal(t, 5, totalEntries, "should have 5 total entries across both streams")
})

t.Run("parsed empty values are added to the parsed labels", func(t *testing.T) {
colTs := semconv.ColumnIdentTimestamp
colMsg := semconv.ColumnIdentMessage
colEnv := semconv.NewIdentifier("env", types.ColumnTypeLabel, types.Loki.String)
colMetadata := semconv.NewIdentifier("metadata", types.ColumnTypeMetadata, types.Loki.String)
colParsedA := semconv.NewIdentifier("Aparsed", types.ColumnTypeParsed, types.Loki.String)
colParsedZ := semconv.NewIdentifier("Zparsed", types.ColumnTypeParsed, types.Loki.String)

schema := arrow.NewSchema(
[]arrow.Field{
semconv.FieldFromIdent(colTs, false),
semconv.FieldFromIdent(colMsg, false),
semconv.FieldFromIdent(colEnv, false),
semconv.FieldFromIdent(colMetadata, false),
semconv.FieldFromIdent(colParsedA, false),
semconv.FieldFromIdent(colParsedZ, false),
},
nil,
)
rows := arrowtest.Rows{
{colTs.FQN(): time.Unix(0, 1620000000000000000).UTC(), colMsg.FQN(): "log line", colEnv.FQN(): "prod", colMetadata.FQN(): "md value", colParsedA.FQN(): "A", colParsedZ.FQN(): "Z"},
{colTs.FQN(): time.Unix(0, 1620000000000000000).UTC(), colMsg.FQN(): "log line", colEnv.FQN(): "prod", colMetadata.FQN(): "", colParsedA.FQN(): "", colParsedZ.FQN(): ""},
}

record := rows.Record(memory.DefaultAllocator, schema)
builder := newStreamsResultBuilder(logproto.BACKWARD)
builder.CollectRecord(record)
record.Release()
require.Equal(t, 2, builder.Len())

md, _ := metadata.NewContext(t.Context())
result := builder.Build(stats.Result{}, md)
streams := result.Data.(logqlmodel.Streams)
require.Equal(t, 2, len(streams), "should have 2 unique streams")

expected := logqlmodel.Streams{
push.Stream{
Labels: labels.FromStrings("Aparsed", "", "env", "prod", "Zparsed", "").String(),
Entries: []logproto.Entry{
{Line: "log line", Timestamp: time.Unix(0, 1620000000000000000), StructuredMetadata: logproto.FromLabelsToLabelAdapters(labels.Labels{}), Parsed: logproto.FromLabelsToLabelAdapters(labels.FromStrings("Aparsed", "", "Zparsed", ""))},
},
},
push.Stream{
Labels: labels.FromStrings("Aparsed", "A", "env", "prod", "metadata", "md value", "Zparsed", "Z").String(),
Entries: []logproto.Entry{
{Line: "log line", Timestamp: time.Unix(0, 1620000000000000000), StructuredMetadata: logproto.FromLabelsToLabelAdapters(labels.FromStrings("metadata", "md value")), Parsed: logproto.FromLabelsToLabelAdapters(labels.FromStrings("Aparsed", "A", "Zparsed", "Z"))},
},
},
}
require.Equal(t, expected, streams)
})
}

func TestVectorResultBuilder(t *testing.T) {
Expand Down
17 changes: 17 additions & 0 deletions pkg/engine/internal/executor/expressions_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -1161,6 +1161,23 @@ func TestEvaluateParseExpression_JSON(t *testing.T) {
},
},
},
{
name: "keys with empty values as empty strings",
schema: arrow.NewSchema([]arrow.Field{
semconv.FieldFromFQN("utf8.builtin.message", true),
}, nil),
input: arrowtest.Rows{
{colMsg: `{"level": "error", "namespace": ""}`},
},
requestedKeys: nil,
expectedFields: 2, // 2 columns: level, namespace
expectedOutput: arrowtest.Rows{
{
"utf8.parsed.level": "error",
"utf8.parsed.namespace": "",
},
},
},
} {
t.Run(tt.name, func(t *testing.T) {
alloc := memory.DefaultAllocator
Expand Down
14 changes: 7 additions & 7 deletions pkg/engine/internal/executor/parse_json.go
Original file line number Diff line number Diff line change
Expand Up @@ -140,14 +140,14 @@ func (j *jsonParser) parseLabelValue(key, value []byte, dataType jsonparser.Valu

// Convert the value to string based on its type
parsedValue := parseValue(value, dataType)
if parsedValue != "" {
// First-wins semantics for duplicates
_, exists := result[keyString]
if exists {
return nil
}
result[keyString] = parsedValue

// Empty keys are always kept for json
// First-wins semantics for duplicates
_, exists := result[keyString]
if exists {
return nil
}
result[keyString] = parsedValue

return nil
}
Expand Down