Files
otelsetup/metrics_test.go
T
argoyleandClaude Opus 5 bbf27f9089
otelsetup / test (push) Skipped
otelsetup / vulnerabilities (push) Skipped
pre-commit / pre-commit (push) Skipped
otelsetup / vulnerabilities (pull_request) Successful in 53s
otelsetup / test (pull_request) Successful in 1m9s
pre-commit / pre-commit (pull_request) Successful in 2m48s
feat(metrics): record eventsourced pg outbox publishes and retries
NewEventsourcedMetrics now handles pg.OutboxPublish (eventsourced.outbox.publish.duration, by event.type and success) and pg.OutboxRetry (eventsourced.outbox.retries, by event.type and permanent), so a service using the pg/v2 outbox can alert on failing or abandoned publishes. OutboxBatch and OutboxCleanup stay ignored.

otelsetup now imports codeberg.org/eventsourced/pg/v2 (v2.1.1). Every current consumer already depends on it; a release raises their minimum to v2.1.1.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_013DJouG8ZZZxKtodF9kDvzj
2026-09-17 13:47:30 +02:00

81 lines
2.4 KiB
Go

package otelsetup
import (
"context"
"sort"
"testing"
"time"
"codeberg.org/eventsourced/eventsourced"
"codeberg.org/eventsourced/pg/v2"
"go.opentelemetry.io/otel"
sdkmetric "go.opentelemetry.io/otel/sdk/metric"
"go.opentelemetry.io/otel/sdk/metric/metricdata"
)
func TestNewEventsourcedMetrics_RecordsContract(t *testing.T) {
reader := sdkmetric.NewManualReader()
otel.SetMeterProvider(sdkmetric.NewMeterProvider(sdkmetric.WithReader(reader)))
r, err := NewEventsourcedMetrics()
if err != nil {
t.Fatalf("NewEventsourcedMetrics returned error: %v", err)
}
if r == nil {
t.Fatal("NewEventsourcedMetrics returned nil recorder")
}
// Recording every known metric type (and an unknown one) must not panic
// and must emit the expected instruments.
for _, m := range []eventsourced.Metric{
eventsourced.CommandDuration{CommandType: "AddEntry", Duration: time.Millisecond, Success: true},
eventsourced.EventStored{AggregateType: "Entry", EventType: "EntryAdded", Duration: time.Millisecond},
eventsourced.EventsLoaded{AggregateType: "Entry", EventCount: 3, Duration: time.Millisecond},
eventsourced.SnapshotStored{AggregateType: "Entry", Duration: time.Millisecond, Success: true},
eventsourced.SnapshotLoaded{AggregateType: "Entry", Found: false, Duration: time.Millisecond},
eventsourced.IdempotencyCheck{AggregateType: "Entry", Hit: true},
pg.OutboxPublish{EventType: "EntryAdded", Success: false, Duration: time.Millisecond},
pg.OutboxRetry{EventType: "EntryAdded", RetryCount: 1},
unknownMetric{},
} {
r.Record(context.Background(), m)
}
var rm metricdata.ResourceMetrics
if err := reader.Collect(context.Background(), &rm); err != nil {
t.Fatalf("collect: %v", err)
}
got := map[string]bool{}
for _, sm := range rm.ScopeMetrics {
for _, md := range sm.Metrics {
got[md.Name] = true
}
}
want := []string{
"eventsourced.command.duration",
"eventsourced.event.store.duration",
"eventsourced.events.loaded",
"eventsourced.event.load.duration",
"eventsourced.snapshot.store.duration",
"eventsourced.snapshot.load.duration",
"eventsourced.idempotency.checks",
"eventsourced.outbox.publish.duration",
"eventsourced.outbox.retries",
}
var missing []string
for _, w := range want {
if !got[w] {
missing = append(missing, w)
}
}
if len(missing) > 0 {
sort.Strings(missing)
t.Errorf("missing expected metrics: %v", missing)
}
}
type unknownMetric struct{}
func (unknownMetric) IsMetric() {}