Skip to content

Commit 1908ea8

Browse files
committed
lint fixes
1 parent 231a181 commit 1908ea8

6 files changed

Lines changed: 365 additions & 299 deletions

File tree

examples/dlq_metrics/main.go

Lines changed: 30 additions & 41 deletions
Original file line numberDiff line numberDiff line change
@@ -27,74 +27,63 @@ var (
2727
ErrProcessingFailed = errors.New("message processing failed")
2828
)
2929

30-
func main() {
31-
ctx, cancel := context.WithCancel(context.Background())
32-
defer cancel()
33-
34-
logger := protoflow.NewSlogServiceLogger(slog.New(slog.NewJSONHandler(os.Stdout, nil)))
35-
36-
// Create DLQ metrics collector using protoflow's built-in type
37-
dlqMetrics := protoflow.NewDLQMetrics(prometheus.DefaultRegisterer)
38-
dlqMetrics.Register()
39-
40-
// Start Prometheus metrics server
41-
go func() {
42-
http.Handle("/metrics", promhttp.Handler())
43-
logger.Info("📊 Prometheus metrics available at http://localhost:2112/metrics", nil)
44-
if err := http.ListenAndServe(":2112", nil); err != nil {
45-
logger.Error("metrics server error", err, nil)
46-
}
47-
}()
48-
49-
cfg := &protoflow.Config{
50-
PubSubSystem: "channel",
51-
PoisonQueue: "dlq.poison",
30+
// setupMetricsServer starts the Prometheus metrics HTTP server.
31+
func setupMetricsServer(logger protoflow.ServiceLogger) {
32+
http.Handle("/metrics", promhttp.Handler())
33+
logger.Info("📊 Prometheus metrics available at http://localhost:2112/metrics", nil)
34+
if err := http.ListenAndServe(":2112", nil); err != nil {
35+
logger.Error("metrics server error", err, nil)
5236
}
37+
}
5338

54-
// Build middleware chain with DLQ tracking
55-
middlewares := []protoflow.MiddlewareRegistration{
39+
// buildMiddlewares creates the middleware chain with DLQ tracking.
40+
func buildMiddlewares(dlqMetrics *protoflow.DLQMetrics) []protoflow.MiddlewareRegistration {
41+
return []protoflow.MiddlewareRegistration{
5642
protoflow.CorrelationIDMiddleware(),
57-
dlqTrackingMiddleware(dlqMetrics), // Track messages going to DLQ
43+
dlqTrackingMiddleware(dlqMetrics),
5844
protoflow.RetryMiddleware(protoflow.RetryMiddlewareConfig{
5945
MaxRetries: 2,
6046
InitialInterval: 100 * time.Millisecond,
6147
}),
6248
protoflow.PoisonQueueMiddleware(nil),
6349
protoflow.RecovererMiddleware(),
6450
}
51+
}
6552

53+
func main() {
54+
ctx, cancel := context.WithCancel(context.Background())
55+
defer cancel()
56+
57+
logger := protoflow.NewSlogServiceLogger(slog.New(slog.NewJSONHandler(os.Stdout, nil)))
58+
59+
dlqMetrics := protoflow.NewDLQMetrics(prometheus.DefaultRegisterer)
60+
if err := dlqMetrics.Register(); err != nil {
61+
logger.Error("failed to register DLQ metrics", err, nil)
62+
}
63+
64+
go setupMetricsServer(logger)
65+
66+
cfg := &protoflow.Config{PubSubSystem: "channel", PoisonQueue: "dlq.poison"}
6667
svc := protoflow.NewService(cfg, logger, ctx, protoflow.ServiceDependencies{
6768
DisableDefaultMiddlewares: true,
68-
Middlewares: middlewares,
69+
Middlewares: buildMiddlewares(dlqMetrics),
6970
})
7071

71-
// Register a handler that occasionally fails
72-
err := protoflow.RegisterMessageHandler(svc, protoflow.MessageHandlerRegistration{
72+
if err := protoflow.RegisterMessageHandler(svc, protoflow.MessageHandlerRegistration{
7373
Name: "flaky-processor",
7474
ConsumeQueue: "dlq.incoming",
7575
Handler: createFlakyHandler(logger),
76-
})
77-
if err != nil {
76+
}); err != nil {
7877
panic(err)
7978
}
8079

81-
// Publish test messages
8280
go publishTestMessages(ctx, svc, logger)
83-
84-
// Periodically print DLQ metrics
8581
go printMetricsPeriodically(ctx, dlqMetrics, logger)
86-
87-
// Run for demonstration
88-
go func() {
89-
time.Sleep(15 * time.Second)
90-
cancel()
91-
}()
82+
go func() { time.Sleep(15 * time.Second); cancel() }()
9283

9384
if err := svc.Start(ctx); err != nil && !errors.Is(err, context.Canceled) {
9485
logger.Error("service stopped with error", err, nil)
9586
}
96-
97-
// Print final metrics
9887
printSnapshot(dlqMetrics, logger)
9988
}
10089

internal/runtime/dlq_metrics.go

Lines changed: 49 additions & 61 deletions
Original file line numberDiff line numberDiff line change
@@ -47,74 +47,62 @@ type DLQMetricsSnapshot struct {
4747
CollectedAt time.Time `json:"collected_at"`
4848
}
4949

50+
// newDLQCounterVec creates a new counter vec with standard protoflow/dlq namespace.
51+
func newDLQCounterVec(name, help string, labels []string) *prometheus.CounterVec {
52+
return prometheus.NewCounterVec(
53+
prometheus.CounterOpts{
54+
Namespace: "protoflow",
55+
Subsystem: "dlq",
56+
Name: name,
57+
Help: help,
58+
},
59+
labels,
60+
)
61+
}
62+
63+
// newDLQGaugeVec creates a new gauge vec with standard protoflow/dlq namespace.
64+
func newDLQGaugeVec(name, help string, labels []string) *prometheus.GaugeVec {
65+
return prometheus.NewGaugeVec(
66+
prometheus.GaugeOpts{
67+
Namespace: "protoflow",
68+
Subsystem: "dlq",
69+
Name: name,
70+
Help: help,
71+
},
72+
labels,
73+
)
74+
}
75+
76+
// newDLQHistogramVec creates a new histogram vec with standard protoflow/dlq namespace.
77+
func newDLQHistogramVec(name, help string, buckets []float64, labels []string) *prometheus.HistogramVec {
78+
return prometheus.NewHistogramVec(
79+
prometheus.HistogramOpts{
80+
Namespace: "protoflow",
81+
Subsystem: "dlq",
82+
Name: name,
83+
Help: help,
84+
Buckets: buckets,
85+
},
86+
labels,
87+
)
88+
}
89+
5090
// NewDLQMetrics creates a new DLQ metrics collector.
5191
func NewDLQMetrics(registerer prometheus.Registerer) *DLQMetrics {
5292
if registerer == nil {
5393
registerer = prometheus.DefaultRegisterer
5494
}
5595

56-
m := &DLQMetrics{
57-
topicCounts: make(map[string]*DLQTopicMetrics),
58-
registerer: registerer,
59-
messagesTotal: prometheus.NewCounterVec(
60-
prometheus.CounterOpts{
61-
Namespace: "protoflow",
62-
Subsystem: "dlq",
63-
Name: "messages_total",
64-
Help: "Total number of messages sent to the dead letter queue",
65-
},
66-
[]string{"topic", "handler"},
67-
),
68-
messagesCurrent: prometheus.NewGaugeVec(
69-
prometheus.GaugeOpts{
70-
Namespace: "protoflow",
71-
Subsystem: "dlq",
72-
Name: "messages_current",
73-
Help: "Current number of messages in the dead letter queue",
74-
},
75-
[]string{"topic"},
76-
),
77-
replayedTotal: prometheus.NewCounterVec(
78-
prometheus.CounterOpts{
79-
Namespace: "protoflow",
80-
Subsystem: "dlq",
81-
Name: "replayed_total",
82-
Help: "Total number of messages replayed from the dead letter queue",
83-
},
84-
[]string{"topic"},
85-
),
86-
purgedTotal: prometheus.NewCounterVec(
87-
prometheus.CounterOpts{
88-
Namespace: "protoflow",
89-
Subsystem: "dlq",
90-
Name: "purged_total",
91-
Help: "Total number of messages purged from the dead letter queue",
92-
},
93-
[]string{"topic"},
94-
),
95-
ageSecondsHist: prometheus.NewHistogramVec(
96-
prometheus.HistogramOpts{
97-
Namespace: "protoflow",
98-
Subsystem: "dlq",
99-
Name: "message_age_seconds",
100-
Help: "Age of messages when moved to DLQ (time since first attempt)",
101-
Buckets: []float64{1, 5, 10, 30, 60, 300, 600, 1800, 3600},
102-
},
103-
[]string{"topic"},
104-
),
105-
retryCountHist: prometheus.NewHistogramVec(
106-
prometheus.HistogramOpts{
107-
Namespace: "protoflow",
108-
Subsystem: "dlq",
109-
Name: "retry_count",
110-
Help: "Number of retries before message was moved to DLQ",
111-
Buckets: []float64{1, 2, 3, 5, 10, 20},
112-
},
113-
[]string{"topic"},
114-
),
96+
return &DLQMetrics{
97+
topicCounts: make(map[string]*DLQTopicMetrics),
98+
registerer: registerer,
99+
messagesTotal: newDLQCounterVec("messages_total", "Total number of messages sent to the dead letter queue", []string{"topic", "handler"}),
100+
messagesCurrent: newDLQGaugeVec("messages_current", "Current number of messages in the dead letter queue", []string{"topic"}),
101+
replayedTotal: newDLQCounterVec("replayed_total", "Total number of messages replayed from the dead letter queue", []string{"topic"}),
102+
purgedTotal: newDLQCounterVec("purged_total", "Total number of messages purged from the dead letter queue", []string{"topic"}),
103+
ageSecondsHist: newDLQHistogramVec("message_age_seconds", "Age of messages when moved to DLQ (time since first attempt)", []float64{1, 5, 10, 30, 60, 300, 600, 1800, 3600}, []string{"topic"}),
104+
retryCountHist: newDLQHistogramVec("retry_count", "Number of retries before message was moved to DLQ", []float64{1, 2, 3, 5, 10, 20}, []string{"topic"}),
115105
}
116-
117-
return m
118106
}
119107

120108
// Register registers the Prometheus collectors. Safe to call multiple times.

internal/runtime/transport/io.go

Lines changed: 55 additions & 40 deletions
Original file line numberDiff line numberDiff line change
@@ -99,6 +99,58 @@ type ioSubscriber struct {
9999
logger watermill.LoggerAdapter
100100
}
101101

102+
// handleEOF handles end-of-file condition by waiting and seeking back to last position.
103+
// Returns true if the goroutine should continue, false if it should exit.
104+
func (s *ioSubscriber) handleEOF(f *os.File, reader *bufio.Reader, lastPos *int64) bool {
105+
currentPos, _ := f.Seek(0, io.SeekCurrent)
106+
currentPos -= int64(reader.Buffered())
107+
108+
if currentPos > *lastPos {
109+
*lastPos = currentPos
110+
}
111+
112+
time.Sleep(50 * time.Millisecond)
113+
114+
if _, err := f.Seek(*lastPos, io.SeekStart); err != nil {
115+
s.logger.Error("Failed to seek file", err, nil)
116+
return false
117+
}
118+
reader.Reset(f)
119+
return true
120+
}
121+
122+
// processMessage parses and sends a message if it matches the topic.
123+
// Returns true if processing should continue.
124+
func (s *ioSubscriber) processMessage(ctx context.Context, out chan<- *message.Message, line []byte, topic string) bool {
125+
var sm storedMessage
126+
if err := json.Unmarshal(line, &sm); err != nil {
127+
s.logger.Error("Failed to unmarshal message", err, nil)
128+
return true // continue processing
129+
}
130+
131+
if sm.Topic != topic {
132+
return true // continue processing
133+
}
134+
135+
msg := message.NewMessage(sm.UUID, sm.Payload)
136+
msg.Metadata = sm.Metadata
137+
138+
select {
139+
case out <- msg:
140+
select {
141+
case <-msg.Acked():
142+
// good
143+
case <-msg.Nacked():
144+
s.logger.Debug("Message nacked", watermill.LogFields{"uuid": msg.UUID})
145+
case <-ctx.Done():
146+
return false
147+
}
148+
case <-ctx.Done():
149+
return false
150+
}
151+
return true
152+
}
153+
102154
func (s *ioSubscriber) Subscribe(ctx context.Context, topic string) (<-chan *message.Message, error) {
103155
out := make(chan *message.Message)
104156

@@ -112,9 +164,7 @@ func (s *ioSubscriber) Subscribe(ctx context.Context, topic string) (<-chan *mes
112164
}
113165
defer f.Close()
114166

115-
// Track our position so we can continue reading new content
116167
var lastPos int64
117-
118168
reader := bufio.NewReader(f)
119169

120170
for {
@@ -125,21 +175,9 @@ func (s *ioSubscriber) Subscribe(ctx context.Context, topic string) (<-chan *mes
125175
line, err := reader.ReadBytes('\n')
126176
if err != nil {
127177
if err == io.EOF {
128-
// Get current position and check if file has grown
129-
currentPos, _ := f.Seek(0, io.SeekCurrent)
130-
// Subtract any buffered but unread data
131-
currentPos -= int64(reader.Buffered())
132-
133-
if currentPos > lastPos {
134-
lastPos = currentPos
178+
if !s.handleEOF(f, reader, &lastPos) {
179+
return
135180
}
136-
137-
// Wait and then check for new content
138-
time.Sleep(50 * time.Millisecond)
139-
140-
// Seek to our last known position and reset reader
141-
f.Seek(lastPos, io.SeekStart)
142-
reader.Reset(f)
143181
continue
144182
}
145183
s.logger.Error("Failed to read file", err, nil)
@@ -150,30 +188,7 @@ func (s *ioSubscriber) Subscribe(ctx context.Context, topic string) (<-chan *mes
150188
currentPos, _ := f.Seek(0, io.SeekCurrent)
151189
lastPos = currentPos - int64(reader.Buffered())
152190

153-
var sm storedMessage
154-
if err := json.Unmarshal(line, &sm); err != nil {
155-
s.logger.Error("Failed to unmarshal message", err, nil)
156-
continue
157-
}
158-
159-
if sm.Topic != topic {
160-
continue
161-
}
162-
163-
msg := message.NewMessage(sm.UUID, sm.Payload)
164-
msg.Metadata = sm.Metadata
165-
166-
select {
167-
case out <- msg:
168-
select {
169-
case <-msg.Acked():
170-
// good
171-
case <-msg.Nacked():
172-
s.logger.Debug("Message nacked", watermill.LogFields{"uuid": msg.UUID})
173-
case <-ctx.Done():
174-
return
175-
}
176-
case <-ctx.Done():
191+
if !s.processMessage(ctx, out, line, topic) {
177192
return
178193
}
179194
}

0 commit comments

Comments
 (0)