SSE: Support for ML query node (#69963)

* introduce a new node-type ML and implement a command outlier that uses ML plugin as a source of data.
* add feature flag mlExpressions that guards the feature
This commit is contained in:
Yuri Tseretyan
2023-07-13 20:37:50 +03:00
committed by GitHub
parent e8b4228f89
commit 541bfe636d
18 changed files with 1196 additions and 11 deletions
+83
View File
@@ -0,0 +1,83 @@
package ml
import (
"strings"
"time"
"github.com/grafana/grafana-plugin-sdk-go/backend"
jsoniter "github.com/json-iterator/go"
)
type CommandConfiguration struct {
Type string `json:"type"`
IntervalMs *uint `json:"intervalMs,omitempty"`
Config jsoniter.RawMessage `json:"config"`
}
type OutlierCommandConfiguration struct {
DatasourceType string `json:"datasource_type"`
DatasourceUID string `json:"datasource_uid,omitempty"`
// If Query is empty it should be contained in a datasource specific format
// inside of QueryParms.
Query string `json:"query,omitempty"`
QueryParams map[string]interface{} `json:"query_params,omitempty"`
Algorithm map[string]interface{} `json:"algorithm"`
ResponseType string `json:"response_type"`
}
// outlierAttributes is outlier command configuration that is sent to Machine learning API
type outlierAttributes struct {
OutlierCommandConfiguration
GrafanaURL string `json:"grafana_url"`
StartEndAttributes timeRangeAndInterval `json:"start_end_attributes"`
}
type outlierData struct {
Attributes outlierAttributes `json:"attributes"`
}
// outlierRequestBody describes a request body that is sent to Outlier API
type outlierRequestBody struct {
Data outlierData `json:"data"`
}
type timeRangeAndInterval struct {
Start mlTime `json:"start"`
End mlTime `json:"end"`
Interval int64 `json:"interval"` // Interval is expected to be in milliseconds
}
func newTimeRangeAndInterval(from, to time.Time, interval time.Duration) timeRangeAndInterval {
return timeRangeAndInterval{
Start: mlTime(from),
End: mlTime(to),
Interval: interval.Milliseconds(),
}
}
// mlTime is a time.Time that is marshalled as a string in a format is supported by Machine Learning API
type mlTime time.Time
func (t *mlTime) UnmarshalJSON(b []byte) error {
s := strings.Trim(string(b), "\"")
parsed, err := time.Parse(timeFormat, s)
if err != nil {
return err
}
*t = mlTime(parsed)
return nil
}
// MarshalJSON implements the Marshaler interface.
func (t mlTime) MarshalJSON() ([]byte, error) {
return []byte("\"" + time.Time(t).Format(timeFormat) + "\""), nil
}
// outlierResponse is a model that represents a response of the outlier proxy API.
type outlierResponse struct {
Status string `json:"status"`
Data *backend.QueryDataResponse `json:"data,omitempty"`
Error string `json:"error,omitempty"`
}
+63
View File
@@ -0,0 +1,63 @@
package ml
import (
"fmt"
"strings"
"time"
"github.com/grafana/grafana-plugin-sdk-go/backend"
jsoniter "github.com/json-iterator/go"
"github.com/grafana/grafana/pkg/api/response"
)
var json = jsoniter.ConfigCompatibleWithStandardLibrary
type CommandType string
const (
Outlier CommandType = "outlier"
// format of the time used by outlier API
timeFormat = "2006-01-02T15:04:05.999999999"
defaultInterval = 1000 * time.Millisecond
)
// Command is an interface implemented by all Machine Learning commands that can be executed against ML API.
type Command interface {
// DatasourceUID returns UID of a data source that is used by machine learning as the source of data
DatasourceUID() string
// Execute creates a payload send request to the ML API by calling the function argument sendRequest, and then parses response.
// Function sendRequest is supposed to abstract the client configuration such creating http request, adding authorization parameters, host etc.
Execute(from, to time.Time, sendRequest func(method string, path string, payload []byte) (response.Response, error)) (*backend.QueryDataResponse, error)
}
// UnmarshalCommand parses a config parameters and creates a command. Requires key `type` to be specified.
// Based on the value of `type` field it parses a Command
func UnmarshalCommand(query []byte, appURL string) (Command, error) {
var expr CommandConfiguration
err := json.Unmarshal(query, &expr)
if err != nil {
return nil, fmt.Errorf("failed to unmarshall Machine learning command: %w", err)
}
if len(expr.Type) == 0 {
return nil, fmt.Errorf("required field 'type' is not specified or empty. Should be one of [%s]", Outlier)
}
if len(expr.Config) == 0 {
return nil, fmt.Errorf("required field 'config' is not specified")
}
var cmd Command
switch mlType := strings.ToLower(expr.Type); mlType {
case string(Outlier):
cmd, err = unmarshalOutlierCommand(expr, appURL)
default:
return nil, fmt.Errorf("unsupported command type. Should be one of [%s]", Outlier)
}
if err != nil {
return nil, fmt.Errorf("failed to unmarshal Machine learning %s command: %w", expr.Type, err)
}
return cmd, nil
}
+167
View File
@@ -0,0 +1,167 @@
package ml
import (
"testing"
"time"
"github.com/google/uuid"
"github.com/stretchr/testify/require"
)
func TestUnmarshalCommand(t *testing.T) {
appURL := "https://grafana.com"
updateJson := func(cmd string, f func(m map[string]interface{})) func(t *testing.T) []byte {
return func(t *testing.T) []byte {
var d map[string]interface{}
require.NoError(t, json.UnmarshalFromString(cmd, &d))
f(d)
data, err := json.Marshal(d)
require.NoError(t, err)
return data
}
}
t.Run("should parse outlier command", func(t *testing.T) {
cmd, err := UnmarshalCommand([]byte(outlierQuery), appURL)
require.NoError(t, err)
require.IsType(t, &OutlierCommand{}, cmd)
outlier := cmd.(*OutlierCommand)
require.Equal(t, 1234*time.Millisecond, outlier.interval)
require.Equal(t, appURL, outlier.appURL)
require.Equal(t, OutlierCommandConfiguration{
DatasourceType: "prometheus",
DatasourceUID: "a4ce599c-4c93-44b9-be5b-76385b8c01be",
QueryParams: map[string]interface{}{
"expr": "go_goroutines{}",
"range": true,
"refId": "A",
},
Algorithm: map[string]interface{}{
"name": "dbscan",
"config": map[string]interface{}{
"epsilon": 7.667,
},
"sensitivity": 0.83,
},
ResponseType: "binary",
}, outlier.config)
})
t.Run("should fallback to default if 'intervalMs' is not specified", func(t *testing.T) {
data := updateJson(outlierQuery, func(m map[string]interface{}) {
delete(m, "intervalMs")
})(t)
cmd, err := UnmarshalCommand(data, appURL)
require.NoError(t, err)
outlier := cmd.(*OutlierCommand)
require.Equal(t, defaultInterval, outlier.interval)
})
t.Run("fails when", func(t *testing.T) {
testCases := []struct {
name string
config func(t *testing.T) []byte
err string
}{
{
name: "field 'type' is missing",
config: updateJson(outlierQuery, func(cmd map[string]interface{}) {
delete(cmd, "type")
}),
err: "required field 'type' is not specified or empty. Should be one of [outlier]",
},
{
name: "field 'type' is not known",
config: updateJson(outlierQuery, func(cmd map[string]interface{}) {
cmd["type"] = uuid.NewString()
}),
err: "unsupported command type. Should be one of [outlier]",
},
{
name: "field 'type' is not string",
config: updateJson(outlierQuery, func(cmd map[string]interface{}) {
cmd["type"] = map[string]interface{}{
"data": 1,
}
}),
err: "failed to unmarshall Machine learning command",
},
{
name: "field 'config' is missing",
config: updateJson(outlierQuery, func(cmd map[string]interface{}) {
delete(cmd, "config")
}),
err: "required field 'config' is not specified",
},
{
name: "field 'intervalMs' is not number",
config: updateJson(outlierQuery, func(cmd map[string]interface{}) {
cmd["intervalMs"] = "test"
}),
err: "failed to unmarshall Machine learning command",
},
{
name: "field 'config.datasource_uid' is not specified",
config: updateJson(outlierQuery, func(cmd map[string]interface{}) {
cfg := cmd["config"].(map[string]interface{})
delete(cfg, "datasource_uid")
}),
err: "required field `config.datasource_uid` is not specified",
},
{
name: "field 'config.algorithm' is not specified",
config: updateJson(outlierQuery, func(cmd map[string]interface{}) {
cfg := cmd["config"].(map[string]interface{})
delete(cfg, "algorithm")
}),
err: "required field `config.algorithm` is not specified",
},
{
name: "field 'config.response_type' is not specified",
config: updateJson(outlierQuery, func(cmd map[string]interface{}) {
cfg := cmd["config"].(map[string]interface{})
delete(cfg, "response_type")
}),
err: "required field `config.response_type` is not specified",
},
{
name: "fields 'config.query' and 'config.query_params' are not specified",
config: updateJson(outlierQuery, func(cmd map[string]interface{}) {
cfg := cmd["config"].(map[string]interface{})
delete(cfg, "query")
delete(cfg, "query_params")
}),
err: "neither of required fields `config.query_params` or `config.query` are specified",
},
}
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
_, err := UnmarshalCommand(testCase.config(t), appURL)
require.ErrorContains(t, err, testCase.err)
})
}
})
}
const outlierQuery = `
{
"type": "outlier",
"intervalMs": 1234,
"config": {
"datasource_uid": "a4ce599c-4c93-44b9-be5b-76385b8c01be",
"datasource_type": "prometheus",
"query_params": {
"expr": "go_goroutines{}",
"range": true,
"refId": "A"
},
"response_type": "binary",
"algorithm": {
"name": "dbscan",
"config": {
"epsilon": 7.667
},
"sensitivity": 0.83
}
}
}
`
+116
View File
@@ -0,0 +1,116 @@
package ml
import (
"fmt"
"net/http"
"time"
"github.com/grafana/grafana-plugin-sdk-go/backend"
"github.com/grafana/grafana/pkg/api/response"
)
// OutlierCommand implements Command that sends a request to outlier proxy API and converts response to backend.QueryDataResponse
type OutlierCommand struct {
config OutlierCommandConfiguration
appURL string
interval time.Duration
}
var _ Command = OutlierCommand{}
func (c OutlierCommand) DatasourceUID() string {
return c.config.DatasourceUID
}
// Execute copies the original configuration JSON and appends (overwrites) a field "start_end_attributes" and "grafana_url" to the root object.
// The value of "start_end_attributes" is JSON object that configures time range and interval.
// The value of "grafana_url" is app URL that should be used by ML to query the data source.
// After payload is generated it sends it to POST /proxy/api/v1/outlier endpoint and parses the response.
// The proxy API normally responds with a structured data. It recognizes status 200 and 204 as successful result.
// Other statuses are considered unsuccessful and result in error. Tries to extract error from the structured payload.
// Otherwise, mentions the full message in error
func (c OutlierCommand) Execute(from, to time.Time, sendRequest func(method string, path string, payload []byte) (response.Response, error)) (*backend.QueryDataResponse, error) {
payload := outlierRequestBody{
Data: outlierData{
Attributes: outlierAttributes{
OutlierCommandConfiguration: c.config,
GrafanaURL: c.appURL,
StartEndAttributes: newTimeRangeAndInterval(from, to, c.interval),
},
},
}
requestBody, err := json.Marshal(payload)
if err != nil {
return nil, err
}
resp, err := sendRequest(http.MethodPost, "/proxy/api/v1/outlier", requestBody)
if err != nil {
return nil, fmt.Errorf("failed to call ML API: %w", err)
}
if resp == nil {
return nil, fmt.Errorf("response is nil")
}
// Outlier proxy API usually returns all responses with this body.
var respData outlierResponse
respBody := resp.Body()
err = json.Unmarshal(respBody, &respData)
if err != nil {
return nil, fmt.Errorf("unexpected format of the response from ML API, status: %d, response: %s", resp.Status(), respBody)
}
if respData.Status == "error" {
return nil, fmt.Errorf("ML API responded with error: %s", respData.Error)
}
if resp.Status() == http.StatusNoContent {
return nil, nil
}
if resp.Status() == http.StatusOK {
return respData.Data, nil
}
return nil, fmt.Errorf("unexpected status %d returned by ML API, response: %s", resp.Status(), respBody)
}
// unmarshalOutlierCommand parses the CommandConfiguration.Config, validates data and produces OutlierCommand.
func unmarshalOutlierCommand(expr CommandConfiguration, appURL string) (*OutlierCommand, error) {
var cfg OutlierCommandConfiguration
err := json.Unmarshal(expr.Config, &cfg)
if err != nil {
return nil, fmt.Errorf("failed to unmarshal outlier command: %w", err)
}
if len(cfg.DatasourceUID) == 0 {
return nil, fmt.Errorf("required field `config.datasource_uid` is not specified")
}
if len(cfg.Query) == 0 && len(cfg.QueryParams) == 0 {
return nil, fmt.Errorf("neither of required fields `config.query_params` or `config.query` are specified")
}
if len(cfg.ResponseType) == 0 {
return nil, fmt.Errorf("required field `config.response_type` is not specified")
}
if len(cfg.Algorithm) == 0 {
return nil, fmt.Errorf("required field `config.algorithm` is not specified")
}
interval := defaultInterval
if expr.IntervalMs != nil {
i := time.Duration(*expr.IntervalMs) * time.Millisecond
if i > 0 {
interval = i
}
}
return &OutlierCommand{
config: cfg,
interval: interval,
appURL: appURL,
}, nil
}
+194
View File
@@ -0,0 +1,194 @@
package ml
import (
"errors"
"fmt"
"net/http"
"testing"
"time"
"github.com/grafana/grafana-plugin-sdk-go/backend"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"golang.org/x/exp/slices"
"github.com/grafana/grafana/pkg/api/response"
)
func TestOutlierExec(t *testing.T) {
outlier := OutlierCommand{
config: OutlierCommandConfiguration{
DatasourceType: "prometheus",
DatasourceUID: "a4ce599c-4c93-44b9-be5b-76385b8c01be",
QueryParams: map[string]interface{}{
"expr": "go_goroutines{}",
"range": true,
"refId": "A",
},
Algorithm: map[string]interface{}{
"name": "dbscan",
"config": map[string]interface{}{
"epsilon": 7.667,
},
"sensitivity": 0.83,
},
ResponseType: "binary",
},
appURL: "https://grafana.com",
interval: 1000 * time.Second,
}
t.Run("should generate expected parameters for request", func(t *testing.T) {
to := time.Now().UTC()
from := to.Add(-10 * time.Hour)
called := false
_, err := outlier.Execute(from, to, func(method string, path string, payload []byte) (response.Response, error) {
require.Equal(t, "POST", method)
require.Equal(t, "/proxy/api/v1/outlier", path)
assert.JSONEq(t, fmt.Sprintf(`{
"data": {
"attributes": {
"datasource_type": "prometheus",
"datasource_uid": "a4ce599c-4c93-44b9-be5b-76385b8c01be",
"query_params": {
"expr": "go_goroutines{}",
"range": true,
"refId": "A"
},
"algorithm": {
"config": {
"epsilon": 7.667
},
"name": "dbscan",
"sensitivity": 0.83
},
"response_type": "binary",
"grafana_url": "https://grafana.com",
"start_end_attributes": {
"start": "%s",
"end": "%s",
"interval": 1000000
}
}
}
}`, from.Format(timeFormat), to.Format(timeFormat)), string(payload))
called = true
return nil, nil
})
require.Truef(t, called, "request function was not called")
require.ErrorContains(t, err, "response is nil")
})
successResponse := `{"status":"success","data":{"results":{"A":{"status":200,"frames":[{"schema":{"name":"test","fields":[{"name":"Time","type":"time","typeInfo":{"frame":"time.Time"}},{"name":"Value","type":"number","typeInfo":{"frame":"float64"},"labels":{"instance":"test"}}]},"data":{"values":[[1686945300000],[0]]}}]}}}}`
testCases := []struct {
name string
response response.Response
assert func(t *testing.T, r *backend.QueryDataResponse, err error)
}{
{
name: "should return parsed frames when 200",
response: response.CreateNormalResponse(nil, []byte(successResponse), http.StatusOK),
assert: func(t *testing.T, r *backend.QueryDataResponse, err error) {
require.NoError(t, err)
require.NotNil(t, r)
require.Contains(t, r.Responses, "A")
require.NoError(t, r.Responses["A"].Error)
require.Equal(t, backend.StatusOK, r.Responses["A"].Status)
require.Len(t, r.Responses["A"].Frames, 1)
},
},
{
name: "should return nil if 204",
response: response.CreateNormalResponse(nil, []byte(`{"status": "success"}`), http.StatusNoContent),
assert: func(t *testing.T, r *backend.QueryDataResponse, err error) {
require.NoError(t, err)
require.Nil(t, r)
},
},
{
name: "should return error if any status and body has status error",
response: response.CreateNormalResponse(nil, []byte(`{"status": "error"}`), 0),
assert: func(t *testing.T, r *backend.QueryDataResponse, err error) {
require.ErrorContains(t, err, "ML API responded with error")
require.Nil(t, r)
},
},
{
name: "should return error with explanations if any status and body has status error",
response: response.CreateNormalResponse(nil, []byte(`{"status": "error", "error": "test-error"}`), 0),
assert: func(t *testing.T, r *backend.QueryDataResponse, err error) {
require.ErrorContains(t, err, "ML API responded with error")
require.ErrorContains(t, err, "test-error")
require.Nil(t, r)
},
},
{
name: "should return error response is empty",
response: response.CreateNormalResponse(nil, []byte(``), 0),
assert: func(t *testing.T, r *backend.QueryDataResponse, err error) {
require.ErrorContains(t, err, "unexpected format of the response from ML API")
require.Nil(t, r)
},
},
{
name: "should return error response is not a valid JSON",
response: response.CreateNormalResponse(nil, []byte(`{`), 0),
assert: func(t *testing.T, r *backend.QueryDataResponse, err error) {
require.ErrorContains(t, err, "unexpected format of the response from ML API")
require.Nil(t, r)
},
},
{
name: "should error if response is nil and no error",
response: nil,
assert: func(t *testing.T, r *backend.QueryDataResponse, err error) {
require.ErrorContains(t, err, "response is nil")
require.Nil(t, r)
},
},
}
for _, tc := range testCases {
t.Run(tc.name, func(t *testing.T) {
to := time.Now()
from := to.Add(-10 * time.Hour)
resp, err := outlier.Execute(from, to, func(method string, path string, payload []byte) (response.Response, error) {
return tc.response, nil
})
tc.assert(t, resp, err)
})
}
t.Run("should return error if status is not known", func(t *testing.T) {
knownStatuses := []int{http.StatusOK, http.StatusNoContent}
for status := 100; status < 600; status++ {
if http.StatusText(status) == "" || slices.Contains(knownStatuses, status) {
continue
}
to := time.Now()
from := to.Add(-10 * time.Hour)
resp, err := outlier.Execute(from, to, func(method string, path string, payload []byte) (response.Response, error) {
return response.CreateNormalResponse(nil, []byte(successResponse), status), nil
})
require.ErrorContains(t, err, fmt.Sprintf("unexpected status %d returned by ML API", status))
require.Nil(t, resp)
}
})
t.Run("should propagate error from request function", func(t *testing.T) {
to := time.Now()
from := to.Add(-10 * time.Hour)
resp, err := outlier.Execute(from, to, func(method string, path string, payload []byte) (response.Response, error) {
return nil, errors.New("test-error")
})
require.ErrorContains(t, err, "test-error")
require.Nil(t, resp)
})
}
+44
View File
@@ -0,0 +1,44 @@
package ml
import (
"time"
"github.com/grafana/grafana-plugin-sdk-go/backend"
"github.com/grafana/grafana/pkg/api/response"
)
type FakeCommand struct {
Method string
Path string
Payload []byte
Response *backend.QueryDataResponse
Error error
Recordings []struct {
From time.Time
To time.Time
Response response.Response
Error error
}
}
var _ Command = &FakeCommand{}
func (f *FakeCommand) DatasourceUID() string {
return "fake-ml-datasource"
}
func (f *FakeCommand) Execute(from, to time.Time, executor func(method string, path string, payload []byte) (response.Response, error)) (*backend.QueryDataResponse, error) {
r, err := executor(f.Method, f.Path, f.Payload)
f.Recordings = append(f.Recordings, struct {
From time.Time
To time.Time
Response response.Response
Error error
}{From: from, To: to, Response: r, Error: err})
if err != nil {
return nil, err
}
return f.Response, f.Error
}