fix: 合并 Prometheus 指标接线
This commit is contained in:
+17
-6
@@ -71,6 +71,10 @@ func runServe(ctx context.Context, cfg config.Config) error {
|
||||
loginLocks := auth.NewLoginLocks()
|
||||
memConns := message.NewMemoryConns()
|
||||
msgLim := message.LimitsFromFullConfig(cfg)
|
||||
metricsReg := metrics.New()
|
||||
db.Queue.OnBatchCommit = func(d time.Duration) {
|
||||
metricsReg.WriteCommitSeconds.Observe(d.Seconds())
|
||||
}
|
||||
|
||||
login := broker.NewLogin(broker.LoginOptions{
|
||||
DB: db,
|
||||
@@ -84,11 +88,13 @@ func runServe(ctx context.Context, cfg config.Config) error {
|
||||
msgApp := message.New(db, msgLim, hashPool,
|
||||
message.WithLocks(loginLocks),
|
||||
message.WithConnRegistry(memConns),
|
||||
message.WithMetrics(metricsReg),
|
||||
)
|
||||
uplink := &appUplink{
|
||||
msg: msgApp,
|
||||
conns: memConns,
|
||||
log: slog.Default(),
|
||||
msg: msgApp,
|
||||
conns: memConns,
|
||||
log: slog.Default(),
|
||||
metrics: metricsReg,
|
||||
}
|
||||
sess := broker.NewSession(broker.SessionOptions{
|
||||
Login: login,
|
||||
@@ -109,6 +115,7 @@ func runServe(ctx context.Context, cfg config.Config) error {
|
||||
Authenticator: login,
|
||||
Uplink: sess,
|
||||
Logger: slog.Default(),
|
||||
Metrics: metricsReg,
|
||||
OnPublishDropped: func(dropCtx context.Context, endpointID string, connID port.ConnID, payload []byte) {
|
||||
if dropErr := msgApp.OnPublishDropped(dropCtx, endpointID, connID, payload); dropErr != nil {
|
||||
slog.Error("on publish dropped", "endpoint", endpointID, "err", dropErr)
|
||||
@@ -125,6 +132,7 @@ func runServe(ctx context.Context, cfg config.Config) error {
|
||||
message.WithLocks(loginLocks),
|
||||
message.WithConnRegistry(memConns),
|
||||
message.WithDownlink(brk),
|
||||
message.WithMetrics(metricsReg),
|
||||
)
|
||||
uplink.msg = msgApp
|
||||
uplink.down = brk
|
||||
@@ -217,7 +225,6 @@ func runServe(ctx context.Context, cfg config.Config) error {
|
||||
},
|
||||
})
|
||||
|
||||
metricsReg := metrics.New()
|
||||
buildHandlers := func(proxies *listener.ProxySet) listener.Handlers {
|
||||
return listener.Handlers{
|
||||
MQTT: brk.WSHandler(proxies),
|
||||
@@ -299,7 +306,7 @@ func runServe(ctx context.Context, cfg config.Config) error {
|
||||
|
||||
loopCtx, loopCancel := context.WithCancel(ctx)
|
||||
defer loopCancel()
|
||||
go messageLoops(loopCtx, msgApp, memConns)
|
||||
go messageLoops(loopCtx, msgApp, memConns, db, hashPool, metricsReg)
|
||||
|
||||
<-ctx.Done()
|
||||
loopCancel()
|
||||
@@ -312,7 +319,7 @@ func runServe(ctx context.Context, cfg config.Config) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
func messageLoops(ctx context.Context, msgApp *message.App, conns *message.MemoryConns) {
|
||||
func messageLoops(ctx context.Context, msgApp *message.App, conns *message.MemoryConns, db *store.DB, hashPool auth.HashPool, met *metrics.Registry) {
|
||||
t := time.NewTicker(time.Second)
|
||||
defer t.Stop()
|
||||
for {
|
||||
@@ -332,6 +339,10 @@ func messageLoops(ctx context.Context, msgApp *message.App, conns *message.Memor
|
||||
if err := msgApp.CleanupOnce(ctx, nowMs); err != nil {
|
||||
slog.Error("cleanup once", "err", err)
|
||||
}
|
||||
if err := metrics.SampleStoreGauges(ctx, met, db.Read); err != nil {
|
||||
slog.Debug("sample store gauges", "err", err)
|
||||
}
|
||||
metrics.SampleQueues(met, db.Queue.Len(), hashPool.QueueLen())
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user