Add OTLP exporter for OpenTelemetry (#47987)
* Add OTLP exporter for OpenTelemtry * Fix lint * Refactore parse settings * Add configuration for propagation + fix tests * Fix tests and lint * Fix alerting tests * Add coments to config * Add propagation to custom.ini
This commit is contained in:
@@ -6,11 +6,16 @@ import (
|
||||
"time"
|
||||
|
||||
"github.com/grafana/grafana/pkg/infra/log"
|
||||
"github.com/grafana/grafana/pkg/infra/log/level"
|
||||
"github.com/grafana/grafana/pkg/setting"
|
||||
"go.etcd.io/etcd/version"
|
||||
jaegerpropagator "go.opentelemetry.io/contrib/propagators/jaeger"
|
||||
"go.opentelemetry.io/otel"
|
||||
"go.opentelemetry.io/otel/attribute"
|
||||
"go.opentelemetry.io/otel/codes"
|
||||
"go.opentelemetry.io/otel/exporters/jaeger"
|
||||
"go.opentelemetry.io/otel/exporters/otlp/otlptrace"
|
||||
"go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracegrpc"
|
||||
"go.opentelemetry.io/otel/propagation"
|
||||
"go.opentelemetry.io/otel/sdk/resource"
|
||||
tracesdk "go.opentelemetry.io/otel/sdk/trace"
|
||||
@@ -18,6 +23,13 @@ import (
|
||||
trace "go.opentelemetry.io/otel/trace"
|
||||
)
|
||||
|
||||
const (
|
||||
jaegerExporter string = "jaeger"
|
||||
otlpExporter string = "otlp"
|
||||
jaegerPropagator string = "jaeger"
|
||||
w3cPropagator string = "w3c"
|
||||
)
|
||||
|
||||
type Tracer interface {
|
||||
Run(context.Context) error
|
||||
Start(ctx context.Context, spanName string, opts ...trace.SpanStartOption) (context.Context, Span)
|
||||
@@ -34,9 +46,10 @@ type Span interface {
|
||||
}
|
||||
|
||||
type Opentelemetry struct {
|
||||
enabled bool
|
||||
address string
|
||||
log log.Logger
|
||||
enabled string
|
||||
address string
|
||||
propagation string
|
||||
log log.Logger
|
||||
|
||||
tracerProvider *tracesdk.TracerProvider
|
||||
tracer trace.Tracer
|
||||
@@ -53,6 +66,12 @@ type EventValue struct {
|
||||
Num int64
|
||||
}
|
||||
|
||||
type otelErrHandler func(err error)
|
||||
|
||||
func (o otelErrHandler) Handle(err error) {
|
||||
o(err)
|
||||
}
|
||||
|
||||
func (ots *Opentelemetry) parseSettingsOpentelemetry() error {
|
||||
section, err := ots.Cfg.Raw.GetSection("tracing.opentelemetry.jaeger")
|
||||
if err != nil {
|
||||
@@ -61,13 +80,26 @@ func (ots *Opentelemetry) parseSettingsOpentelemetry() error {
|
||||
|
||||
ots.address = section.Key("address").MustString("")
|
||||
if ots.address != "" {
|
||||
ots.enabled = true
|
||||
ots.enabled = jaegerExporter
|
||||
return nil
|
||||
}
|
||||
ots.propagation = section.Key("propagation").MustString("")
|
||||
|
||||
section, err = ots.Cfg.Raw.GetSection("tracing.opentelemetry.otlp")
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
ots.address = section.Key("address").MustString("")
|
||||
if ots.address != "" {
|
||||
ots.enabled = otlpExporter
|
||||
}
|
||||
ots.propagation = section.Key("propagation").MustString("")
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (ots *Opentelemetry) initTracerProvider() (*tracesdk.TracerProvider, error) {
|
||||
func (ots *Opentelemetry) initJaegerTracerProvider() (*tracesdk.TracerProvider, error) {
|
||||
// Create the Jaeger exporter
|
||||
exp, err := jaeger.New(jaeger.WithCollectorEndpoint(jaeger.WithEndpoint(ots.address)))
|
||||
if err != nil {
|
||||
@@ -86,18 +118,69 @@ func (ots *Opentelemetry) initTracerProvider() (*tracesdk.TracerProvider, error)
|
||||
return tp, nil
|
||||
}
|
||||
|
||||
func (ots *Opentelemetry) initOpentelemetryTracer() error {
|
||||
tp, err := ots.initTracerProvider()
|
||||
func (ots *Opentelemetry) initOTLPTracerProvider() (*tracesdk.TracerProvider, error) {
|
||||
client := otlptracegrpc.NewClient(otlptracegrpc.WithEndpoint(ots.address), otlptracegrpc.WithInsecure())
|
||||
exp, err := otlptrace.New(context.Background(), client)
|
||||
if err != nil {
|
||||
return err
|
||||
return nil, err
|
||||
}
|
||||
|
||||
res, err := resource.New(
|
||||
context.Background(),
|
||||
resource.WithAttributes(
|
||||
semconv.ServiceNameKey.String("grafana"),
|
||||
semconv.ServiceVersionKey.String(version.Version),
|
||||
),
|
||||
resource.WithProcessRuntimeDescription(),
|
||||
resource.WithTelemetrySDK(),
|
||||
)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
tp := tracesdk.NewTracerProvider(
|
||||
tracesdk.WithBatcher(exp),
|
||||
tracesdk.WithSampler(tracesdk.ParentBased(
|
||||
tracesdk.AlwaysSample(),
|
||||
)),
|
||||
tracesdk.WithResource(res),
|
||||
)
|
||||
return tp, nil
|
||||
}
|
||||
|
||||
func (ots *Opentelemetry) initOpentelemetryTracer() error {
|
||||
var tp *tracesdk.TracerProvider
|
||||
var err error
|
||||
switch ots.enabled {
|
||||
case jaegerExporter:
|
||||
tp, err = ots.initJaegerTracerProvider()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
case otlpExporter:
|
||||
tp, err = ots.initOTLPTracerProvider()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
default:
|
||||
ots.log.Error("invalid trace exporter")
|
||||
}
|
||||
|
||||
// Register our TracerProvider as the global so any imported
|
||||
// instrumentation in the future will default to using it
|
||||
// only if tracing is enabled
|
||||
if ots.enabled {
|
||||
if ots.enabled != "" {
|
||||
otel.SetTracerProvider(tp)
|
||||
}
|
||||
|
||||
switch ots.propagation {
|
||||
case w3cPropagator:
|
||||
otel.SetTextMapPropagator(propagation.TraceContext{})
|
||||
case jaegerPropagator:
|
||||
otel.SetTextMapPropagator(jaegerpropagator.Jaeger{})
|
||||
default:
|
||||
otel.SetTextMapPropagator(propagation.TraceContext{})
|
||||
}
|
||||
ots.tracerProvider = tp
|
||||
ots.tracer = otel.GetTracerProvider().Tracer("component-main")
|
||||
|
||||
@@ -105,6 +188,12 @@ func (ots *Opentelemetry) initOpentelemetryTracer() error {
|
||||
}
|
||||
|
||||
func (ots *Opentelemetry) Run(ctx context.Context) error {
|
||||
otel.SetErrorHandler(otelErrHandler(func(err error) {
|
||||
err = level.Error(ots.log).Log("msg", "OpenTelemetry handler returned an error", "err", err)
|
||||
if err != nil {
|
||||
ots.log.Error("OpenTelemetry log returning error", err)
|
||||
}
|
||||
}))
|
||||
<-ctx.Done()
|
||||
|
||||
ots.log.Info("Closing tracing")
|
||||
|
||||
@@ -1,7 +1,9 @@
|
||||
package tracing
|
||||
|
||||
func InitializeTracerForTest() (Tracer, error) {
|
||||
ots := &Opentelemetry{}
|
||||
ots := &Opentelemetry{
|
||||
enabled: "jaeger",
|
||||
}
|
||||
err := ots.initOpentelemetryTracer()
|
||||
if err != nil {
|
||||
return ots, err
|
||||
@@ -10,7 +12,9 @@ func InitializeTracerForTest() (Tracer, error) {
|
||||
}
|
||||
|
||||
func InitializeForBus() Tracer {
|
||||
ots := &Opentelemetry{}
|
||||
ots := &Opentelemetry{
|
||||
enabled: "jaeger",
|
||||
}
|
||||
_ = ots.initOpentelemetryTracer()
|
||||
return ots
|
||||
}
|
||||
|
||||
@@ -28,29 +28,37 @@ const (
|
||||
)
|
||||
|
||||
func ProvideService(cfg *setting.Cfg) (Tracer, error) {
|
||||
ts := &Opentracing{
|
||||
Cfg: cfg,
|
||||
log: log.New("tracing"),
|
||||
}
|
||||
|
||||
if err := ts.parseSettings(); err != nil {
|
||||
ts, ots, err := parseSettings(cfg)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
if ts.enabled {
|
||||
return ts, ts.initGlobalTracer()
|
||||
return ts, ts.initJaegerGlobalTracer()
|
||||
}
|
||||
|
||||
return ots, ots.initOpentelemetryTracer()
|
||||
}
|
||||
|
||||
func parseSettings(cfg *setting.Cfg) (*Opentracing, *Opentelemetry, error) {
|
||||
ts := &Opentracing{
|
||||
Cfg: cfg,
|
||||
log: log.New("tracing"),
|
||||
}
|
||||
err := ts.parseSettings()
|
||||
if err != nil {
|
||||
return ts, nil, err
|
||||
}
|
||||
if ts.enabled {
|
||||
return ts, nil, nil
|
||||
}
|
||||
|
||||
ots := &Opentelemetry{
|
||||
Cfg: cfg,
|
||||
log: log.New("tracing"),
|
||||
}
|
||||
|
||||
if err := ots.parseSettingsOpentelemetry(); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return ots, ots.initOpentelemetryTracer()
|
||||
err = ots.parseSettingsOpentelemetry()
|
||||
return ts, ots, err
|
||||
}
|
||||
|
||||
type traceKey struct{}
|
||||
@@ -136,7 +144,7 @@ func (ts *Opentracing) initJaegerCfg() (jaegercfg.Configuration, error) {
|
||||
return cfg, nil
|
||||
}
|
||||
|
||||
func (ts *Opentracing) initGlobalTracer() error {
|
||||
func (ts *Opentracing) initJaegerGlobalTracer() error {
|
||||
cfg, err := ts.initJaegerCfg()
|
||||
if err != nil {
|
||||
return err
|
||||
|
||||
@@ -27,6 +27,11 @@ import (
|
||||
"github.com/grafana/grafana/pkg/tests/testinfra"
|
||||
)
|
||||
|
||||
type Response struct {
|
||||
Message string `json:"message"`
|
||||
TraceID string `json:"traceID"`
|
||||
}
|
||||
|
||||
func TestAMConfigAccess(t *testing.T) {
|
||||
_, err := tracing.InitializeTracerForTest()
|
||||
require.NoError(t, err)
|
||||
@@ -880,11 +885,11 @@ func TestAlertRuleCRUD(t *testing.T) {
|
||||
// Now, let's try to create some invalid alert rules.
|
||||
{
|
||||
testCases := []struct {
|
||||
desc string
|
||||
rulegroup string
|
||||
interval model.Duration
|
||||
rule apimodels.PostableExtendedRuleNode
|
||||
expectedResponse string
|
||||
desc string
|
||||
rulegroup string
|
||||
interval model.Duration
|
||||
rule apimodels.PostableExtendedRuleNode
|
||||
expectedMessage string
|
||||
}{
|
||||
{
|
||||
desc: "alert rule without queries and expressions",
|
||||
@@ -900,7 +905,7 @@ func TestAlertRuleCRUD(t *testing.T) {
|
||||
Data: []ngmodels.AlertQuery{},
|
||||
},
|
||||
},
|
||||
expectedResponse: `{"message": "invalid rule specification at index [0]: invalid alert rule: no queries or expressions are found", "traceID":"00000000000000000000000000000000"}`,
|
||||
expectedMessage: "invalid rule specification at index [0]: invalid alert rule: no queries or expressions are found",
|
||||
},
|
||||
{
|
||||
desc: "alert rule with empty title",
|
||||
@@ -930,7 +935,7 @@ func TestAlertRuleCRUD(t *testing.T) {
|
||||
},
|
||||
},
|
||||
},
|
||||
expectedResponse: `{"message": "invalid rule specification at index [0]: alert rule title cannot be empty", "traceID":"00000000000000000000000000000000"}`,
|
||||
expectedMessage: "invalid rule specification at index [0]: alert rule title cannot be empty",
|
||||
},
|
||||
{
|
||||
desc: "alert rule with too long name",
|
||||
@@ -960,7 +965,7 @@ func TestAlertRuleCRUD(t *testing.T) {
|
||||
},
|
||||
},
|
||||
},
|
||||
expectedResponse: `{"message": "invalid rule specification at index [0]: alert rule title is too long. Max length is 190", "traceID":"00000000000000000000000000000000"}`,
|
||||
expectedMessage: "invalid rule specification at index [0]: alert rule title is too long. Max length is 190",
|
||||
},
|
||||
{
|
||||
desc: "alert rule with too long rulegroup",
|
||||
@@ -990,7 +995,7 @@ func TestAlertRuleCRUD(t *testing.T) {
|
||||
},
|
||||
},
|
||||
},
|
||||
expectedResponse: `{"message": "rule group name is too long. Max length is 190", "traceID":"00000000000000000000000000000000"}`,
|
||||
expectedMessage: "rule group name is too long. Max length is 190",
|
||||
},
|
||||
{
|
||||
desc: "alert rule with invalid interval",
|
||||
@@ -1021,8 +1026,7 @@ func TestAlertRuleCRUD(t *testing.T) {
|
||||
},
|
||||
},
|
||||
},
|
||||
expectedResponse: `{"message": "rule evaluation interval (1 second) should be positive ` +
|
||||
`number that is multiple of the base interval of 10 seconds", "traceID":"00000000000000000000000000000000"}`,
|
||||
expectedMessage: "rule evaluation interval (1 second) should be positive number that is multiple of the base interval of 10 seconds",
|
||||
},
|
||||
{
|
||||
desc: "alert rule with unknown datasource",
|
||||
@@ -1052,8 +1056,7 @@ func TestAlertRuleCRUD(t *testing.T) {
|
||||
},
|
||||
},
|
||||
},
|
||||
expectedResponse: `{"message": "invalid rule specification at index [0]: failed to validate condition of alert rule AlwaysFiring:` +
|
||||
` invalid query A: data source not found: unknown", "traceID":"00000000000000000000000000000000"}`,
|
||||
expectedMessage: "invalid rule specification at index [0]: failed to validate condition of alert rule AlwaysFiring: invalid query A: data source not found: unknown",
|
||||
},
|
||||
{
|
||||
desc: "alert rule with invalid condition",
|
||||
@@ -1083,8 +1086,7 @@ func TestAlertRuleCRUD(t *testing.T) {
|
||||
},
|
||||
},
|
||||
},
|
||||
expectedResponse: `{"message": "invalid rule specification at index [0]: failed to validate condition of alert rule AlwaysFiring: ` +
|
||||
`condition B not found in any query or expression: it should be one of: [A]", "traceID":"00000000000000000000000000000000"}`,
|
||||
expectedMessage: "invalid rule specification at index [0]: failed to validate condition of alert rule AlwaysFiring: condition B not found in any query or expression: it should be one of: [A]",
|
||||
},
|
||||
}
|
||||
|
||||
@@ -1113,8 +1115,14 @@ func TestAlertRuleCRUD(t *testing.T) {
|
||||
b, err := ioutil.ReadAll(resp.Body)
|
||||
require.NoError(t, err)
|
||||
|
||||
res := &Response{}
|
||||
err = json.Unmarshal(b, &res)
|
||||
require.NoError(t, err)
|
||||
|
||||
assert.Equal(t, res.Message, tc.expectedMessage)
|
||||
assert.NotEmpty(t, res.TraceID)
|
||||
|
||||
assert.Equal(t, resp.StatusCode, http.StatusBadRequest)
|
||||
require.JSONEq(t, tc.expectedResponse, string(b))
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -2266,6 +2274,7 @@ func TestEval(t *testing.T) {
|
||||
payload string
|
||||
expectedStatusCode int
|
||||
expectedResponse string
|
||||
expectedMessage string
|
||||
}{
|
||||
{
|
||||
desc: "alerting condition",
|
||||
@@ -2414,8 +2423,7 @@ func TestEval(t *testing.T) {
|
||||
}
|
||||
`,
|
||||
expectedStatusCode: http.StatusBadRequest,
|
||||
expectedResponse: `{"message": "invalid condition: condition B not found in any query or expression: it should be one of: [A]",` +
|
||||
`"traceID": "00000000000000000000000000000000"}`,
|
||||
expectedMessage: "invalid condition: condition B not found in any query or expression: it should be one of: [A]",
|
||||
},
|
||||
{
|
||||
desc: "unknown query datasource",
|
||||
@@ -2440,7 +2448,7 @@ func TestEval(t *testing.T) {
|
||||
}
|
||||
`,
|
||||
expectedStatusCode: http.StatusBadRequest,
|
||||
expectedResponse: `{"message": "invalid condition: invalid query A: data source not found: unknown", "traceID": "00000000000000000000000000000000"}`,
|
||||
expectedMessage: "invalid condition: invalid query A: data source not found: unknown",
|
||||
},
|
||||
}
|
||||
|
||||
@@ -2457,9 +2465,18 @@ func TestEval(t *testing.T) {
|
||||
})
|
||||
b, err := ioutil.ReadAll(resp.Body)
|
||||
require.NoError(t, err)
|
||||
res := Response{}
|
||||
err = json.Unmarshal(b, &res)
|
||||
require.NoError(t, err)
|
||||
|
||||
assert.Equal(t, tc.expectedStatusCode, resp.StatusCode)
|
||||
require.JSONEq(t, tc.expectedResponse, string(b))
|
||||
if tc.expectedResponse != "" {
|
||||
require.JSONEq(t, tc.expectedResponse, string(b))
|
||||
}
|
||||
if tc.expectedMessage != "" {
|
||||
assert.Equal(t, tc.expectedMessage, res.Message)
|
||||
assert.NotEmpty(t, res.TraceID)
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
@@ -2469,6 +2486,7 @@ func TestEval(t *testing.T) {
|
||||
payload string
|
||||
expectedStatusCode int
|
||||
expectedResponse string
|
||||
expectedMessage string
|
||||
}{
|
||||
{
|
||||
desc: "alerting condition",
|
||||
@@ -2596,8 +2614,7 @@ func TestEval(t *testing.T) {
|
||||
}
|
||||
`,
|
||||
expectedStatusCode: http.StatusBadRequest,
|
||||
expectedResponse: `{"message": "invalid queries or expressions: invalid query A: data source not found: unknown",` +
|
||||
`"traceID": "00000000000000000000000000000000"}`,
|
||||
expectedMessage: "invalid queries or expressions: invalid query A: data source not found: unknown",
|
||||
},
|
||||
}
|
||||
|
||||
@@ -2614,9 +2631,19 @@ func TestEval(t *testing.T) {
|
||||
})
|
||||
b, err := ioutil.ReadAll(resp.Body)
|
||||
require.NoError(t, err)
|
||||
res := Response{}
|
||||
err = json.Unmarshal(b, &res)
|
||||
require.NoError(t, err)
|
||||
|
||||
assert.Equal(t, tc.expectedStatusCode, resp.StatusCode)
|
||||
require.JSONEq(t, tc.expectedResponse, string(b))
|
||||
if tc.expectedResponse != "" {
|
||||
require.JSONEq(t, tc.expectedResponse, string(b))
|
||||
}
|
||||
|
||||
if tc.expectedMessage != "" {
|
||||
require.Equal(t, tc.expectedMessage, res.Message)
|
||||
require.NotEmpty(t, res.TraceID)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
@@ -62,7 +62,10 @@ func TestTestReceivers(t *testing.T) {
|
||||
|
||||
b, err := ioutil.ReadAll(resp.Body)
|
||||
require.NoError(t, err)
|
||||
require.JSONEq(t, `{"traceID":"00000000000000000000000000000000"}`, string(b))
|
||||
res := Response{}
|
||||
err = json.Unmarshal(b, &res)
|
||||
require.NoError(t, err)
|
||||
require.NotEmpty(t, res.TraceID)
|
||||
})
|
||||
|
||||
t.Run("assert working receiver returns OK", func(t *testing.T) {
|
||||
|
||||
Reference in New Issue
Block a user