Move metrics to metrics/ pacakge

This commit is contained in:
binwiederhier
2026-07-17 22:08:21 +02:00
parent ac63a2eea0
commit f8d2fcd7a6
14 changed files with 256 additions and 186 deletions
-1
View File
@@ -150,7 +150,6 @@ type Config struct {
TwilioVerifyBaseURL string
TwilioVerifyService string
TwilioCallFormat *template.Template
MetricsEnable bool
MetricsListenHTTP string
ProfileListenHTTP string
MessageDelayMin time.Duration
+15 -16
View File
@@ -16,22 +16,21 @@ import (
// Log tags
const (
tagStartup = "startup"
tagHTTP = "http"
tagPublish = "publish"
tagSubscribe = "subscribe"
tagFirebase = "firebase"
tagSMTP = "smtp" // Receive email
tagEmail = "email" // Send email
tagTwilio = "twilio"
tagMessageCache = "message_cache"
tagStripe = "stripe"
tagAccount = "account"
tagManager = "manager"
tagResetter = "resetter"
tagWebsocket = "websocket"
tagMatrix = "matrix"
tagWebPush = "webpush"
tagStartup = "startup"
tagHTTP = "http"
tagPublish = "publish"
tagSubscribe = "subscribe"
tagFirebase = "firebase"
tagSMTP = "smtp" // Receive email
tagEmail = "email" // Send email
tagTwilio = "twilio"
tagStripe = "stripe"
tagAccount = "account"
tagManager = "manager"
tagResetter = "resetter"
tagWebsocket = "websocket"
tagMatrix = "matrix"
tagWebPush = "webpush"
)
var (
+15 -20
View File
@@ -36,6 +36,7 @@ import (
"heckel.io/ntfy/v2/log"
"heckel.io/ntfy/v2/mail"
"heckel.io/ntfy/v2/message"
"heckel.io/ntfy/v2/metrics"
"heckel.io/ntfy/v2/model"
"heckel.io/ntfy/v2/payments"
"heckel.io/ntfy/v2/twilio"
@@ -411,13 +412,11 @@ func (s *Server) Run() error {
}()
}
if s.config.MetricsListenHTTP != "" {
initMetrics()
s.httpMetricsServer = &http.Server{Addr: s.config.MetricsListenHTTP, Handler: promhttp.Handler()}
go func() {
errChan <- s.httpMetricsServer.ListenAndServe()
}()
} else if s.config.EnableMetrics {
initMetrics()
s.metricsHandler = promhttp.Handler()
}
if s.config.ProfileListenHTTP != "" {
@@ -505,9 +504,7 @@ func (s *Server) handle(w http.ResponseWriter, r *http.Request) {
s.handleError(w, r, v, err)
return
}
if metricHTTPRequests != nil {
metricHTTPRequests.WithLabelValues("200", "20000", r.Method).Inc()
}
metrics.HTTPRequests.WithLabelValues("200", "20000", r.Method).Inc()
}).
Debug("HTTP request finished")
}
@@ -517,9 +514,7 @@ func (s *Server) handleError(w http.ResponseWriter, r *http.Request, v *visitor,
if !ok {
httpErr = errHTTPInternalError
}
if metricHTTPRequests != nil {
metricHTTPRequests.WithLabelValues(fmt.Sprintf("%d", httpErr.HTTPCode), fmt.Sprintf("%d", httpErr.Code), r.Method).Inc()
}
metrics.HTTPRequests.WithLabelValues(strconv.Itoa(httpErr.HTTPCode), strconv.Itoa(httpErr.Code), r.Method).Inc()
isRateLimiting := util.Contains(rateLimitingErrorCodes, httpErr.HTTPCode)
isNormalError := strings.Contains(err.Error(), "i/o timeout") || util.Contains(normalErrorCodes, httpErr.HTTPCode)
ev := logvr(v, r).Err(err)
@@ -940,27 +935,27 @@ func (s *Server) handlePublishInternal(r *http.Request, v *visitor) (*model.Mess
s.messages++
s.mu.Unlock()
if unifiedpush {
minc(metricUnifiedPushPublishedSuccess)
metrics.UnifiedPushPublishedSuccess.Inc()
}
mset(metricMessagePublishDurationMillis, time.Since(start).Milliseconds())
metrics.MessagePublishDurationMillis.Set(float64(time.Since(start).Milliseconds()))
return m, nil
}
func (s *Server) handlePublish(w http.ResponseWriter, r *http.Request, v *visitor) error {
m, err := s.handlePublishInternal(r, v)
if err != nil {
minc(metricMessagesPublishedFailure)
metrics.MessagesPublishedFailure.Inc()
return err
}
minc(metricMessagesPublishedSuccess)
metrics.MessagesPublishedSuccess.Inc()
return s.writeJSON(w, m.ForJSON())
}
func (s *Server) handlePublishMatrix(w http.ResponseWriter, r *http.Request, v *visitor) error {
_, err := s.handlePublishInternal(r, v)
if err != nil {
minc(metricMessagesPublishedFailure)
minc(metricMatrixPublishedFailure)
metrics.MessagesPublishedFailure.Inc()
metrics.MatrixPublishedFailure.Inc()
if e, ok := err.(*errHTTP); ok && e.HTTPCode == errHTTPInsufficientStorageUnifiedPush.HTTPCode {
topic, err := fromContext[*topic](r, contextTopic)
if err != nil {
@@ -976,8 +971,8 @@ func (s *Server) handlePublishMatrix(w http.ResponseWriter, r *http.Request, v *
}
return err
}
minc(metricMessagesPublishedSuccess)
minc(metricMatrixPublishedSuccess)
metrics.MessagesPublishedSuccess.Inc()
metrics.MatrixPublishedSuccess.Inc()
return writeMatrixSuccess(w)
}
@@ -1049,7 +1044,7 @@ func (s *Server) handleActionMessage(w http.ResponseWriter, r *http.Request, v *
func (s *Server) sendToFirebase(v *visitor, m *model.Message) {
logvm(v, m).Tag(tagFirebase).Debug("Publishing to Firebase")
if err := s.firebaseClient.Send(v, m); err != nil {
minc(metricFirebasePublishedFailure)
metrics.FirebasePublishedFailure.Inc()
if errors.Is(err, errFirebaseTemporarilyBanned) {
logvm(v, m).Tag(tagFirebase).Err(err).Debug("Unable to publish to Firebase: %v", err.Error())
} else {
@@ -1057,17 +1052,17 @@ func (s *Server) sendToFirebase(v *visitor, m *model.Message) {
}
return
}
minc(metricFirebasePublishedSuccess)
metrics.FirebasePublishedSuccess.Inc()
}
func (s *Server) sendEmail(v *visitor, m *model.Message, email string) {
logvm(v, m).Tag(tagEmail).Field("email", email).Info("Sending email to %s", email)
if err := s.mailer.SendNotification(email, m, v.ip.String()); err != nil {
logvm(v, m).Tag(tagEmail).Field("email", email).Err(err).Warn("Unable to send email to %s: %v", email, err.Error())
minc(metricEmailsPublishedFailure)
metrics.EmailsPublishedFailure.Inc()
return
}
minc(metricEmailsPublishedSuccess)
metrics.EmailsPublishedSuccess.Inc()
}
func (s *Server) forwardPollRequest(v *visitor, m *model.Message) {
+3 -2
View File
@@ -423,8 +423,9 @@
# doing, and/or secure access to the endpoint in your reverse proxy.
#
# - enable-metrics enables the /metrics endpoint for the default ntfy server (i.e. HTTP, HTTPS and/or Unix socket)
# - metrics-listen-http exposes the metrics endpoint via a dedicated [IP]:port. If set, this option implicitly
# enables metrics as well, e.g. "10.0.1.1:9090" or ":9090"
# - metrics-listen-http moves the metrics endpoint to a dedicated [IP]:port, e.g. "10.0.1.1:9090" or ":9090".
# It implicitly enables metrics. If set, the metrics are served only on that dedicated port, and the default
# ntfy server does not serve /metrics, even if enable-metrics is also set.
#
# enable-metrics: false
# metrics-listen-http:
+7 -6
View File
@@ -2,6 +2,7 @@ package server
import (
"heckel.io/ntfy/v2/log"
"heckel.io/ntfy/v2/metrics"
"heckel.io/ntfy/v2/util"
)
@@ -93,13 +94,13 @@ func (s *Server) execManager() {
"emails_sent_failure": sentMailFailure,
}).
Info("Server stats")
mset(metricMessagesCached, messagesCached)
mset(metricVisitors, visitorsCount)
mset(metricUsers, usersCount)
mset(metricSubscribers, subscribers)
mset(metricTopics, topicsCount)
metrics.MessagesCached.Set(float64(messagesCached))
metrics.Visitors.Set(float64(visitorsCount))
metrics.Users.Set(float64(usersCount))
metrics.Subscribers.Set(float64(subscribers))
metrics.Topics.Set(float64(topicsCount))
if s.attachment != nil {
mset(metricAttachmentsTotalSize, s.attachment.Size())
metrics.AttachmentsTotalSize.Set(float64(s.attachment.Size()))
}
}
-132
View File
@@ -1,132 +0,0 @@
package server
import (
"github.com/prometheus/client_golang/prometheus"
)
var (
metricMessagesPublishedSuccess prometheus.Counter
metricMessagesPublishedFailure prometheus.Counter
metricMessagesCached prometheus.Gauge
metricMessagePublishDurationMillis prometheus.Gauge
metricFirebasePublishedSuccess prometheus.Counter
metricFirebasePublishedFailure prometheus.Counter
metricEmailsPublishedSuccess prometheus.Counter
metricEmailsPublishedFailure prometheus.Counter
metricEmailsReceivedSuccess prometheus.Counter
metricEmailsReceivedFailure prometheus.Counter
metricCallsMadeSuccess prometheus.Counter
metricCallsMadeFailure prometheus.Counter
metricUnifiedPushPublishedSuccess prometheus.Counter
metricMatrixPublishedSuccess prometheus.Counter
metricMatrixPublishedFailure prometheus.Counter
metricAttachmentsTotalSize prometheus.Gauge
metricVisitors prometheus.Gauge
metricSubscribers prometheus.Gauge
metricTopics prometheus.Gauge
metricUsers prometheus.Gauge
metricHTTPRequests *prometheus.CounterVec
)
func initMetrics() {
metricMessagesPublishedSuccess = prometheus.NewCounter(prometheus.CounterOpts{
Name: "ntfy_messages_published_success",
})
metricMessagesPublishedFailure = prometheus.NewCounter(prometheus.CounterOpts{
Name: "ntfy_messages_published_failure",
})
metricMessagesCached = prometheus.NewGauge(prometheus.GaugeOpts{
Name: "ntfy_messages_cached_total",
})
metricMessagePublishDurationMillis = prometheus.NewGauge(prometheus.GaugeOpts{
Name: "ntfy_message_publish_duration_ms",
})
metricFirebasePublishedSuccess = prometheus.NewCounter(prometheus.CounterOpts{
Name: "ntfy_firebase_published_success",
})
metricFirebasePublishedFailure = prometheus.NewCounter(prometheus.CounterOpts{
Name: "ntfy_firebase_published_failure",
})
metricEmailsPublishedSuccess = prometheus.NewCounter(prometheus.CounterOpts{
Name: "ntfy_emails_sent_success",
})
metricEmailsPublishedFailure = prometheus.NewCounter(prometheus.CounterOpts{
Name: "ntfy_emails_sent_failure",
})
metricEmailsReceivedSuccess = prometheus.NewCounter(prometheus.CounterOpts{
Name: "ntfy_emails_received_success",
})
metricEmailsReceivedFailure = prometheus.NewCounter(prometheus.CounterOpts{
Name: "ntfy_emails_received_failure",
})
metricCallsMadeSuccess = prometheus.NewCounter(prometheus.CounterOpts{
Name: "ntfy_calls_made_success",
})
metricCallsMadeFailure = prometheus.NewCounter(prometheus.CounterOpts{
Name: "ntfy_calls_made_failure",
})
metricUnifiedPushPublishedSuccess = prometheus.NewCounter(prometheus.CounterOpts{
Name: "ntfy_unifiedpush_published_success",
})
metricMatrixPublishedSuccess = prometheus.NewCounter(prometheus.CounterOpts{
Name: "ntfy_matrix_published_success",
})
metricMatrixPublishedFailure = prometheus.NewCounter(prometheus.CounterOpts{
Name: "ntfy_matrix_published_failure",
})
metricAttachmentsTotalSize = prometheus.NewGauge(prometheus.GaugeOpts{
Name: "ntfy_attachments_total_size",
})
metricVisitors = prometheus.NewGauge(prometheus.GaugeOpts{
Name: "ntfy_visitors_total",
})
metricUsers = prometheus.NewGauge(prometheus.GaugeOpts{
Name: "ntfy_users_total",
})
metricSubscribers = prometheus.NewGauge(prometheus.GaugeOpts{
Name: "ntfy_subscribers_total",
})
metricTopics = prometheus.NewGauge(prometheus.GaugeOpts{
Name: "ntfy_topics_total",
})
metricHTTPRequests = prometheus.NewCounterVec(prometheus.CounterOpts{
Name: "ntfy_http_requests_total",
}, []string{"http_code", "ntfy_code", "http_method"})
prometheus.MustRegister(
metricMessagesPublishedSuccess,
metricMessagesPublishedFailure,
metricMessagesCached,
metricMessagePublishDurationMillis,
metricFirebasePublishedSuccess,
metricFirebasePublishedFailure,
metricEmailsPublishedSuccess,
metricEmailsPublishedFailure,
metricEmailsReceivedSuccess,
metricEmailsReceivedFailure,
metricCallsMadeSuccess,
metricCallsMadeFailure,
metricUnifiedPushPublishedSuccess,
metricMatrixPublishedSuccess,
metricMatrixPublishedFailure,
metricAttachmentsTotalSize,
metricVisitors,
metricUsers,
metricSubscribers,
metricTopics,
metricHTTPRequests,
)
}
// minc increments a prometheus.Counter if it is non-nil
func minc(counter prometheus.Counter) {
if counter != nil {
counter.Inc()
}
}
// mset sets a prometheus.Gauge if it is non-nil
func mset[T int | int64 | float64](gauge prometheus.Gauge, value T) {
if gauge != nil {
gauge.Set(float64(value))
}
}
+38
View File
@@ -22,6 +22,7 @@ import (
"testing"
"time"
"github.com/prometheus/client_golang/prometheus/promhttp"
"github.com/stretchr/testify/require"
"golang.org/x/crypto/bcrypt"
dbtest "heckel.io/ntfy/v2/db/test"
@@ -323,6 +324,40 @@ func TestServer_WebEnabled(t *testing.T) {
require.Equal(t, 200, rr.Code)
})
}
// TestServer_MetricsEnabled ensures that the /metrics endpoint serves the registered ntfy metrics
// once the metrics handler is set (as Serve does when enable-metrics is configured).
func TestServer_MetricsEnabled(t *testing.T) {
forEachBackend(t, func(t *testing.T, databaseURL string) {
s := newTestServer(t, newTestConfig(t, databaseURL))
s.metricsHandler = promhttp.Handler() // Serve sets this when enable-metrics is configured
// Count at least one request first: Prometheus only reports a CounterVec such as
// ntfy_http_requests_total once it has children
request(t, s, "GET", "/v1/health", "", nil)
rr := request(t, s, "GET", "/metrics", "", nil)
require.Equal(t, 200, rr.Code)
require.Contains(t, rr.Body.String(), "ntfy_messages_published_success")
require.Contains(t, rr.Body.String(), "ntfy_http_requests_total")
})
}
// TestServer_MetricsDisabled ensures that the ntfy metrics are not exposed when the metrics handler
// is unset (the default). The collectors are always registered with the Prometheus registry, so a
// nil metrics handler is the only thing keeping them off the wire.
func TestServer_MetricsDisabled(t *testing.T) {
forEachBackend(t, func(t *testing.T, databaseURL string) {
conf := newTestConfig(t, databaseURL)
conf.WebRoot = "" // Disable the web app, so its catch-all does not mask the /metrics route
s := newTestServer(t, conf)
rr := request(t, s, "GET", "/metrics", "", nil)
require.Equal(t, 404, rr.Code)
require.NotContains(t, rr.Body.String(), "ntfy_messages_published_success")
})
}
func TestServer_PublishLargeMessage(t *testing.T) {
forEachBackend(t, func(t *testing.T, databaseURL string) {
c := newTestConfig(t, databaseURL)
@@ -1684,6 +1719,7 @@ func TestServer_PublishEmailVerify_BoolValueUsesPrimary(t *testing.T) {
"Authorization": util.BasicAuth("phil", "phil"),
})
require.Equal(t, 200, response.Code)
waitFor(t, func() bool { return mailer.LastTo() != "" }) // E-Mail publishing happens in a Go routine
require.Equal(t, "zzz@example.com", mailer.LastTo())
})
}
@@ -1710,6 +1746,7 @@ func TestServer_PublishEmailVerify_BoolValueNoVerifyUsesPrimary(t *testing.T) {
"Authorization": util.BasicAuth("phil", "phil"),
})
require.Equal(t, 200, response.Code)
waitFor(t, func() bool { return mailer.LastTo() != "" }) // E-Mail publishing happens in a Go routine
require.Equal(t, "zzz@example.com", mailer.LastTo())
})
}
@@ -1751,6 +1788,7 @@ func TestServer_PublishEmailVerify_BoolValueProvisionedUsesPrimary(t *testing.T)
"Authorization": util.BasicAuth("prov", "provpass"),
})
require.Equal(t, 200, response.Code)
waitFor(t, func() bool { return mailer.LastTo() != "" }) // E-Mail publishing happens in a Go routine
require.Equal(t, "zzz@example.com", mailer.LastTo())
})
}
+3 -2
View File
@@ -1,6 +1,7 @@
package server
import (
"heckel.io/ntfy/v2/metrics"
"heckel.io/ntfy/v2/model"
"heckel.io/ntfy/v2/twilio"
"heckel.io/ntfy/v2/user"
@@ -46,8 +47,8 @@ func (s *Server) callPhone(v *visitor, m *model.Message, to string) {
})
if err != nil {
logvm(v, m).Tag(tagTwilio).Field("twilio_to", to).Err(err).Warn("Unable to call phone %s: %v", to, err.Error())
minc(metricCallsMadeFailure)
metrics.CallsMadeFailure.Inc()
return
}
minc(metricCallsMadeSuccess)
metrics.CallsMadeSuccess.Inc()
}
+3 -2
View File
@@ -19,6 +19,7 @@ import (
"github.com/emersion/go-smtp"
"github.com/microcosm-cc/bluemonday"
"heckel.io/ntfy/v2/metrics"
"heckel.io/ntfy/v2/model"
)
@@ -180,7 +181,7 @@ func (s *smtpSession) Data(r io.Reader) error {
s.backend.mu.Lock()
s.backend.success++
s.backend.mu.Unlock()
minc(metricEmailsReceivedSuccess)
metrics.EmailsReceivedSuccess.Inc()
return nil
})
}
@@ -238,7 +239,7 @@ func (s *smtpSession) withFailCount(fn func() error) error {
// We do not want to spam the log with WARN messages.
logem(s.conn).Err(err).Debug("Incoming mail error")
s.backend.failure++
minc(metricEmailsReceivedFailure)
metrics.EmailsReceivedFailure.Inc()
}
return err
}