feat(pkger): add service logging and tracing middlewares

pull/16553/head
Johnny Steenbergen 2020-01-15 11:00:31 -08:00 committed by Johnny Steenbergen
parent 540e785eb4
commit 93c8a2a104
4 changed files with 145 additions and 0 deletions

View File

@ -827,6 +827,8 @@ func (m *Launcher) run(ctx context.Context) (err error) {
pkger.WithTelegrafSVC(authorizer.NewTelegrafConfigService(b.TelegrafService, b.UserResourceMappingService)),
pkger.WithVariableSVC(authorizer.NewVariableService(b.VariableService)),
)
pkgSVC = pkger.MWTracing()(pkgSVC)
pkgSVC = pkger.MWLogging(pkgerLogger)(pkgSVC)
}
var pkgHTTPServer *http.HandlerPkg

View File

@ -24,6 +24,9 @@ type SVC interface {
Apply(ctx context.Context, orgID, userID influxdb.ID, pkg *Pkg, opts ...ApplyOptFn) (Summary, error)
}
// SVCMiddleware is a service middleware func.
type SVCMiddleware func(SVC) SVC
type serviceOpt struct {
logger *zap.Logger

99
pkger/service_logging.go Normal file
View File

@ -0,0 +1,99 @@
package pkger
import (
"context"
"time"
"github.com/influxdata/influxdb"
"go.uber.org/zap"
)
type loggingMW struct {
logger *zap.Logger
next SVC
}
// MWLogging adds logging functionality for the service.
func MWLogging(log *zap.Logger) SVCMiddleware {
return func(svc SVC) SVC {
return &loggingMW{
logger: log,
next: svc,
}
}
}
var _ SVC = (*loggingMW)(nil)
func (s *loggingMW) CreatePkg(ctx context.Context, setters ...CreatePkgSetFn) (pkg *Pkg, err error) {
defer func(start time.Time) {
dur := zap.Duration("took", time.Since(start))
if err != nil {
s.logger.Error("failed to create pkg", zap.Error(err), dur)
return
}
s.logger.Info("pkg create", append(s.summaryLogFields(pkg.Summary()), dur)...)
}(time.Now())
return s.next.CreatePkg(ctx, setters...)
}
func (s *loggingMW) DryRun(ctx context.Context, orgID, userID influxdb.ID, pkg *Pkg) (sum Summary, diff Diff, err error) {
defer func(start time.Time) {
dur := zap.Duration("took", time.Since(start))
if err != nil {
s.logger.Error("failed to dry run pkg",
zap.String("orgID", orgID.String()),
zap.String("userID", userID.String()),
zap.Error(err),
dur,
)
return
}
s.logger.Info("pkg dry run successful", append(s.summaryLogFields(sum), dur)...)
}(time.Now())
return s.next.DryRun(ctx, orgID, userID, pkg)
}
func (s *loggingMW) Apply(ctx context.Context, orgID, userID influxdb.ID, pkg *Pkg, opts ...ApplyOptFn) (sum Summary, err error) {
defer func(start time.Time) {
dur := zap.Duration("took", time.Since(start))
if err != nil {
s.logger.Error("failed to apply pkg",
zap.String("orgID", orgID.String()),
zap.String("userID", userID.String()),
zap.Error(err),
dur,
)
}
s.logger.Info("pkg apply successful", append(s.summaryLogFields(sum), dur)...)
}(time.Now())
return s.next.Apply(ctx, orgID, userID, pkg, opts...)
}
func (s *loggingMW) summaryLogFields(sum Summary) []zap.Field {
potentialFields := []struct {
key string
val int
}{
{key: "buckets", val: len(sum.Buckets)},
{key: "checks", val: len(sum.Checks)},
{key: "dashboards", val: len(sum.Dashboards)},
{key: "endpoints", val: len(sum.NotificationEndpoints)},
{key: "labels", val: len(sum.Labels)},
{key: "label_mappings", val: len(sum.LabelMappings)},
{key: "rules", val: len(sum.NotificationRules)},
{key: "secrets", val: len(sum.MissingSecrets)},
{key: "tasks", val: len(sum.Tasks)},
{key: "telegrafs", val: len(sum.TelegrafConfigs)},
{key: "variables", val: len(sum.Variables)},
}
var fields []zap.Field
for _, f := range potentialFields {
if f.val > 0 {
fields = append(fields, zap.Int("num_"+f.key, f.val))
}
}
return fields
}

41
pkger/service_tracing.go Normal file
View File

@ -0,0 +1,41 @@
package pkger
import (
"context"
"github.com/influxdata/influxdb"
"github.com/influxdata/influxdb/kit/tracing"
)
type traceMW struct {
next SVC
}
// MWTracing adds tracing functionality for the service.
func MWTracing() SVCMiddleware {
return func(svc SVC) SVC {
return &traceMW{next: svc}
}
}
var _ SVC = (*traceMW)(nil)
func (s *traceMW) CreatePkg(ctx context.Context, setters ...CreatePkgSetFn) (pkg *Pkg, err error) {
span, ctx := tracing.StartSpanFromContextWithOperationName(ctx, "CreatePkg")
defer span.Finish()
return s.next.CreatePkg(ctx, setters...)
}
func (s *traceMW) DryRun(ctx context.Context, orgID, userID influxdb.ID, pkg *Pkg) (sum Summary, diff Diff, err error) {
span, ctx := tracing.StartSpanFromContextWithOperationName(ctx, "DryRun")
span.LogKV("orgID", orgID.String(), "userID", userID.String())
defer span.Finish()
return s.next.DryRun(ctx, orgID, userID, pkg)
}
func (s *traceMW) Apply(ctx context.Context, orgID, userID influxdb.ID, pkg *Pkg, opts ...ApplyOptFn) (sum Summary, err error) {
span, ctx := tracing.StartSpanFromContextWithOperationName(ctx, "Apply")
span.LogKV("orgID", orgID.String(), "userID", userID.String())
defer span.Finish()
return s.next.Apply(ctx, orgID, userID, pkg, opts...)
}