feat(go): expose nats bridge pending metrics

This commit is contained in:
lingniu
2026-07-02 19:27:36 +08:00
parent 1f4a23ff69
commit df96ec8879
3 changed files with 54 additions and 2 deletions

View File

@@ -7,6 +7,7 @@ import (
"strings"
"testing"
"github.com/nats-io/nats.go"
"github.com/segmentio/kafka-go"
"lingniu-vehicle-ingest/go/vehicle-gateway/internal/envelope"
@@ -90,6 +91,28 @@ func TestBridgeBatchRecordsMetrics(t *testing.T) {
}
}
func TestRecordNATSConsumerInfoMetrics(t *testing.T) {
registry := metrics.NewRegistry()
cfg := config{NATSStream: "VEHICLE_INGEST", NATSDurable: "vehicle-kafka-bridge"}
recordNATSConsumerInfoMetrics(registry, cfg, &nats.ConsumerInfo{
NumPending: 17,
NumAckPending: 3,
NumWaiting: 2,
})
text := registry.Render()
for _, want := range []string{
`vehicle_bridge_nats_consumer_pending{consumer="vehicle-kafka-bridge",stream="VEHICLE_INGEST"} 17`,
`vehicle_bridge_nats_consumer_ack_pending{consumer="vehicle-kafka-bridge",stream="VEHICLE_INGEST"} 3`,
`vehicle_bridge_nats_consumer_waiting{consumer="vehicle-kafka-bridge",stream="VEHICLE_INGEST"} 2`,
} {
if !strings.Contains(text, want) {
t.Fatalf("metrics missing %s:\n%s", want, text)
}
}
}
type recordingBridgeWriter struct {
messages []kafka.Message
err error