| package common |
|
|
| import ( |
| "context" |
| "fmt" |
| "log/slog" |
| "net/http" |
| "os" |
| "strings" |
| "time" |
|
|
| health "github.com/AppsFlyer/go-sundheit" |
| healthhttp "github.com/AppsFlyer/go-sundheit/http" |
| "github.com/go-chi/chi/v5" |
| "github.com/go-chi/chi/v5/middleware" |
| "github.com/go-slog/otelslog" |
| "github.com/google/uuid" |
| "github.com/google/wire" |
| "github.com/oklog/ulid/v2" |
| "github.com/prometheus/client_golang/prometheus" |
| "github.com/prometheus/client_golang/prometheus/collectors" |
| "github.com/prometheus/client_golang/prometheus/promhttp" |
| slogmulti "github.com/samber/slog-multi" |
| realotelslog "go.opentelemetry.io/contrib/bridges/otelslog" |
| "go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp" |
| "go.opentelemetry.io/contrib/instrumentation/runtime" |
| "go.opentelemetry.io/otel/log" |
| "go.opentelemetry.io/otel/metric" |
| "go.opentelemetry.io/otel/propagation" |
| sdklog "go.opentelemetry.io/otel/sdk/log" |
| sdkmetric "go.opentelemetry.io/otel/sdk/metric" |
| "go.opentelemetry.io/otel/sdk/resource" |
| sdktrace "go.opentelemetry.io/otel/sdk/trace" |
| semconv "go.opentelemetry.io/otel/semconv/v1.27.0" |
| "go.opentelemetry.io/otel/trace" |
|
|
| "github.com/openmeterio/openmeter/app/config" |
| "github.com/openmeterio/openmeter/openmeter/server" |
| "github.com/openmeterio/openmeter/pkg/contextx" |
| "github.com/openmeterio/openmeter/pkg/framework/transport/httptransport" |
| "github.com/openmeterio/openmeter/pkg/gosundheit" |
| ) |
|
|
| var TelemetryWithoutServer = wire.NewSet( |
| NewTelemetryResource, |
|
|
| NewLoggerProvider, |
| wire.Bind(new(log.LoggerProvider), new(*sdklog.LoggerProvider)), |
| NewLogger, |
|
|
| NewMeterProvider, |
| wire.Bind(new(metric.MeterProvider), new(*sdkmetric.MeterProvider)), |
| NewMeter, |
| NewTracerProvider, |
| wire.Bind(new(trace.TracerProvider), new(*sdktrace.TracerProvider)), |
| NewTracer, |
|
|
| NewRuntimeMetricsCollector, |
| ) |
|
|
| var Telemetry = wire.NewSet( |
| NewTelemetryResource, |
|
|
| NewLoggerProvider, |
| wire.Bind(new(log.LoggerProvider), new(*sdklog.LoggerProvider)), |
| NewLogger, |
|
|
| NewMeterProvider, |
| wire.Bind(new(metric.MeterProvider), new(*sdkmetric.MeterProvider)), |
| NewMeter, |
| NewTracerProvider, |
| wire.Bind(new(trace.TracerProvider), new(*sdktrace.TracerProvider)), |
| NewTracer, |
|
|
| NewHealthChecker, |
|
|
| NewTelemetryHandler, |
| NewTelemetryServer, |
|
|
| NewRuntimeMetricsCollector, |
| ) |
|
|
| |
| |
| |
| func init() { |
| slog.SetDefault(slog.New(config.NewJSONHandler(os.Stderr, slog.LevelInfo))) |
| } |
|
|
| const ( |
| DefaultShutdownTimeout = 5 * time.Second |
| ) |
|
|
| func NewTelemetryResource(metadata Metadata) *resource.Resource { |
| extraResources, _ := resource.New( |
| |
| context.Background(), |
| resource.WithContainer(), |
| resource.WithAttributes( |
| semconv.ServiceName(metadata.ServiceName), |
| semconv.ServiceVersion(metadata.Version), |
| semconv.DeploymentEnvironmentName(metadata.Environment), |
| ), |
| ) |
|
|
| res, _ := resource.Merge( |
| resource.Default(), |
| extraResources, |
| ) |
|
|
| return res |
| } |
|
|
| func NewLoggerProvider(ctx context.Context, conf config.LogTelemetryConfig, res *resource.Resource) (*sdklog.LoggerProvider, func(), error) { |
| loggerProvider, err := conf.NewLoggerProvider(ctx, res) |
| if err != nil { |
| return nil, nil, fmt.Errorf("failed to initialize OpenTelemetry Trace provider: %w", err) |
| } |
|
|
| return loggerProvider, func() { |
| |
| |
| ctx, cancel := context.WithTimeout(context.Background(), DefaultShutdownTimeout) |
| defer cancel() |
|
|
| if err := loggerProvider.ForceFlush(ctx); err != nil { |
| |
| slog.Error("flushing logger provider", slog.Any("error", err)) |
| } |
|
|
| if err := loggerProvider.Shutdown(ctx); err != nil { |
| |
| slog.Error("shutting down logger provider", slog.Any("error", err)) |
| } |
| }, nil |
| } |
|
|
| func NewLogger(conf config.LogTelemetryConfig, res *resource.Resource, loggerProvider log.LoggerProvider, metadata Metadata, additionalMiddlewares []slogmulti.Middleware) *slog.Logger { |
| baseMiddlewares := []slogmulti.Middleware{ |
| otelslog.ResourceMiddleware(res), |
| otelslog.NewHandler, |
| } |
|
|
| baseMiddlewares = append(baseMiddlewares, additionalMiddlewares...) |
|
|
| |
| stdoutLogger := slogmulti. |
| Pipe(baseMiddlewares...). |
| Handler(conf.NewHandler(os.Stdout)) |
|
|
| |
| |
| otelLogger := NewLevelHandler( |
| realotelslog.NewHandler(metadata.OpenTelemetryName, realotelslog.WithLoggerProvider(loggerProvider)), |
| conf.Level, |
| ) |
|
|
| |
| out := slogmulti.Fanout( |
| stdoutLogger, |
| otelLogger, |
| ) |
|
|
| |
| middlewares := slogmulti.Pipe( |
| contextx.NewLogHandler, |
| ) |
|
|
| return slog.New(middlewares.Handler(out)) |
| } |
|
|
| func TelemetryLoggerNoAdditionalMiddlewares() []slogmulti.Middleware { |
| return nil |
| } |
|
|
| func NewMeterProvider(ctx context.Context, conf config.MetricsTelemetryConfig, res *resource.Resource, logger *slog.Logger) (*sdkmetric.MeterProvider, func(), error) { |
| meterProvider, err := conf.NewMeterProvider(ctx, res) |
| if err != nil { |
| return nil, nil, fmt.Errorf("failed to initialize OpenTelemetry Metrics provider: %w", err) |
| } |
|
|
| return meterProvider, func() { |
| |
| |
| ctx, cancel := context.WithTimeout(context.Background(), DefaultShutdownTimeout) |
| defer cancel() |
|
|
| if err := meterProvider.Shutdown(ctx); err != nil { |
| logger.Error("shutting down meter provider", slog.Any("error", err)) |
| } |
| }, nil |
| } |
|
|
| func NewMeter(meterProvider metric.MeterProvider, metadata Metadata) metric.Meter { |
| return meterProvider.Meter(metadata.OpenTelemetryName) |
| } |
|
|
| func NewTracerProvider(ctx context.Context, conf config.TraceTelemetryConfig, res *resource.Resource, logger *slog.Logger) (*sdktrace.TracerProvider, func(), error) { |
| tracerProvider, err := conf.NewTracerProvider(ctx, res) |
| if err != nil { |
| return nil, nil, fmt.Errorf("failed to initialize OpenTelemetry Trace provider: %w", err) |
| } |
|
|
| |
| httptransport.EnableOperationSpans(conf.OperationSpans) |
|
|
| return tracerProvider, func() { |
| |
| |
| ctx, cancel := context.WithTimeout(context.Background(), DefaultShutdownTimeout) |
| defer cancel() |
|
|
| if err := tracerProvider.Shutdown(ctx); err != nil { |
| logger.Error("shutting down tracer provider", slog.Any("error", err)) |
| } |
| }, nil |
| } |
|
|
| func NewTracer(tracerProvider trace.TracerProvider, metadata Metadata) trace.Tracer { |
| return tracerProvider.Tracer(metadata.OpenTelemetryName) |
| } |
|
|
| func NewDefaultTextMapPropagator() propagation.TextMapPropagator { |
| return propagation.TraceContext{} |
| } |
|
|
| func NewHealthChecker(logger *slog.Logger) health.Health { |
| return health.New(health.WithCheckListeners(gosundheit.NewLogger(logger.With(slog.String("component", "healthcheck"))))) |
| } |
|
|
| type TelemetryHandler http.Handler |
|
|
| func NewTelemetryHandler( |
| metricsConf config.MetricsTelemetryConfig, |
| healthChecker health.Health, |
| logger *slog.Logger, |
| ) TelemetryHandler { |
| if metricsConf.Exporters.Prometheus.DisableDefaultCollectors { |
| prometheus.Unregister(collectors.NewGoCollector()) |
| prometheus.Unregister(collectors.NewProcessCollector(collectors.ProcessCollectorOpts{})) |
| } |
|
|
| router := chi.NewRouter() |
| router.Mount("/debug", middleware.Profiler()) |
|
|
| if metricsConf.Exporters.Prometheus.Enabled { |
| router.Handle("/metrics", promhttp.Handler()) |
| } |
|
|
| |
| { |
| handler := healthhttp.HandleHealthJSON(healthChecker) |
| router.Handle("/healthz", handler) |
|
|
| |
| router.HandleFunc("/healthz/live", func(w http.ResponseWriter, _ *http.Request) { |
| _, _ = w.Write([]byte("ok")) |
| }) |
| router.Handle("/healthz/ready", handler) |
| } |
|
|
| return router |
| } |
|
|
| type TelemetryServer = *http.Server |
|
|
| func NewTelemetryServer(conf config.TelemetryConfig, handler TelemetryHandler) (TelemetryServer, func()) { |
| server := &http.Server{ |
| Addr: conf.Address, |
| Handler: handler, |
| ReadHeaderTimeout: conf.ReadHeaderTimeout, |
| } |
|
|
| return server, func() { server.Close() } |
| } |
|
|
| type TelemetryMiddlewareHook server.MiddlewareHook |
|
|
| func NewTelemetryRouterHook(meterProvider metric.MeterProvider, tracerProvider trace.TracerProvider) TelemetryMiddlewareHook { |
| return func(m server.MiddlewareManager) { |
| m.Use(func(h http.Handler) http.Handler { |
| return otelhttp.NewHandler( |
| http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { |
| h.ServeHTTP(w, r) |
|
|
| routePattern := chi.RouteContext(r.Context()).RoutePattern() |
|
|
| |
| |
| |
| |
| |
| span := trace.SpanFromContext(r.Context()) |
| span.SetAttributes(semconv.URLPath(r.URL.Path), semconv.HTTPRoute(routePattern)) |
|
|
| if labeler, ok := otelhttp.LabelerFromContext(r.Context()); ok { |
| labeler.Add(semconv.HTTPRoute(routePattern)) |
| } |
| }), |
| "", |
| otelhttp.WithMeterProvider(meterProvider), |
| otelhttp.WithTracerProvider(tracerProvider), |
| otelhttp.WithSpanNameFormatter(func(_ string, r *http.Request) string { |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| route := r.Pattern |
| if route == "" { |
| if rctx := chi.RouteContext(r.Context()); rctx != nil { |
| route = rctx.RoutePattern() |
| } |
| } |
| if route == "" || strings.Contains(route, "*") { |
| route = lowCardinalityPath(r.URL.Path) |
| } |
| return r.Method + " " + route |
| }), |
| ) |
| }) |
| } |
| } |
|
|
| const ( |
| |
| |
| |
| maxRouteSegmentLen = 32 |
| |
| |
| maxRouteSegments = 12 |
| ) |
|
|
| |
| |
| |
| |
| |
| |
| |
| func lowCardinalityPath(path string) string { |
| segments := strings.Split(path, "/") |
|
|
| truncated := false |
| if len(segments) > maxRouteSegments { |
| segments = segments[:maxRouteSegments] |
| truncated = true |
| } |
|
|
| for i, seg := range segments { |
| if isHighCardinalitySegment(seg) { |
| segments[i] = ":id" |
| } |
| } |
|
|
| out := strings.Join(segments, "/") |
| if truncated { |
| out += "/..." |
| } |
| return out |
| } |
|
|
| |
| |
| func isHighCardinalitySegment(seg string) bool { |
| if len(seg) > maxRouteSegmentLen { |
| return true |
| } |
| if isAllDigits(seg) { |
| return true |
| } |
| if _, err := ulid.ParseStrict(seg); err == nil { |
| return true |
| } |
| if _, err := uuid.Parse(seg); err == nil { |
| return true |
| } |
| return false |
| } |
|
|
| func isAllDigits(s string) bool { |
| if s == "" { |
| return false |
| } |
| for _, r := range s { |
| if r < '0' || r > '9' { |
| return false |
| } |
| } |
| return true |
| } |
|
|
| type RuntimeMetricsCollector struct{} |
|
|
| func NewRuntimeMetricsCollector( |
| meterProvider metric.MeterProvider, |
| logger *slog.Logger, |
| ) (RuntimeMetricsCollector, error) { |
| err := runtime.Start( |
| runtime.WithMinimumReadMemStatsInterval(time.Second), |
| runtime.WithMeterProvider(meterProvider), |
| ) |
| if err != nil { |
| return RuntimeMetricsCollector{}, fmt.Errorf("starting runtime metrics: %w", err) |
| } |
|
|
| logger.Debug("started collecting runtime metrics") |
| return RuntimeMetricsCollector{}, nil |
| } |
|
|
| |
| var _ slog.Handler = (*LevelHandler)(nil) |
|
|
| |
| func NewLevelHandler(handler slog.Handler, level slog.Leveler) *LevelHandler { |
| return &LevelHandler{ |
| handler: handler, |
| level: level, |
| } |
| } |
|
|
| |
| type LevelHandler struct { |
| handler slog.Handler |
| level slog.Leveler |
| } |
|
|
| func (h *LevelHandler) Enabled(ctx context.Context, level slog.Level) bool { |
| |
| return level >= h.level.Level() && h.handler.Enabled(ctx, level) |
| } |
|
|
| func (h *LevelHandler) WithGroup(name string) slog.Handler { |
| return NewLevelHandler(h.handler.WithGroup(name), h.level) |
| } |
|
|
| func (h *LevelHandler) WithAttrs(attrs []slog.Attr) slog.Handler { |
| return NewLevelHandler(h.handler.WithAttrs(attrs), h.level) |
| } |
|
|
| func (h *LevelHandler) Handle(ctx context.Context, record slog.Record) error { |
| return h.handler.Handle(ctx, record) |
| } |
|
|