Loki: backend: use json-field (#48486)
* fixed strings * loki: backend: use json-field
This commit is contained in:
@@ -1,6 +1,7 @@
|
||||
package loki
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"hash/fnv"
|
||||
"sort"
|
||||
@@ -88,12 +89,12 @@ func adjustLogsFrame(frame *data.Frame, query *lokiQuery) error {
|
||||
timeField := fields[1]
|
||||
lineField := fields[2]
|
||||
|
||||
if (timeField.Type() != data.FieldTypeTime) || (lineField.Type() != data.FieldTypeString) || (labelsField.Type() != data.FieldTypeString) {
|
||||
return fmt.Errorf("invalid fields in metric frame")
|
||||
if (timeField.Type() != data.FieldTypeTime) || (lineField.Type() != data.FieldTypeString) || (labelsField.Type() != data.FieldTypeJSON) {
|
||||
return fmt.Errorf("invalid fields in logs frame")
|
||||
}
|
||||
|
||||
if (timeField.Len() != lineField.Len()) || (timeField.Len() != labelsField.Len()) {
|
||||
return fmt.Errorf("invalid fields in metric frame")
|
||||
return fmt.Errorf("invalid fields in logs frame")
|
||||
}
|
||||
|
||||
if frame.Meta == nil {
|
||||
@@ -127,8 +128,9 @@ func makeStringTimeField(timeField *data.Field) *data.Field {
|
||||
return data.NewField("tsNs", timeField.Labels.Copy(), stringTimestamps)
|
||||
}
|
||||
|
||||
func calculateCheckSum(time string, line string, labels string) (string, error) {
|
||||
input := []byte(line + "_" + labels)
|
||||
func calculateCheckSum(time string, line string, labels []byte) (string, error) {
|
||||
input := []byte(line + "_")
|
||||
input = append(input, labels...)
|
||||
hash := fnv.New32()
|
||||
_, err := hash.Write(input)
|
||||
if err != nil {
|
||||
@@ -147,7 +149,7 @@ func makeIdField(stringTimeField *data.Field, lineField *data.Field, labelsField
|
||||
for i := 0; i < length; i++ {
|
||||
time := stringTimeField.At(i).(string)
|
||||
line := lineField.At(i).(string)
|
||||
labels := labelsField.At(i).(string)
|
||||
labels := labelsField.At(i).(json.RawMessage)
|
||||
|
||||
sum, err := calculateCheckSum(time, line, labels)
|
||||
if err != nil {
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
package loki
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
@@ -39,11 +40,11 @@ func TestFormatName(t *testing.T) {
|
||||
func TestAdjustFrame(t *testing.T) {
|
||||
t.Run("logs-frame metadata should be set correctly", func(t *testing.T) {
|
||||
frame := data.NewFrame("",
|
||||
data.NewField("labels", nil, []string{
|
||||
`{"level":"info"}`,
|
||||
`{"level":"error"}`,
|
||||
`{"level":"error"}`,
|
||||
`{"level":"info"}`,
|
||||
data.NewField("labels", nil, []json.RawMessage{
|
||||
json.RawMessage(`{"level":"info"}`),
|
||||
json.RawMessage(`{"level":"error"}`),
|
||||
json.RawMessage(`{"level":"error"}`),
|
||||
json.RawMessage(`{"level":"info"}`),
|
||||
}),
|
||||
data.NewField("time", nil, []time.Time{
|
||||
time.Date(2022, 1, 2, 3, 4, 5, 6, time.UTC),
|
||||
@@ -86,7 +87,7 @@ func TestAdjustFrame(t *testing.T) {
|
||||
require.Equal(t, "1641092765000000006_948c1a7d_A", idField.At(3))
|
||||
})
|
||||
|
||||
t.Run("logs-frame id and string-time fields should be created", func(t *testing.T) {
|
||||
t.Run("naming inside metric fields should be correct", func(t *testing.T) {
|
||||
field1 := data.NewField("", nil, make([]time.Time, 0))
|
||||
field2 := data.NewField("", nil, make([]float64, 0))
|
||||
field2.Labels = data.Labels{"app": "Application", "tag2": "tag2"}
|
||||
|
||||
@@ -1,8 +1,8 @@
|
||||
package loki
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"sort"
|
||||
"time"
|
||||
|
||||
"github.com/grafana/grafana-plugin-sdk-go/data"
|
||||
@@ -109,35 +109,23 @@ func lokiVectorToDataFrames(vector loghttp.Vector, query *lokiQuery, stats []dat
|
||||
}
|
||||
|
||||
// we serialize the labels as an ordered list of pairs
|
||||
func labelsToString(labels data.Labels) (string, error) {
|
||||
keys := make([]string, 0, len(labels))
|
||||
for k := range labels {
|
||||
keys = append(keys, k)
|
||||
}
|
||||
sort.Strings(keys)
|
||||
|
||||
labelArray := make([][2]string, 0, len(labels))
|
||||
|
||||
for _, k := range keys {
|
||||
pair := [2]string{k, labels[k]}
|
||||
labelArray = append(labelArray, pair)
|
||||
}
|
||||
|
||||
bytes, err := jsoniter.Marshal(labelArray)
|
||||
func labelsToRawJson(labels data.Labels) (json.RawMessage, error) {
|
||||
// data.Labels when converted to JSON keep the fields sorted
|
||||
bytes, err := jsoniter.Marshal(labels)
|
||||
if err != nil {
|
||||
return "", err
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return string(bytes), nil
|
||||
return json.RawMessage(bytes), nil
|
||||
}
|
||||
|
||||
func lokiStreamsToDataFrames(streams loghttp.Streams, query *lokiQuery, stats []data.QueryStat) (data.Frames, error) {
|
||||
var timeVector []time.Time
|
||||
var values []string
|
||||
var labelsVector []string
|
||||
var labelsVector []json.RawMessage
|
||||
|
||||
for _, v := range streams {
|
||||
labelsText, err := labelsToString(v.Labels.Map())
|
||||
labelsJson, err := labelsToRawJson(v.Labels.Map())
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
@@ -145,17 +133,13 @@ func lokiStreamsToDataFrames(streams loghttp.Streams, query *lokiQuery, stats []
|
||||
for _, k := range v.Entries {
|
||||
timeVector = append(timeVector, k.Timestamp.UTC())
|
||||
values = append(values, k.Line)
|
||||
labelsVector = append(labelsVector, labelsText)
|
||||
labelsVector = append(labelsVector, labelsJson)
|
||||
}
|
||||
}
|
||||
|
||||
timeField := data.NewField(data.TimeSeriesTimeFieldName, nil, timeVector)
|
||||
valueField := data.NewField("Line", nil, values)
|
||||
labelsField := data.NewField("labels", nil, labelsVector)
|
||||
labelsField.Config = &data.FieldConfig{
|
||||
// we should have a native json-field-type
|
||||
Custom: map[string]interface{}{"json": true},
|
||||
}
|
||||
|
||||
frame := data.NewFrame("", labelsField, timeField, valueField)
|
||||
frame.SetMeta(&data.FrameMeta{
|
||||
|
||||
+13
-13
File diff suppressed because one or more lines are too long
Reference in New Issue
Block a user