-
Notifications
You must be signed in to change notification settings - Fork 41
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
Signed-off-by: Vyom Yadav <[email protected]>
- Loading branch information
1 parent
de9b376
commit b92b0de
Showing
5 changed files
with
201 additions
and
6 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,80 @@ | ||
package metrics | ||
|
||
import ( | ||
"context" | ||
"go.opentelemetry.io/otel/metric" | ||
) | ||
|
||
// Default bucket boundaries in seconds for the send delay histogram | ||
var sendDelayBuckets = []float64{ | ||
0, // immediate | ||
10, // 10 seconds | ||
20, // 20 seconds | ||
40, // 40 seconds | ||
80, // 1m 20s | ||
160, // 2m 40s | ||
320, // 5m 20s | ||
640, // 10m 40s | ||
1280, // 21m 20s | ||
} | ||
|
||
type Metrics struct { | ||
// Time between when a reminder became eligible and when it was sent | ||
SendDelay metric.Float64Histogram | ||
|
||
// Current number of reminders in the batch | ||
BatchSize metric.Int64Gauge | ||
|
||
// Average batch size (updated on each batch) | ||
AvgBatchSize metric.Float64Gauge | ||
|
||
// For tracking average calculation | ||
// TODO: consider persisting this to avoid reset on restart (maybe) | ||
totalBatches int64 | ||
totalReminders int64 | ||
} | ||
|
||
func NewMetrics(meter metric.Meter) (*Metrics, error) { | ||
sendDelay, err := meter.Float64Histogram( | ||
"reminder_send_delay", | ||
metric.WithDescription("Time between reminder becoming eligible and actual send (seconds)"), | ||
metric.WithUnit("s"), | ||
metric.WithExplicitBucketBoundaries(sendDelayBuckets...), | ||
) | ||
if err != nil { | ||
return nil, err | ||
} | ||
|
||
batchSize, err := meter.Int64Gauge( | ||
"reminder_batch_size", | ||
metric.WithDescription("Current number of reminders in the batch"), | ||
) | ||
if err != nil { | ||
return nil, err | ||
} | ||
|
||
avgBatchSize, err := meter.Float64Gauge( | ||
"reminder_avg_batch_size", | ||
metric.WithDescription("Average number of reminders per batch"), | ||
) | ||
if err != nil { | ||
return nil, err | ||
} | ||
|
||
return &Metrics{ | ||
SendDelay: sendDelay, | ||
BatchSize: batchSize, | ||
AvgBatchSize: avgBatchSize, | ||
}, nil | ||
} | ||
|
||
func (m *Metrics) RecordBatch(ctx context.Context, size int64) { | ||
// Update current batch size | ||
m.BatchSize.Record(ctx, size) | ||
|
||
// Update running average | ||
m.totalBatches++ | ||
m.totalReminders += size | ||
avgSize := float64(m.totalReminders) / float64(m.totalBatches) | ||
m.AvgBatchSize.Record(ctx, avgSize) | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,80 @@ | ||
package reminder | ||
|
||
import ( | ||
"context" | ||
"fmt" | ||
"github.com/prometheus/client_golang/prometheus/promhttp" | ||
"github.com/rs/zerolog" | ||
"go.opentelemetry.io/otel" | ||
"go.opentelemetry.io/otel/exporters/prometheus" | ||
sdkmetric "go.opentelemetry.io/otel/sdk/metric" | ||
"go.opentelemetry.io/otel/sdk/resource" | ||
semconv "go.opentelemetry.io/otel/semconv/v1.17.0" | ||
"net/http" | ||
"time" | ||
) | ||
|
||
const ( | ||
metricsPath = "/metrics" | ||
readHeaderTimeout = 2 * time.Second | ||
) | ||
|
||
func (r *reminder) startMetricServer(ctx context.Context) error { | ||
logger := zerolog.Ctx(ctx) | ||
|
||
prometheusExporter, err := prometheus.New( | ||
prometheus.WithNamespace("reminder"), | ||
) | ||
if err != nil { | ||
return fmt.Errorf("failed to create Prometheus exporter: %w", err) | ||
} | ||
|
||
res := resource.NewWithAttributes( | ||
semconv.SchemaURL, | ||
semconv.ServiceName("reminder"), | ||
// TODO: Make this auto-generated | ||
semconv.ServiceVersion("v0.1.0"), | ||
) | ||
|
||
mp := sdkmetric.NewMeterProvider( | ||
sdkmetric.WithReader(prometheusExporter), | ||
sdkmetric.WithResource(res), | ||
) | ||
|
||
otel.SetMeterProvider(mp) | ||
|
||
mux := http.NewServeMux() | ||
mux.Handle(metricsPath, promhttp.Handler()) | ||
|
||
server := &http.Server{ | ||
Addr: fmt.Sprintf("%s:%d", r.cfg.MetricsConfig.Host, r.cfg.MetricsConfig.Port), | ||
Handler: mux, | ||
ReadHeaderTimeout: readHeaderTimeout, | ||
} | ||
|
||
logger.Info().Msgf("starting metrics server on %s", server.Addr) | ||
|
||
errCh := make(chan error) | ||
go func() { | ||
errCh <- server.ListenAndServe() | ||
}() | ||
|
||
select { | ||
case err := <-errCh: | ||
return err | ||
case <-ctx.Done(): | ||
case <-r.stop: | ||
} | ||
|
||
// shutdown the metrics server when either the context is done or when reminder is stopped | ||
shutdownCtx, shutdownRelease := context.WithTimeout(context.Background(), 5*time.Second) | ||
defer shutdownRelease() | ||
|
||
logger.Info().Msg("shutting down metrics server") | ||
|
||
if err := mp.Shutdown(shutdownCtx); err != nil { | ||
logger.Err(err).Msg("error shutting down metrics provider") | ||
} | ||
|
||
return server.Shutdown(shutdownCtx) | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,7 @@ | ||
package reminder | ||
|
||
type MetricsConfig struct { | ||
Enabled bool `mapstructure:"enabled" default:"true"` | ||
Host string `mapstructure:"host" default:"127.0.0.1"` | ||
Port int `mapstructure:"port" default:"8080"` | ||
} |