diff --git a/internal/component/common/loki/client/consumer_wal.go b/internal/component/common/loki/client/consumer_wal.go index 5641c047116..184f438d931 100644 --- a/internal/component/common/loki/client/consumer_wal.go +++ b/internal/component/common/loki/client/consumer_wal.go @@ -21,7 +21,7 @@ func NewWALConsumer(logger *slog.Logger, reg prometheus.Registerer, walCfg wal.C return nil, fmt.Errorf("at least one endpoint config must be provided") } - writer, err := wal.NewWriter(walCfg, logger, reg) + writer, err := wal.NewWriter(walCfg, logger, reg, wal.NewWriterMetrics(reg)) if err != nil { return nil, fmt.Errorf("error creating wal writer: %w", err) } diff --git a/internal/component/common/loki/wal/writer.go b/internal/component/common/loki/wal/writer.go index f6de1262730..4d1ff03436f 100644 --- a/internal/component/common/loki/wal/writer.go +++ b/internal/component/common/loki/wal/writer.go @@ -55,15 +55,13 @@ type Writer struct { writeSubscribersLock sync.RWMutex writeSubscribers []WriteEventSubscriber - reclaimedOldSegmentsSpaceCounter *prometheus.CounterVec - lastReclaimedSegment *prometheus.GaugeVec - lastWrittenTimestamp *prometheus.GaugeVec + metrics *WriterMetrics closeCleaner chan struct{} } // NewWriter creates a new Writer. -func NewWriter(walCfg Config, logger *slog.Logger, reg prometheus.Registerer) (*Writer, error) { +func NewWriter(walCfg Config, logger *slog.Logger, reg prometheus.Registerer, metrics *WriterMetrics) (*Writer, error) { // Start WAL wl, err := New(Config{ Dir: walCfg.Dir, @@ -80,32 +78,7 @@ func NewWriter(walCfg Config, logger *slog.Logger, reg prometheus.Registerer) (* wal: wl, entryWriter: newEntryWriter(), closeCleaner: make(chan struct{}, 1), - } - - wrt.reclaimedOldSegmentsSpaceCounter = prometheus.NewCounterVec(prometheus.CounterOpts{ - Namespace: "loki_write", - Subsystem: "wal_writer", - Name: "reclaimed_space", - Help: "Number of bytes reclaimed from storage.", - }, []string{}) - - wrt.lastReclaimedSegment = prometheus.NewGaugeVec(prometheus.GaugeOpts{ - Namespace: "loki_write", - Subsystem: "wal_writer", - Name: "last_reclaimed_segment", - Help: "Last reclaimed segment number", - }, []string{}) - wrt.lastWrittenTimestamp = prometheus.NewGaugeVec(prometheus.GaugeOpts{ - Namespace: "loki_write", - Subsystem: "wal_writer", - Name: "last_written_timestamp", - Help: "Latest timestamp that was written to the WAL", - }, []string{}) - - if reg != nil { - _ = reg.Register(wrt.reclaimedOldSegmentsSpaceCounter) - _ = reg.Register(wrt.lastReclaimedSegment) - _ = reg.Register(wrt.lastWrittenTimestamp) + metrics: metrics, } return wrt, nil @@ -122,7 +95,7 @@ func (wrt *Writer) Start(maxSegmentAge time.Duration) { } // emit metric with latest written timestamp, to be able to track delay from writer to watcher - wrt.lastWrittenTimestamp.WithLabelValues().Set(float64(e.Timestamp.Unix())) + wrt.metrics.lastWrittenTimestamp.WithLabelValues().Set(float64(e.Timestamp.Unix())) wrt.writeSubscribersLock.RLock() for _, s := range wrt.writeSubscribers { @@ -202,7 +175,7 @@ func (wrt *Writer) cleanSegments(maxAge time.Duration) error { wrt.logger.Error("Error old wal segment", "err", err, "segmentNum", segment.number) } wrt.logger.Debug("Deleted old wal segment", "segmentNum", segment.number) - wrt.reclaimedOldSegmentsSpaceCounter.WithLabelValues().Add(float64(segment.size)) + wrt.metrics.reclaimedOldSegmentsSpaceCounter.WithLabelValues().Add(float64(segment.size)) // keep track of the largest segment number reclaimed if segment.number > maxReclaimed { maxReclaimed = segment.number @@ -216,7 +189,7 @@ func (wrt *Writer) cleanSegments(maxAge time.Duration) error { for _, subscriber := range wrt.cleanupSubscribers { subscriber.SeriesReset(maxReclaimed) } - wrt.lastReclaimedSegment.WithLabelValues().Set(float64(maxReclaimed)) + wrt.metrics.lastReclaimedSegment.WithLabelValues().Set(float64(maxReclaimed)) } return nil } diff --git a/internal/component/common/loki/wal/writer_metrics.go b/internal/component/common/loki/wal/writer_metrics.go new file mode 100644 index 00000000000..7a46dffe7a5 --- /dev/null +++ b/internal/component/common/loki/wal/writer_metrics.go @@ -0,0 +1,43 @@ +package wal + +import ( + "github.com/grafana/alloy/internal/util" + "github.com/prometheus/client_golang/prometheus" +) + +type WriterMetrics struct { + lastReclaimedSegment *prometheus.GaugeVec + lastWrittenTimestamp *prometheus.GaugeVec + reclaimedOldSegmentsSpaceCounter *prometheus.CounterVec +} + +func NewWriterMetrics(reg prometheus.Registerer) *WriterMetrics { + m := &WriterMetrics{ + lastReclaimedSegment: prometheus.NewGaugeVec(prometheus.GaugeOpts{ + Namespace: "loki_write", + Subsystem: "wal_writer", + Name: "last_reclaimed_segment", + Help: "Last reclaimed segment number", + }, []string{}), + lastWrittenTimestamp: prometheus.NewGaugeVec(prometheus.GaugeOpts{ + Namespace: "loki_write", + Subsystem: "wal_writer", + Name: "last_written_timestamp", + Help: "Latest timestamp that was written to the WAL", + }, []string{}), + reclaimedOldSegmentsSpaceCounter: prometheus.NewCounterVec(prometheus.CounterOpts{ + Namespace: "loki_write", + Subsystem: "wal_writer", + Name: "reclaimed_space", + Help: "Number of bytes reclaimed from storage.", + }, []string{}), + } + + if reg != nil { + m.lastReclaimedSegment = util.MustRegisterOrGet(reg, m.lastReclaimedSegment).(*prometheus.GaugeVec) + m.lastWrittenTimestamp = util.MustRegisterOrGet(reg, m.lastWrittenTimestamp).(*prometheus.GaugeVec) + m.reclaimedOldSegmentsSpaceCounter = util.MustRegisterOrGet(reg, m.reclaimedOldSegmentsSpaceCounter).(*prometheus.CounterVec) + } + + return m +} diff --git a/internal/component/common/loki/wal/writer_test.go b/internal/component/common/loki/wal/writer_test.go index 534ca2a2b75..c6a4e1d09be 100644 --- a/internal/component/common/loki/wal/writer_test.go +++ b/internal/component/common/loki/wal/writer_test.go @@ -4,13 +4,16 @@ import ( "fmt" "os" "path/filepath" + "strings" "sync" "testing" "time" "github.com/grafana/loki/pkg/push" "github.com/prometheus/client_golang/prometheus" + "github.com/prometheus/client_golang/prometheus/testutil" "github.com/prometheus/common/model" + "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" "github.com/grafana/alloy/internal/component/common/loki" @@ -20,13 +23,16 @@ import ( ) func TestWriter_EntriesAreWrittenToWAL(t *testing.T) { - dir := t.TempDir() + var ( + dir = t.TempDir() + reg = prometheus.NewRegistry() + ) writer, err := NewWriter(Config{ Dir: dir, Enabled: true, MaxSegmentAge: time.Minute, - }, autil.TestAlloyLogger(t).Slog(), prometheus.NewRegistry()) + }, autil.TestAlloyLogger(t).Slog(), reg, NewWriterMetrics(reg)) require.NoError(t, err) defer func() { writer.Stop() @@ -62,6 +68,58 @@ func TestWriter_EntriesAreWrittenToWAL(t *testing.T) { require.Equal(t, testLabels, readEntries[0].Labels) } +func TestWriter_MetricsWorkAfterRecreation(t *testing.T) { + var ( + dir = t.TempDir() + reg = prometheus.NewRegistry() + ) + + writer, err := NewWriter(Config{ + Dir: dir, + Enabled: true, + MaxSegmentAge: time.Minute, + }, logging.NewSlogNop(), reg, NewWriterMetrics(reg)) + require.NoError(t, err) + + writer.Start(time.Minute) + entry := loki.NewEntry(model.LabelSet{"foo": "bar"}, push.Entry{Timestamp: time.Now(), Line: "line"}) + writer.Chan() <- entry + + expected := fmt.Sprintf(` + # HELP loki_write_wal_writer_last_written_timestamp Latest timestamp that was written to the WAL + # TYPE loki_write_wal_writer_last_written_timestamp gauge + loki_write_wal_writer_last_written_timestamp %d + `, entry.Timestamp.Unix()) + + require.EventuallyWithT(t, func(c *assert.CollectT) { + require.NoError(c, testutil.GatherAndCompare(reg, strings.NewReader(expected), "loki_write_wal_writer_last_written_timestamp")) + }, 2*time.Second, 100*time.Millisecond) + + writer.Stop() + + writer, err = NewWriter(Config{ + Dir: dir, + Enabled: true, + MaxSegmentAge: time.Minute, + }, logging.NewSlogNop(), reg, NewWriterMetrics(reg)) + require.NoError(t, err) + writer.Start(time.Minute) + defer writer.Stop() + + newEntry := loki.NewEntry(model.LabelSet{"foo": "bar"}, push.Entry{Timestamp: time.Now().Add(1 * time.Second), Line: "line"}) + writer.Chan() <- newEntry + + expected = fmt.Sprintf(` + # HELP loki_write_wal_writer_last_written_timestamp Latest timestamp that was written to the WAL + # TYPE loki_write_wal_writer_last_written_timestamp gauge + loki_write_wal_writer_last_written_timestamp %d + `, newEntry.Timestamp.Unix()) + + require.EventuallyWithT(t, func(c *assert.CollectT) { + require.NoError(c, testutil.GatherAndCompare(reg, strings.NewReader(expected), "loki_write_wal_writer_last_written_timestamp")) + }, 2*time.Second, 100*time.Millisecond) +} + type notifySegmentsCleanedFunc func(num int) func (n notifySegmentsCleanedFunc) NotifyWrite() { @@ -72,18 +130,19 @@ func (n notifySegmentsCleanedFunc) SeriesReset(segmentNum int) { } func TestWriter_OldSegmentsAreCleanedUp(t *testing.T) { - dir := t.TempDir() - - maxSegmentAge := time.Second * 2 - - subscriber1 := []int{} - subscriber2 := []int{} + var ( + dir = t.TempDir() + reg = prometheus.NewRegistry() + maxSegmentAge = time.Second * 2 + subscriber1 = []int{} + subscriber2 = []int{} + ) writer, err := NewWriter(Config{ Dir: dir, Enabled: true, MaxSegmentAge: maxSegmentAge, - }, autil.TestAlloyLogger(t).Slog(), prometheus.NewRegistry()) + }, autil.TestAlloyLogger(t).Slog(), reg, NewWriterMetrics(reg)) require.NoError(t, err) writer.Start(maxSegmentAge) defer func() { @@ -166,17 +225,18 @@ func TestWriter_OldSegmentsAreCleanedUp(t *testing.T) { } func TestWriter_NoSegmentIsCleanedUpIfTheresOnlyOne(t *testing.T) { - dir := t.TempDir() - - maxSegmentAge := time.Second * 2 - - segmentsReclaimedNotificationsReceived := []int{} + var ( + dir = t.TempDir() + reg = prometheus.NewRegistry() + maxSegmentAge = time.Second * 2 + segmentsReclaimedNotificationsReceived = []int{} + ) writer, err := NewWriter(Config{ Dir: dir, Enabled: true, MaxSegmentAge: maxSegmentAge, - }, autil.TestAlloyLogger(t).Slog(), prometheus.NewRegistry()) + }, autil.TestAlloyLogger(t).Slog(), reg, NewWriterMetrics(reg)) require.NoError(t, err) writer.Start(maxSegmentAge) defer func() { @@ -342,13 +402,16 @@ func BenchmarkWriter_WriteEntries(b *testing.B) { } func benchWriteEntries(b *testing.B, lines, labelSetCount int) { - dir := b.TempDir() + var ( + dir = b.TempDir() + reg = prometheus.NewRegistry() + ) writer, err := NewWriter(Config{ Dir: dir, Enabled: true, MaxSegmentAge: time.Minute, - }, logging.NewSlogNop(), prometheus.NewRegistry()) + }, logging.NewSlogNop(), reg, NewWriterMetrics(reg)) require.NoError(b, err) writer.Start(time.Minute) defer func() {