Influx: Support flux in the influx datasource (#25308)
* add flux * add token to datasource config editor * add backend for flux * make the interpolated query available in query inspector * go mod tidy * Chore: fixes a couple of strict null errors in influxdb plugin Co-authored-by: kyle <kyle@grafana.com> Co-authored-by: Lukas Siatka <lukasz.siatka@grafana.com>
This commit is contained in:
co-authored by
kyle
Lukas Siatka
parent
c7aac1fd40
commit
5f1f820bb9
@@ -0,0 +1,182 @@
|
||||
package flux
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
|
||||
"github.com/grafana/grafana-plugin-sdk-go/data"
|
||||
influxdb2 "github.com/influxdata/influxdb-client-go"
|
||||
)
|
||||
|
||||
// Copied from: (Apache 2 license)
|
||||
// https://github.com/influxdata/influxdb-client-go/blob/master/query.go#L30
|
||||
const (
|
||||
stringDatatype = "string"
|
||||
doubleDatatype = "double"
|
||||
boolDatatype = "bool"
|
||||
longDatatype = "long"
|
||||
uLongDatatype = "unsignedLong"
|
||||
durationDatatype = "duration"
|
||||
base64BinaryDataType = "base64Binary"
|
||||
timeDatatypeRFC = "dateTime:RFC3339"
|
||||
timeDatatypeRFCNano = "dateTime:RFC3339Nano"
|
||||
)
|
||||
|
||||
type columnInfo struct {
|
||||
name string
|
||||
converter *data.FieldConverter
|
||||
}
|
||||
|
||||
// This is an interface to help testing
|
||||
type FrameBuilder struct {
|
||||
tableId int64
|
||||
active *data.Frame
|
||||
frames []*data.Frame
|
||||
value *data.FieldConverter
|
||||
columns []columnInfo
|
||||
labels []string
|
||||
maxPoints int // max points in a series
|
||||
maxSeries int // max number of series
|
||||
totalSeries int
|
||||
isTimeSeries bool
|
||||
}
|
||||
|
||||
func isTag(schk string) bool {
|
||||
return (schk != "result" && schk != "table" && schk[0] != '_')
|
||||
}
|
||||
|
||||
func getConverter(t string) (*data.FieldConverter, error) {
|
||||
switch t {
|
||||
case stringDatatype:
|
||||
return &AnyToOptionalString, nil
|
||||
case timeDatatypeRFC:
|
||||
return &Int64ToOptionalInt64, nil
|
||||
case timeDatatypeRFCNano:
|
||||
return &Int64ToOptionalInt64, nil
|
||||
case durationDatatype:
|
||||
return &Int64ToOptionalInt64, nil
|
||||
case doubleDatatype:
|
||||
return &Float64ToOptionalFloat64, nil
|
||||
case boolDatatype:
|
||||
return &BoolToOptionalBool, nil
|
||||
case longDatatype:
|
||||
return &Int64ToOptionalInt64, nil
|
||||
case uLongDatatype:
|
||||
return &UInt64ToOptionalUInt64, nil
|
||||
case base64BinaryDataType:
|
||||
return &AnyToOptionalString, nil
|
||||
}
|
||||
|
||||
return nil, fmt.Errorf("No matching converter found for [%v]", t)
|
||||
}
|
||||
|
||||
// Init initializes the frame to be returned
|
||||
// fields points at entries in the frame, and provides easier access
|
||||
// names indexes the columns encountered
|
||||
func (fb *FrameBuilder) Init(metadata *influxdb2.FluxTableMetadata) error {
|
||||
columns := metadata.Columns()
|
||||
fb.frames = make([]*data.Frame, 0)
|
||||
fb.tableId = -1
|
||||
fb.value = nil
|
||||
fb.columns = make([]columnInfo, 0)
|
||||
fb.isTimeSeries = false
|
||||
|
||||
for _, col := range columns {
|
||||
switch {
|
||||
case col.Name() == "_value":
|
||||
if fb.value != nil {
|
||||
return fmt.Errorf("multiple values found")
|
||||
}
|
||||
converter, err := getConverter(col.DataType())
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
fb.value = converter
|
||||
case col.Name() == "_measurement":
|
||||
fb.isTimeSeries = true
|
||||
case isTag(col.Name()):
|
||||
fb.labels = append(fb.labels, col.Name())
|
||||
}
|
||||
}
|
||||
|
||||
if !fb.isTimeSeries {
|
||||
fb.labels = make([]string, 0)
|
||||
for _, col := range columns {
|
||||
converter, err := getConverter(col.DataType())
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
fb.columns = append(fb.columns, columnInfo{
|
||||
name: col.Name(),
|
||||
converter: converter,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// Append appends a single entry from an influxdb2 record to a data frame
|
||||
// Values are appended to _value
|
||||
// Tags are appended as labels
|
||||
// _measurement holds the dataframe name
|
||||
// _field holds the field name.
|
||||
func (fb *FrameBuilder) Append(record *influxdb2.FluxRecord) error {
|
||||
table, ok := record.ValueByKey("table").(int64)
|
||||
if ok && table != fb.tableId {
|
||||
fb.totalSeries++
|
||||
if fb.totalSeries > fb.maxSeries {
|
||||
return fmt.Errorf("reached max series limit (%d)", fb.maxSeries)
|
||||
}
|
||||
|
||||
if fb.isTimeSeries {
|
||||
// Series Data
|
||||
labels := make(map[string]string)
|
||||
for _, name := range fb.labels {
|
||||
labels[name] = record.ValueByKey(name).(string)
|
||||
}
|
||||
fb.active = data.NewFrame(
|
||||
record.Measurement(),
|
||||
data.NewFieldFromFieldType(data.FieldTypeTime, 0),
|
||||
data.NewFieldFromFieldType(fb.value.OutputFieldType, 0),
|
||||
)
|
||||
|
||||
fb.active.Fields[0].Name = "Time"
|
||||
fb.active.Fields[1].Name = record.Field()
|
||||
fb.active.Fields[1].Labels = labels
|
||||
} else {
|
||||
fields := make([]*data.Field, len(fb.columns))
|
||||
for idx, col := range fb.columns {
|
||||
fields[idx] = data.NewFieldFromFieldType(col.converter.OutputFieldType, 0)
|
||||
fields[idx].Name = col.name
|
||||
}
|
||||
fb.active = data.NewFrame("", fields...)
|
||||
}
|
||||
|
||||
fb.frames = append(fb.frames, fb.active)
|
||||
fb.tableId = table
|
||||
}
|
||||
|
||||
if fb.isTimeSeries {
|
||||
val, err := fb.value.Converter(record.Value())
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
fb.active.Fields[0].Append(record.Time())
|
||||
fb.active.Fields[1].Append(val)
|
||||
} else {
|
||||
// Table view
|
||||
for idx, col := range fb.columns {
|
||||
val, err := col.converter.Converter(record.ValueByKey(col.name))
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
fb.active.Fields[idx].Append(val)
|
||||
}
|
||||
}
|
||||
|
||||
if fb.active.Fields[0].Len() > fb.maxPoints {
|
||||
return fmt.Errorf("returned too many points in a series: %d", fb.maxPoints)
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
@@ -0,0 +1,46 @@
|
||||
package flux
|
||||
|
||||
import (
|
||||
"regexp"
|
||||
"testing"
|
||||
)
|
||||
|
||||
var isField = regexp.MustCompile(`^_(time|value|measurement|field|start|stop)$`)
|
||||
|
||||
func TestColumnIdentification(t *testing.T) {
|
||||
t.Run("Test Field Identification", func(t *testing.T) {
|
||||
fieldNames := []string{"_time", "_value", "_measurement", "_field", "_start", "_stop"}
|
||||
for _, item := range fieldNames {
|
||||
if !isField.MatchString(item) {
|
||||
t.Fatal("Field", item, "Expected field, but got false")
|
||||
}
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("Test Not Field Identification", func(t *testing.T) {
|
||||
fieldNames := []string{"_header", "_cpu", "_hello", "_doctor"}
|
||||
for _, item := range fieldNames {
|
||||
if isField.MatchString(item) {
|
||||
t.Fatal("Field", item, "Expected NOT a field, but got true")
|
||||
}
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("Test Tag Identification", func(t *testing.T) {
|
||||
tagNames := []string{"header", "value", "tag"}
|
||||
for _, item := range tagNames {
|
||||
if !isTag(item) {
|
||||
t.Fatal("Tag", item, "Expected tag, but got false")
|
||||
}
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("Test Special Case Tag Identification", func(t *testing.T) {
|
||||
notTagNames := []string{"table", "result"}
|
||||
for _, item := range notTagNames {
|
||||
if isTag(item) {
|
||||
t.Fatal("Special tag", item, "Expected NOT a tag, but got true")
|
||||
}
|
||||
}
|
||||
})
|
||||
}
|
||||
@@ -0,0 +1,168 @@
|
||||
package flux
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"strconv"
|
||||
"time"
|
||||
|
||||
"github.com/grafana/grafana-plugin-sdk-go/data"
|
||||
)
|
||||
|
||||
// Int64NOOP .....
|
||||
var Int64NOOP = data.FieldConverter{
|
||||
OutputFieldType: data.FieldTypeInt64,
|
||||
}
|
||||
|
||||
// BoolNOOP .....
|
||||
var BoolNOOP = data.FieldConverter{
|
||||
OutputFieldType: data.FieldTypeBool,
|
||||
}
|
||||
|
||||
// Float64NOOP .....
|
||||
var Float64NOOP = data.FieldConverter{
|
||||
OutputFieldType: data.FieldTypeFloat64,
|
||||
}
|
||||
|
||||
// StringNOOP value is already in the proper format
|
||||
var StringNOOP = data.FieldConverter{
|
||||
OutputFieldType: data.FieldTypeString,
|
||||
}
|
||||
|
||||
// AnyToOptionalString any value as a string
|
||||
var AnyToOptionalString = data.FieldConverter{
|
||||
OutputFieldType: data.FieldTypeNullableString,
|
||||
Converter: func(v interface{}) (interface{}, error) {
|
||||
if v == nil {
|
||||
return nil, nil
|
||||
}
|
||||
str := fmt.Sprintf("%+v", v) // the +v adds field names
|
||||
return &str, nil
|
||||
},
|
||||
}
|
||||
|
||||
// Float64ToOptionalFloat64 optional float value
|
||||
var Float64ToOptionalFloat64 = data.FieldConverter{
|
||||
OutputFieldType: data.FieldTypeNullableFloat64,
|
||||
Converter: func(v interface{}) (interface{}, error) {
|
||||
if v == nil {
|
||||
return nil, nil
|
||||
}
|
||||
val, ok := v.(float64)
|
||||
if !ok { // or return some default value instead of erroring
|
||||
return nil, fmt.Errorf("[float] expected float64 input but got type %T", v)
|
||||
}
|
||||
return &val, nil
|
||||
},
|
||||
}
|
||||
|
||||
// Int64ToOptionalInt64 optional int value
|
||||
var Int64ToOptionalInt64 = data.FieldConverter{
|
||||
OutputFieldType: data.FieldTypeNullableInt64,
|
||||
Converter: func(v interface{}) (interface{}, error) {
|
||||
if v == nil {
|
||||
return nil, nil
|
||||
}
|
||||
val, ok := v.(int64)
|
||||
if !ok { // or return some default value instead of erroring
|
||||
return nil, fmt.Errorf("[int] expected int64 input but got type %T", v)
|
||||
}
|
||||
return &val, nil
|
||||
},
|
||||
}
|
||||
|
||||
// UInt64ToOptionalUInt64 optional int value
|
||||
var UInt64ToOptionalUInt64 = data.FieldConverter{
|
||||
OutputFieldType: data.FieldTypeNullableUint64,
|
||||
Converter: func(v interface{}) (interface{}, error) {
|
||||
if v == nil {
|
||||
return nil, nil
|
||||
}
|
||||
val, ok := v.(uint64)
|
||||
if !ok { // or return some default value instead of erroring
|
||||
return nil, fmt.Errorf("[uint] expected uint64 input but got type %T", v)
|
||||
}
|
||||
return &val, nil
|
||||
},
|
||||
}
|
||||
|
||||
// BoolToOptionalBool optional int value
|
||||
var BoolToOptionalBool = data.FieldConverter{
|
||||
OutputFieldType: data.FieldTypeNullableBool,
|
||||
Converter: func(v interface{}) (interface{}, error) {
|
||||
if v == nil {
|
||||
return nil, nil
|
||||
}
|
||||
val, ok := v.(bool)
|
||||
if !ok { // or return some default value instead of erroring
|
||||
return nil, fmt.Errorf("[bool] expected bool input but got type %T", v)
|
||||
}
|
||||
return &val, nil
|
||||
},
|
||||
}
|
||||
|
||||
// RFC3339StringToNullableTime .....
|
||||
func RFC3339StringToNullableTime(s string) (*time.Time, error) {
|
||||
if s == "" {
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
rv, err := time.Parse(time.RFC3339, s)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
u := rv.UTC()
|
||||
return &u, nil
|
||||
}
|
||||
|
||||
// StringToOptionalFloat64 string to float
|
||||
var StringToOptionalFloat64 = data.FieldConverter{
|
||||
OutputFieldType: data.FieldTypeNullableFloat64,
|
||||
Converter: func(v interface{}) (interface{}, error) {
|
||||
if v == nil {
|
||||
return nil, nil
|
||||
}
|
||||
val, ok := v.(string)
|
||||
if !ok { // or return some default value instead of erroring
|
||||
return nil, fmt.Errorf("[floatz] expected string input but got type %T", v)
|
||||
}
|
||||
fV, err := strconv.ParseFloat(val, 64)
|
||||
return &fV, err
|
||||
},
|
||||
}
|
||||
|
||||
// Float64EpochSecondsToTime numeric seconds to time
|
||||
var Float64EpochSecondsToTime = data.FieldConverter{
|
||||
OutputFieldType: data.FieldTypeTime,
|
||||
Converter: func(v interface{}) (interface{}, error) {
|
||||
fV, ok := v.(float64)
|
||||
if !ok { // or return some default value instead of erroring
|
||||
return nil, fmt.Errorf("[seconds] expected float64 input but got type %T", v)
|
||||
}
|
||||
return time.Unix(int64(fV), 0).UTC(), nil
|
||||
},
|
||||
}
|
||||
|
||||
// Float64EpochMillisToTime convert to time
|
||||
var Float64EpochMillisToTime = data.FieldConverter{
|
||||
OutputFieldType: data.FieldTypeTime,
|
||||
Converter: func(v interface{}) (interface{}, error) {
|
||||
fV, ok := v.(float64)
|
||||
if !ok { // or return some default value instead of erroring
|
||||
return nil, fmt.Errorf("[ms] expected float64 input but got type %T", v)
|
||||
}
|
||||
return time.Unix(0, int64(fV)*int64(time.Millisecond)).UTC(), nil
|
||||
},
|
||||
}
|
||||
|
||||
// Boolean ...
|
||||
var Boolean = data.FieldConverter{
|
||||
OutputFieldType: data.FieldTypeBool,
|
||||
Converter: func(v interface{}) (interface{}, error) {
|
||||
fV, ok := v.(bool)
|
||||
if !ok { // or return some default value instead of erroring
|
||||
return nil, fmt.Errorf("[ms] expected bool input but got type %T", v)
|
||||
}
|
||||
return fV, nil
|
||||
},
|
||||
}
|
||||
@@ -0,0 +1,88 @@
|
||||
package flux
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
|
||||
"github.com/grafana/grafana-plugin-sdk-go/backend"
|
||||
"github.com/grafana/grafana-plugin-sdk-go/data"
|
||||
influxdb2 "github.com/influxdata/influxdb-client-go"
|
||||
)
|
||||
|
||||
// ExecuteQuery runs a flux query using the QueryModel to interpolate the query and the runner to execute it.
|
||||
// maxSeries somehow limits the response.
|
||||
func ExecuteQuery(ctx context.Context, query QueryModel, runner queryRunner, maxSeries int) (dr backend.DataResponse) {
|
||||
dr = backend.DataResponse{}
|
||||
|
||||
flux, err := Interpolate(query)
|
||||
if err != nil {
|
||||
dr.Error = err
|
||||
return
|
||||
}
|
||||
|
||||
glog.Debug("Flux", "interpolated query", flux)
|
||||
|
||||
tables, err := runner.runQuery(ctx, flux)
|
||||
if err != nil {
|
||||
dr.Error = err
|
||||
return
|
||||
}
|
||||
|
||||
dr = readDataFrames(tables, int(float64(query.MaxDataPoints)*1.5), maxSeries)
|
||||
|
||||
for _, frame := range dr.Frames {
|
||||
if frame.Meta == nil {
|
||||
frame.Meta = &data.FrameMeta{}
|
||||
}
|
||||
frame.Meta.ExecutedQueryString = flux
|
||||
}
|
||||
|
||||
return dr
|
||||
}
|
||||
|
||||
func readDataFrames(result *influxdb2.QueryTableResult, maxPoints int, maxSeries int) (dr backend.DataResponse) {
|
||||
dr = backend.DataResponse{}
|
||||
|
||||
builder := &FrameBuilder{
|
||||
maxPoints: maxPoints,
|
||||
maxSeries: maxSeries,
|
||||
}
|
||||
|
||||
for result.Next() {
|
||||
// Observe when there is new grouping key producing new table
|
||||
if result.TableChanged() {
|
||||
if builder.frames != nil {
|
||||
for _, frame := range builder.frames {
|
||||
dr.Frames = append(dr.Frames, frame)
|
||||
}
|
||||
}
|
||||
err := builder.Init(result.TableMetadata())
|
||||
if err != nil {
|
||||
dr.Error = err
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
if builder.frames == nil {
|
||||
dr.Error = fmt.Errorf("Invalid state")
|
||||
return dr
|
||||
}
|
||||
|
||||
err := builder.Append(result.Record())
|
||||
if err != nil {
|
||||
dr.Error = err
|
||||
break
|
||||
}
|
||||
}
|
||||
|
||||
// Add the inprogress record
|
||||
if builder.frames != nil {
|
||||
for _, frame := range builder.frames {
|
||||
dr.Frames = append(dr.Frames, frame)
|
||||
}
|
||||
}
|
||||
|
||||
// Attach any errors (may be null)
|
||||
dr.Error = result.Err()
|
||||
return dr
|
||||
}
|
||||
@@ -0,0 +1,172 @@
|
||||
package flux
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"io/ioutil"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
influxdb2 "github.com/influxdata/influxdb-client-go"
|
||||
)
|
||||
|
||||
//--------------------------------------------------------------
|
||||
// TestData -- reads result from saved files
|
||||
//--------------------------------------------------------------
|
||||
|
||||
// MockRunner reads local file path for testdata.
|
||||
type MockRunner struct {
|
||||
testDataPath string
|
||||
}
|
||||
|
||||
func (r *MockRunner) runQuery(ctx context.Context, q string) (*influxdb2.QueryTableResult, error) {
|
||||
bytes, err := ioutil.ReadFile("./testdata/" + r.testDataPath)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
time.Sleep(100 * time.Millisecond)
|
||||
if r.Method == http.MethodPost {
|
||||
w.WriteHeader(http.StatusOK)
|
||||
_, _ = w.Write(bytes)
|
||||
} else {
|
||||
w.WriteHeader(http.StatusNotFound)
|
||||
}
|
||||
}))
|
||||
defer server.Close()
|
||||
|
||||
client := influxdb2.NewClient(server.URL, "a")
|
||||
return client.QueryApi("x").Query(ctx, q)
|
||||
}
|
||||
|
||||
func TestExecuteSimple(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
|
||||
t.Run("Simple Test", func(t *testing.T) {
|
||||
runner := &MockRunner{
|
||||
testDataPath: "simple.csv",
|
||||
}
|
||||
|
||||
dr := ExecuteQuery(ctx, QueryModel{MaxDataPoints: 100}, runner, 50)
|
||||
|
||||
if dr.Error != nil {
|
||||
t.Fatal(dr.Error)
|
||||
}
|
||||
|
||||
if len(dr.Frames) != 1 {
|
||||
t.Fatalf("Expected 1 frame, received [%d] frames", len(dr.Frames))
|
||||
}
|
||||
|
||||
if !strings.Contains(dr.Frames[0].Name, "test") {
|
||||
t.Fatalf("Frame must match _measurement column. Expected [%s] Got [%s]", "test", dr.Frames[0].Name)
|
||||
}
|
||||
|
||||
if len(dr.Frames[0].Fields[1].Labels) != 2 {
|
||||
t.Fatalf("Error parsing labels. Expected [%d] Got [%d]", 2, len(dr.Frames[0].Fields[1].Labels))
|
||||
}
|
||||
|
||||
if dr.Frames[0].Fields[0].Name != "Time" {
|
||||
t.Fatalf("Error parsing fields. Field 1 should always be time. Got name [%s]", dr.Frames[0].Fields[0].Name)
|
||||
}
|
||||
|
||||
st, _ := dr.Frames[0].StringTable(-1, -1)
|
||||
fmt.Println(st)
|
||||
fmt.Println("----------------------")
|
||||
})
|
||||
}
|
||||
|
||||
func TestExecuteMultiple(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
|
||||
t.Run("Multiple Test", func(t *testing.T) {
|
||||
runner := &MockRunner{
|
||||
testDataPath: "multiple.csv",
|
||||
}
|
||||
|
||||
dr := ExecuteQuery(ctx, QueryModel{MaxDataPoints: 100}, runner, 50)
|
||||
|
||||
if dr.Error != nil {
|
||||
t.Fatal(dr.Error)
|
||||
}
|
||||
|
||||
if len(dr.Frames) != 4 {
|
||||
t.Fatalf("Expected 4 frames, received [%d] frames", len(dr.Frames))
|
||||
}
|
||||
|
||||
if !strings.Contains(dr.Frames[0].Name, "test") {
|
||||
t.Fatalf("Frame must include _measurement column. Expected [%s] Got [%s]", "test", dr.Frames[0].Name)
|
||||
}
|
||||
|
||||
if len(dr.Frames[0].Fields[1].Labels) != 2 {
|
||||
t.Fatalf("Error parsing labels. Expected [%d] Got [%d]", 2, len(dr.Frames[0].Fields[1].Labels))
|
||||
}
|
||||
|
||||
if dr.Frames[0].Fields[0].Name != "Time" {
|
||||
t.Fatalf("Error parsing fields. Field 1 should always be time. Got name [%s]", dr.Frames[0].Fields[0].Name)
|
||||
}
|
||||
|
||||
st, _ := dr.Frames[0].StringTable(-1, -1)
|
||||
fmt.Println(st)
|
||||
fmt.Println("----------------------")
|
||||
})
|
||||
}
|
||||
|
||||
func TestExecuteGrouping(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
|
||||
t.Run("Grouping Test", func(t *testing.T) {
|
||||
runner := &MockRunner{
|
||||
testDataPath: "grouping.csv",
|
||||
}
|
||||
|
||||
dr := ExecuteQuery(ctx, QueryModel{MaxDataPoints: 100}, runner, 50)
|
||||
|
||||
if dr.Error != nil {
|
||||
t.Fatal(dr.Error)
|
||||
}
|
||||
|
||||
if len(dr.Frames) != 3 {
|
||||
t.Fatalf("Expected 3 frames, received [%d] frames", len(dr.Frames))
|
||||
}
|
||||
|
||||
if !strings.Contains(dr.Frames[0].Name, "system") {
|
||||
t.Fatalf("Frame must match _measurement column. Expected [%s] Got [%s]", "test", dr.Frames[0].Name)
|
||||
}
|
||||
|
||||
if len(dr.Frames[0].Fields[1].Labels) != 1 {
|
||||
t.Fatalf("Error parsing labels. Expected [%d] Got [%d]", 1, len(dr.Frames[0].Fields[1].Labels))
|
||||
}
|
||||
|
||||
if dr.Frames[0].Fields[0].Name != "Time" {
|
||||
t.Fatalf("Error parsing fields. Field 1 should always be time. Got name [%s]", dr.Frames[0].Fields[0].Name)
|
||||
}
|
||||
|
||||
st, _ := dr.Frames[0].StringTable(-1, -1)
|
||||
fmt.Println(st)
|
||||
fmt.Println("----------------------")
|
||||
})
|
||||
}
|
||||
|
||||
func TestBuckets(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
|
||||
t.Run("Buckes", func(t *testing.T) {
|
||||
runner := &MockRunner{
|
||||
testDataPath: "buckets.csv",
|
||||
}
|
||||
|
||||
dr := ExecuteQuery(ctx, QueryModel{MaxDataPoints: 100}, runner, 50)
|
||||
|
||||
if dr.Error != nil {
|
||||
t.Fatal(dr.Error)
|
||||
}
|
||||
|
||||
st, _ := dr.Frames[0].StringTable(-1, -1)
|
||||
fmt.Println(st)
|
||||
fmt.Println("----------------------")
|
||||
})
|
||||
}
|
||||
@@ -0,0 +1,102 @@
|
||||
package flux
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
|
||||
"github.com/grafana/grafana-plugin-sdk-go/backend"
|
||||
"github.com/grafana/grafana/pkg/infra/log"
|
||||
"github.com/grafana/grafana/pkg/models"
|
||||
"github.com/grafana/grafana/pkg/tsdb"
|
||||
influxdb2 "github.com/influxdata/influxdb-client-go"
|
||||
)
|
||||
|
||||
var (
|
||||
glog log.Logger
|
||||
)
|
||||
|
||||
func init() {
|
||||
glog = log.New("tsdb.influx_flux")
|
||||
}
|
||||
|
||||
// Query builds flux queries, executes them, and returns the results.
|
||||
func Query(ctx context.Context, dsInfo *models.DataSource, tsdbQuery *tsdb.TsdbQuery) (*tsdb.Response, error) {
|
||||
tRes := &tsdb.Response{
|
||||
Results: make(map[string]*tsdb.QueryResult),
|
||||
}
|
||||
runner, err := RunnerFromDataSource(dsInfo)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
for _, query := range tsdbQuery.Queries {
|
||||
|
||||
qm, err := GetQueryModelTSDB(query, tsdbQuery.TimeRange, dsInfo)
|
||||
if err != nil {
|
||||
tRes.Results[query.RefId] = &tsdb.QueryResult{Error: err}
|
||||
continue
|
||||
}
|
||||
|
||||
res := ExecuteQuery(context.Background(), *qm, runner, 10)
|
||||
|
||||
tRes.Results[query.RefId] = backendDataResponseToTSDBResponse(&res, query.RefId)
|
||||
}
|
||||
return tRes, nil
|
||||
}
|
||||
|
||||
// Runner is an influxdb2 Client with an attached org property and is used
|
||||
// for running flux queries.
|
||||
type Runner struct {
|
||||
client influxdb2.Client
|
||||
org string
|
||||
}
|
||||
|
||||
// This is an interface to help testing
|
||||
type queryRunner interface {
|
||||
runQuery(ctx context.Context, q string) (*influxdb2.QueryTableResult, error)
|
||||
}
|
||||
|
||||
// runQuery executes fluxQuery against the Runner's organization and returns an flux typed result.
|
||||
func (r *Runner) runQuery(ctx context.Context, fluxQuery string) (*influxdb2.QueryTableResult, error) {
|
||||
return r.client.QueryApi(r.org).Query(ctx, fluxQuery)
|
||||
}
|
||||
|
||||
// RunnerFromDataSource creates a runner from the datasource model (the datasource instance's configuration).
|
||||
func RunnerFromDataSource(dsInfo *models.DataSource) (*Runner, error) {
|
||||
org := dsInfo.JsonData.Get("organization").MustString("")
|
||||
if org == "" {
|
||||
return nil, fmt.Errorf("missing organization in datasource configuration")
|
||||
}
|
||||
|
||||
url := dsInfo.Url
|
||||
if url == "" {
|
||||
return nil, fmt.Errorf("missing url from datasource configuration")
|
||||
}
|
||||
token, found := dsInfo.SecureJsonData.DecryptedValue("token")
|
||||
if !found {
|
||||
return nil, fmt.Errorf("token is missing from datasource configuration and is needed to use Flux")
|
||||
}
|
||||
|
||||
return &Runner{
|
||||
client: influxdb2.NewClient(url, token),
|
||||
org: org,
|
||||
}, nil
|
||||
|
||||
}
|
||||
|
||||
// backendDataResponseToTSDBResponse takes the SDK's style response and changes it into a
|
||||
// tsdb.QueryResult. This is a wrapper so less of existing code needs to be changed. This should
|
||||
// be able to be removed in the near future https://github.com/grafana/grafana/pull/25472.
|
||||
func backendDataResponseToTSDBResponse(dr *backend.DataResponse, refID string) *tsdb.QueryResult {
|
||||
qr := &tsdb.QueryResult{RefId: refID}
|
||||
|
||||
if dr.Error != nil {
|
||||
qr.Error = dr.Error
|
||||
return qr
|
||||
}
|
||||
|
||||
if dr.Frames != nil {
|
||||
qr.Dataframes = tsdb.NewDecodedDataFrames(dr.Frames)
|
||||
}
|
||||
return qr
|
||||
}
|
||||
@@ -0,0 +1,45 @@
|
||||
package flux
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"regexp"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/grafana/grafana-plugin-sdk-go/backend"
|
||||
)
|
||||
|
||||
const variableFilter = `(?m)([a-zA-Z]+)\.([a-zA-Z]+)`
|
||||
|
||||
// Interpolate processes macros
|
||||
func Interpolate(query QueryModel) (string, error) {
|
||||
|
||||
flux := query.RawQuery
|
||||
|
||||
variableFilterExp, err := regexp.Compile(variableFilter)
|
||||
matches := variableFilterExp.FindAllStringSubmatch(flux, -1)
|
||||
if matches != nil {
|
||||
timeRange := query.TimeRange
|
||||
from := timeRange.From.UTC().Format(time.RFC3339)
|
||||
to := timeRange.To.UTC().Format(time.RFC3339)
|
||||
for _, match := range matches {
|
||||
switch match[2] {
|
||||
case "timeRangeStart":
|
||||
flux = strings.ReplaceAll(flux, match[0], from)
|
||||
case "timeRangeStop":
|
||||
flux = strings.ReplaceAll(flux, match[0], to)
|
||||
case "windowPeriod":
|
||||
flux = strings.ReplaceAll(flux, match[0], query.Interval.String())
|
||||
case "bucket":
|
||||
flux = strings.ReplaceAll(flux, match[0], "\""+query.Options.Bucket+"\"")
|
||||
case "defaultBucket":
|
||||
flux = strings.ReplaceAll(flux, match[0], "\""+query.Options.DefaultBucket+"\"")
|
||||
case "organization":
|
||||
flux = strings.ReplaceAll(flux, match[0], "\""+query.Options.Organization+"\"")
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
backend.Logger.Info(fmt.Sprintf("%s => %v", flux, query.Options))
|
||||
return flux, err
|
||||
}
|
||||
@@ -0,0 +1,56 @@
|
||||
package flux
|
||||
|
||||
import (
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/google/go-cmp/cmp"
|
||||
"github.com/grafana/grafana-plugin-sdk-go/backend"
|
||||
)
|
||||
|
||||
func TestInterpolate(t *testing.T) {
|
||||
// Unix sec: 1500376552
|
||||
// Unix ms: 1500376552001
|
||||
|
||||
timeRange := backend.TimeRange{
|
||||
From: time.Unix(0, 0),
|
||||
To: time.Unix(0, 0),
|
||||
}
|
||||
|
||||
options := QueryOptions{
|
||||
Organization: "grafana1",
|
||||
Bucket: "grafana2",
|
||||
DefaultBucket: "grafana3",
|
||||
}
|
||||
|
||||
tests := []struct {
|
||||
name string
|
||||
before string
|
||||
after string
|
||||
}{
|
||||
{
|
||||
name: "interpolate flux variables",
|
||||
before: `v.timeRangeStart, something.timeRangeStop, XYZ.bucket, uuUUu.defaultBucket, aBcDefG.organization, window.windowPeriod, a91{}.bucket`,
|
||||
after: `1970-01-01T00:00:00Z, 1970-01-01T00:00:00Z, "grafana2", "grafana3", "grafana1", 1s, a91{}.bucket`,
|
||||
},
|
||||
}
|
||||
for _, tt := range tests {
|
||||
t.Run(tt.name, func(t *testing.T) {
|
||||
|
||||
query := QueryModel{
|
||||
RawQuery: tt.before,
|
||||
Options: options,
|
||||
TimeRange: timeRange,
|
||||
MaxDataPoints: 1,
|
||||
Interval: 1000 * 1000 * 1000,
|
||||
}
|
||||
interpolatedQuery, err := Interpolate(query)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if diff := cmp.Diff(tt.after, interpolatedQuery); diff != "" {
|
||||
t.Fatalf("Result mismatch (-want +got):\n%s", diff)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,88 @@
|
||||
package flux
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"time"
|
||||
|
||||
"github.com/grafana/grafana-plugin-sdk-go/backend"
|
||||
"github.com/grafana/grafana/pkg/models"
|
||||
"github.com/grafana/grafana/pkg/tsdb"
|
||||
)
|
||||
|
||||
// QueryOptions represents datasource configuration options
|
||||
type QueryOptions struct {
|
||||
Bucket string `json:"bucket"`
|
||||
DefaultBucket string `json:"defaultBucket"`
|
||||
Organization string `json:"organization"`
|
||||
}
|
||||
|
||||
// QueryModel represents a spreadsheet query.
|
||||
type QueryModel struct {
|
||||
RawQuery string `json:"query"`
|
||||
Options QueryOptions `json:"options"`
|
||||
|
||||
// Not from JSON
|
||||
TimeRange backend.TimeRange `json:"-"`
|
||||
MaxDataPoints int64 `json:"-"`
|
||||
Interval time.Duration `json:"-"`
|
||||
}
|
||||
|
||||
// The following is commented out but kept as it should be useful when
|
||||
// restoring this code to be closer to the SDK's models.
|
||||
|
||||
// func GetQueryModel(query backend.DataQuery) (*QueryModel, error) {
|
||||
// model := &QueryModel{}
|
||||
|
||||
// err := json.Unmarshal(query.JSON, &model)
|
||||
// if err != nil {
|
||||
// return nil, fmt.Errorf("error reading query: %s", err.Error())
|
||||
// }
|
||||
|
||||
// // Copy directly from the well typed query
|
||||
// model.TimeRange = query.TimeRange
|
||||
// model.MaxDataPoints = query.MaxDataPoints
|
||||
// model.Interval = query.Interval
|
||||
// return model, nil
|
||||
// }
|
||||
|
||||
// GetQueryModelTSDB builds a QueryModel from tsdb.Query information and datasource configuration (dsInfo).
|
||||
func GetQueryModelTSDB(query *tsdb.Query, timeRange *tsdb.TimeRange, dsInfo *models.DataSource) (*QueryModel, error) {
|
||||
model := &QueryModel{}
|
||||
queryBytes, err := query.Model.Encode()
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("failed to re-encode the flux query into JSON: %w", err)
|
||||
}
|
||||
|
||||
err = json.Unmarshal(queryBytes, &model)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("error reading query: %s", err.Error())
|
||||
}
|
||||
if model.Options.DefaultBucket == "" {
|
||||
model.Options.DefaultBucket = dsInfo.JsonData.Get("defaultBucket").MustString("")
|
||||
}
|
||||
if model.Options.Bucket == "" {
|
||||
model.Options.Bucket = model.Options.DefaultBucket
|
||||
}
|
||||
if model.Options.Organization == "" {
|
||||
model.Options.Organization = dsInfo.JsonData.Get("organization").MustString("")
|
||||
}
|
||||
startTime, err := timeRange.ParseFrom()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
endTime, err := timeRange.ParseTo()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// Copy directly from the well typed query
|
||||
model.TimeRange = backend.TimeRange{
|
||||
From: startTime,
|
||||
To: endTime,
|
||||
}
|
||||
model.MaxDataPoints = query.MaxDataPoints
|
||||
model.Interval = time.Millisecond * time.Duration(query.IntervalMs)
|
||||
return model, nil
|
||||
}
|
||||
+7
@@ -0,0 +1,7 @@
|
||||
#datatype,string,long,string,string,string,string,long
|
||||
#group,false,false,false,false,true,false,false
|
||||
#default,_result,,,,,,
|
||||
,result,table,name,id,organizationID,retentionPolicy,retentionPeriod
|
||||
,,0,grafana,059b46a59abab001,059b46a59abab000,,604800000000000
|
||||
,,0,_tasks,059b46a59abab002,059b46a59abab000,,259200000000000
|
||||
,,0,_monitoring,059b46a59abab003,059b46a59abab000,,604800000000000
|
||||
|
+21
@@ -0,0 +1,21 @@
|
||||
#datatype,string,long,dateTime:RFC3339,dateTime:RFC3339,string,string,string,double,dateTime:RFC3339
|
||||
#group,false,false,true,true,true,true,true,false,false
|
||||
#default,mean,,,,,,,,
|
||||
,result,table,_start,_stop,_field,_measurement,host,_value,_time
|
||||
,,0,2020-05-05T18:38:47.207881833Z,2020-05-05T19:38:47.207881833Z,load1,system,hostname,,2020-05-05T18:38:50Z
|
||||
,,0,2020-05-05T18:38:47.207881833Z,2020-05-05T19:38:47.207881833Z,load1,system,hostname,3.56,2020-05-05T18:39:00Z
|
||||
,,0,2020-05-05T18:38:47.207881833Z,2020-05-05T19:38:47.207881833Z,load1,system,hostname,,2020-05-05T19:38:47.207881833Z
|
||||
,,1,2020-05-05T18:38:47.207881833Z,2020-05-05T19:38:47.207881833Z,load15,system,hostname,,2020-05-05T18:38:50Z
|
||||
,,1,2020-05-05T18:38:47.207881833Z,2020-05-05T19:38:47.207881833Z,load15,system,hostname,2.51,2020-05-05T18:39:00Z
|
||||
,,1,2020-05-05T18:38:47.207881833Z,2020-05-05T19:38:47.207881833Z,load15,system,hostname,1.74,2020-05-05T19:38:40Z
|
||||
,,1,2020-05-05T18:38:47.207881833Z,2020-05-05T19:38:47.207881833Z,load15,system,hostname,,2020-05-05T19:38:47.207881833Z
|
||||
,,2,2020-05-05T18:38:47.207881833Z,2020-05-05T19:38:47.207881833Z,load5,system,hostname,,2020-05-05T18:38:50Z
|
||||
,,2,2020-05-05T18:38:47.207881833Z,2020-05-05T19:38:47.207881833Z,load5,system,hostname,3.14,2020-05-05T18:39:00Z
|
||||
,,2,2020-05-05T18:38:47.207881833Z,2020-05-05T19:38:47.207881833Z,load5,system,hostname,3.04,2020-05-05T18:39:10Z
|
||||
,,2,2020-05-05T18:38:47.207881833Z,2020-05-05T19:38:47.207881833Z,load5,system,hostname,1.8,2020-05-05T19:37:50Z
|
||||
,,2,2020-05-05T18:38:47.207881833Z,2020-05-05T19:38:47.207881833Z,load5,system,hostname,1.76,2020-05-05T19:38:00Z
|
||||
,,2,2020-05-05T18:38:47.207881833Z,2020-05-05T19:38:47.207881833Z,load5,system,hostname,1.75,2020-05-05T19:38:10Z
|
||||
,,2,2020-05-05T18:38:47.207881833Z,2020-05-05T19:38:47.207881833Z,load5,system,hostname,1.71,2020-05-05T19:38:20Z
|
||||
,,2,2020-05-05T18:38:47.207881833Z,2020-05-05T19:38:47.207881833Z,load5,system,hostname,1.77,2020-05-05T19:38:30Z
|
||||
,,2,2020-05-05T18:38:47.207881833Z,2020-05-05T19:38:47.207881833Z,load5,system,hostname,1.71,2020-05-05T19:38:40Z
|
||||
,,2,2020-05-05T18:38:47.207881833Z,2020-05-05T19:38:47.207881833Z,load5,system,hostname,,2020-05-05T19:38:47.207881833Z
|
||||
|
+25
@@ -0,0 +1,25 @@
|
||||
#datatype,string,long,dateTime:RFC3339,dateTime:RFC3339,dateTime:RFC3339,double,string,string,string,string
|
||||
#group,false,false,true,true,false,false,true,true,true,true
|
||||
#default,_result,,,,,,,,,
|
||||
,result,table,_start,_stop,_time,_value,_field,_measurement,a,b
|
||||
,,0,2020-02-17T22:19:49.747562847Z,2020-02-18T22:19:49.747562847Z,2020-02-18T10:34:08.135814545Z,1.4,f,test,1,adsfasdf
|
||||
,,0,2020-02-17T22:19:49.747562847Z,2020-02-18T22:19:49.747562847Z,2020-02-18T22:08:44.850214724Z,6.6,f,test,1,adsfasdf
|
||||
#datatype,string,long,dateTime:RFC3339,dateTime:RFC3339,dateTime:RFC3339,long,string,string,string,string
|
||||
#group,false,false,true,true,false,false,true,true,true,true
|
||||
#default,_result,,,,,,,,,
|
||||
,result,table,_start,_stop,_time,_value,_field,_measurement,a,b
|
||||
,,1,2020-02-17T22:19:49.747562847Z,2020-02-18T22:19:49.747562847Z,2020-02-18T10:34:08.135814545Z,4,i,test,1,adsfasdf
|
||||
,,1,2020-02-17T22:19:49.747562847Z,2020-02-18T22:19:49.747562847Z,2020-02-18T22:08:44.850214724Z,-1,i,test,1,adsfasdf
|
||||
#datatype,string,long,dateTime:RFC3339,dateTime:RFC3339,dateTime:RFC3339,bool,string,string,string,string
|
||||
#group,false,false,true,true,false,false,true,true,true,true
|
||||
#default,_result,,,,,,,,,
|
||||
,result,table,_start,_stop,_time,_value,_field,_measurement,a,b
|
||||
,,2,2020-02-17T22:19:49.747562847Z,2020-02-18T22:19:49.747562847Z,2020-02-18T22:08:44.62797864Z,false,f,test,0,adsfasdf
|
||||
,,2,2020-02-17T22:19:49.747562847Z,2020-02-18T22:19:49.747562847Z,2020-02-18T22:08:44.969100374Z,true,f,test,0,adsfasdf
|
||||
#datatype,string,long,dateTime:RFC3339Nano,dateTime:RFC3339Nano,dateTime:RFC3339Nano,unsignedLong,string,string,string,string
|
||||
#group,false,false,true,true,false,false,true,true,true,true
|
||||
#default,_result,,,,,,,,,
|
||||
,result,table,_start,_stop,_time,_value,_field,_measurement,a,b
|
||||
,,3,2020-02-17T22:19:49.747562847Z,2020-02-18T22:19:49.747562847Z,2020-02-18T22:08:44.62797864Z,0,i,test,0,adsfasdf
|
||||
,,3,2020-02-17T22:19:49.747562847Z,2020-02-18T22:19:49.747562847Z,2020-02-18T22:08:44.969100374Z,2,i,test,0,adsfasdf
|
||||
|
||||
|
+6
@@ -0,0 +1,6 @@
|
||||
#datatype,string,long,dateTime:RFC3339,dateTime:RFC3339,dateTime:RFC3339,double,string,string,string,string
|
||||
#group,false,false,true,true,false,false,true,true,true,true
|
||||
#default,_result,,,,,,,,,
|
||||
,result,table,_start,_stop,_time,_value,_field,_measurement,a,b
|
||||
,,0,2020-02-17T22:19:49.747562847Z,2020-02-18T22:19:49.747562847Z,2020-02-18T10:34:08.135814545Z,1.4,f,test,1,adsfasdf
|
||||
,,0,2020-02-17T22:19:49.747562847Z,2020-02-18T22:19:49.747562847Z,2020-02-18T22:08:44.850214724Z,6.6,f,test,1,adsfasdf
|
||||
|
@@ -14,6 +14,7 @@ import (
|
||||
"github.com/grafana/grafana/pkg/models"
|
||||
"github.com/grafana/grafana/pkg/setting"
|
||||
"github.com/grafana/grafana/pkg/tsdb"
|
||||
"github.com/grafana/grafana/pkg/tsdb/influxdb/flux"
|
||||
"golang.org/x/net/context/ctxhttp"
|
||||
)
|
||||
|
||||
@@ -42,8 +43,35 @@ func init() {
|
||||
tsdb.RegisterTsdbQueryEndpoint("influxdb", NewInfluxDBExecutor)
|
||||
}
|
||||
|
||||
func AllFlux(queries *tsdb.TsdbQuery) (bool, error) {
|
||||
var hasFlux bool
|
||||
var allFlux bool
|
||||
for i, q := range queries.Queries {
|
||||
qType := q.Model.Get("queryType").MustString("")
|
||||
if qType == "Flux" {
|
||||
hasFlux = true
|
||||
if i == 0 && hasFlux {
|
||||
allFlux = true
|
||||
continue
|
||||
}
|
||||
}
|
||||
if allFlux && qType != "Flux" {
|
||||
return true, fmt.Errorf("when using flux, all queries must be a flux query")
|
||||
}
|
||||
}
|
||||
return allFlux, nil
|
||||
}
|
||||
|
||||
func (e *InfluxDBExecutor) Query(ctx context.Context, dsInfo *models.DataSource, tsdbQuery *tsdb.TsdbQuery) (*tsdb.Response, error) {
|
||||
result := &tsdb.Response{}
|
||||
allFlux, err := AllFlux(tsdbQuery)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
if allFlux {
|
||||
return flux.Query(ctx, dsInfo, tsdbQuery)
|
||||
}
|
||||
|
||||
query, err := e.getQuery(dsInfo, tsdbQuery.Queries, tsdbQuery)
|
||||
if err != nil {
|
||||
|
||||
Reference in New Issue
Block a user