metrics.go 3.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102
  1. package observability
  2. import (
  3. "github.com/prometheus/client_golang/prometheus"
  4. )
  5. // NewRegistry returns a fresh Prometheus registry. Each service gets
  6. // its own so the metric labels are scoped correctly.
  7. func NewRegistry(serviceName string) (*prometheus.Registry, *IngestdMetrics) {
  8. reg := prometheus.NewRegistry()
  9. reg.MustRegister(
  10. prometheus.NewGoCollector(),
  11. prometheus.NewProcessCollector(prometheus.ProcessCollectorOpts{}),
  12. )
  13. return reg, NewIngestdMetrics(reg, serviceName)
  14. }
  15. // IngestdMetrics groups the counters/histograms declared in SPEC §22
  16. // for the ingest tier. Other tiers get their own metric groups.
  17. type IngestdMetrics struct {
  18. AlertsReceived *prometheus.CounterVec // result=accepted|invalid|rate_limited|payload_too_large|quarantined|circuit_open
  19. PayloadBytes prometheus.Histogram
  20. RateLimitHits *prometheus.CounterVec // scope=source|company
  21. Quarantines *prometheus.CounterVec
  22. CBState *prometheus.GaugeVec
  23. PublishLatency prometheus.Histogram
  24. // MQTTMessages is the M4 per-message counter; labels mirror
  25. // the same result taxonomy as AlertsReceived (accepted,
  26. // deduped, bad_topic, bad_signature, unknown_source,
  27. // rate_limited_source, rate_limited_company, invalid,
  28. // invalid_json, broker_unavailable) plus a "received" label
  29. // for every message that survived parseIncomingTopic.
  30. MQTTMessages *prometheus.CounterVec
  31. }
  32. // NewIngestdMetrics registers and returns the ingestd metrics.
  33. func NewIngestdMetrics(reg prometheus.Registerer, serviceName string) *IngestdMetrics {
  34. m := &IngestdMetrics{
  35. AlertsReceived: prometheus.NewCounterVec(prometheus.CounterOpts{
  36. Namespace: "ba",
  37. Subsystem: "ingestd",
  38. Name: "alerts_received_total",
  39. Help: "Number of inbound alerts by result.",
  40. ConstLabels: prometheus.Labels{"service": serviceName},
  41. }, []string{"result"}),
  42. PayloadBytes: prometheus.NewHistogram(prometheus.HistogramOpts{
  43. Namespace: "ba",
  44. Subsystem: "ingestd",
  45. Name: "payload_bytes",
  46. Help: "Accepted alert payload size in bytes.",
  47. Buckets: prometheus.ExponentialBuckets(64, 4, 8), // 64..1MB
  48. ConstLabels: prometheus.Labels{"service": serviceName},
  49. }),
  50. RateLimitHits: prometheus.NewCounterVec(prometheus.CounterOpts{
  51. Namespace: "ba",
  52. Subsystem: "ingestd",
  53. Name: "rate_limited_total",
  54. Help: "Rate-limit rejections by scope.",
  55. ConstLabels: prometheus.Labels{"service": serviceName},
  56. }, []string{"scope"}),
  57. Quarantines: prometheus.NewCounterVec(prometheus.CounterOpts{
  58. Namespace: "ba",
  59. Subsystem: "ingestd",
  60. Name: "source_quarantined_total",
  61. Help: "Source quarantines triggered.",
  62. ConstLabels: prometheus.Labels{"service": serviceName},
  63. }, []string{"source_id", "company_id"}),
  64. CBState: prometheus.NewGaugeVec(prometheus.GaugeOpts{
  65. Namespace: "ba",
  66. Subsystem: "ingestd",
  67. Name: "circuit_breaker_state",
  68. Help: "0=closed, 1=half_open, 2=open.",
  69. ConstLabels: prometheus.Labels{"service": serviceName},
  70. }, []string{"component"}),
  71. PublishLatency: prometheus.NewHistogram(prometheus.HistogramOpts{
  72. Namespace: "ba",
  73. Subsystem: "ingestd",
  74. Name: "publish_latency_seconds",
  75. Help: "Time to publish an accepted alert to NATS.",
  76. Buckets: prometheus.DefBuckets,
  77. ConstLabels: prometheus.Labels{"service": serviceName},
  78. }),
  79. MQTTMessages: prometheus.NewCounterVec(prometheus.CounterOpts{
  80. Namespace: "ba",
  81. Subsystem: "ingestd",
  82. Name: "mqtt_messages_total",
  83. Help: "Inbound MQTT messages by result (M4).",
  84. ConstLabels: prometheus.Labels{"service": serviceName},
  85. }, []string{"result"}),
  86. }
  87. reg.MustRegister(
  88. m.AlertsReceived,
  89. m.PayloadBytes,
  90. m.RateLimitHits,
  91. m.Quarantines,
  92. m.CBState,
  93. m.PublishLatency,
  94. m.MQTTMessages,
  95. )
  96. m.AlertsReceived.WithLabelValues("accepted")
  97. return m
  98. }