diff --git a/cmd/nixmsg/serve.go b/cmd/nixmsg/serve.go index 876288a..fd94958 100644 --- a/cmd/nixmsg/serve.go +++ b/cmd/nixmsg/serve.go @@ -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 @@ -184,7 +192,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), @@ -266,7 +273,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() @@ -279,7 +286,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 { @@ -299,6 +306,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()) } } } diff --git a/cmd/nixmsg/uplink.go b/cmd/nixmsg/uplink.go index 9ca0dca..1f429ff 100644 --- a/cmd/nixmsg/uplink.go +++ b/cmd/nixmsg/uplink.go @@ -11,6 +11,7 @@ import ( "git.asio.asia/nixevol/NixMsg/internal/app/message" "git.asio.asia/nixevol/NixMsg/internal/app/port" "git.asio.asia/nixevol/NixMsg/internal/app/presence" + "git.asio.asia/nixevol/NixMsg/internal/metrics" "git.asio.asia/nixevol/NixMsg/internal/protocol" ) @@ -23,6 +24,7 @@ type appUplink struct { conns *message.MemoryConns down port.Downlink log *slog.Logger + metrics *metrics.Registry } func (u *appUplink) OnSessionEstablished(ctx context.Context, conn port.ConnInfo) error { @@ -255,6 +257,9 @@ func (u *appUplink) replyErr(ctx context.Context, conn port.ConnInfo, rid, code, if rid == "" { rid = "0" } + if u.metrics != nil && code != "" { + u.metrics.ErrorsTotal.WithLabelValues(code).Inc() + } resp := protocol.Resp{ V: protocol.Version, Type: protocol.TypeResp, diff --git a/docs/DEVIATIONS.md b/docs/DEVIATIONS.md index 3ebcd47..9824601 100644 --- a/docs/DEVIATIONS.md +++ b/docs/DEVIATIONS.md @@ -1092,3 +1092,12 @@ - 原因:同一连接上 `group.create`/`group.add` 同步向本连接注入下行时,与 mochi InlineClient 互相等待,`resp` 回不去(`TestUplinkDMOfflineGroupRecall` 在清掉测试客户端 dial deadline 后稳定复现)。 - 备选方案:broker 层对 Inline 发布做无锁队列。 - 影响:`group_event` 可能略晚于 `resp` 到达;业务结果仍以 `resp` 为准。 + +### fix-issue-5 + +1. **指标在真实事件点打点,不新造名字** + - 原条款:PRD F22 / DEVELOPMENT 4.3 / issue #5;DEVIATIONS P4 已定名但从未接线。 + - 实际做法:`nixmsg_connections{transport}` 在 broker `OnSessionEstablished`/`OnDisconnect` 末尾 Inc/Dec;`endpoints`/`deliveries_pending`/`messages_scheduled` 与写队列、哈希排队在 `messageLoops` 每秒按库/队列真实长度采样;`dispatch_to_push`/`ack` 直方图在成功推送与确认路径 Observe;写批提交耗时经 `store.Queue.OnBatchCommit`;`errors_total` 仅在上行 `replyErr` 时按错误码递增。门禁不变。 + - 原因:空指标等于监控未交付;采样避免在每条写路径上改大段分发逻辑,并减小与 #3/#6 的合并面。 + - 备选方案:全部改为纯事件加减(pending 等需在每处状态迁移维护计数)。 + - 影响:仪表盘按既有名字即可看在线连接与待投递;无对应事件时计数保持 0,不做假数。 diff --git a/internal/app/message/ack.go b/internal/app/message/ack.go index 5f03f11..7c4c1b4 100644 --- a/internal/app/message/ack.go +++ b/internal/app/message/ack.go @@ -17,6 +17,8 @@ func (a *App) Ack(ctx context.Context, endpointID string, req *protocol.Ack) (Ac nowMs := a.now().UnixMilli() var out AckResult var seq int64 + var ackLatencySec float64 + var observeAck bool err := a.db.Queue.Do(ctx, func(tx *sql.Tx) error { err := tx.QueryRow(`SELECT seq FROM messages WHERE sender_id = ? AND id = ?`, req.From, req.ID).Scan(&seq) if err == sql.ErrNoRows { @@ -25,6 +27,10 @@ func (a *App) Ack(ctx context.Context, endpointID string, req *protocol.Ack) (Ac if err != nil { return err } + var pushedAt sql.NullInt64 + _ = tx.QueryRow(` +SELECT pushed_at FROM deliveries +WHERE seq = ? AND endpoint_id = ? AND state = 'pending'`, seq, endpointID).Scan(&pushedAt) res, err := tx.Exec(` UPDATE deliveries SET state = ?, reason = '', pushed_conn = NULL, updated_at = ? WHERE seq = ? AND endpoint_id = ? AND state = 'pending'`, @@ -35,6 +41,10 @@ WHERE seq = ? AND endpoint_id = ? AND state = 'pending'`, aff, _ := res.RowsAffected() if aff > 0 { out.Result = DeliveryAccepted + if pushedAt.Valid && pushedAt.Int64 > 0 && nowMs >= pushedAt.Int64 { + ackLatencySec = float64(nowMs-pushedAt.Int64) / 1000.0 + observeAck = true + } if e := insertReceiptTx(tx, req.From, seq, endpointID, DeliveryAccepted, "", nowMs); e != nil { return e } @@ -55,6 +65,9 @@ SELECT state FROM deliveries WHERE seq = ? AND endpoint_id = ?`, seq, endpointID if err != nil { return out, err } + if observeAck && a.met != nil { + a.met.AckSeconds.Observe(ackLatencySec) + } if out.Result == DeliveryAccepted { a.releaseLarge(seq, endpointID) } diff --git a/internal/app/message/app.go b/internal/app/message/app.go index 89930ef..ae9d39f 100644 --- a/internal/app/message/app.go +++ b/internal/app/message/app.go @@ -7,6 +7,7 @@ import ( "git.asio.asia/nixevol/NixMsg/internal/app/port" "git.asio.asia/nixevol/NixMsg/internal/auth" "git.asio.asia/nixevol/NixMsg/internal/config" + "git.asio.asia/nixevol/NixMsg/internal/metrics" "git.asio.asia/nixevol/NixMsg/internal/protocol" "git.asio.asia/nixevol/NixMsg/internal/store" ) @@ -86,6 +87,7 @@ type App struct { down port.Downlink conns ConnRegistry + met *metrics.Registry mu sync.Mutex largeSem chan struct{} @@ -117,6 +119,11 @@ func WithConnRegistry(c ConnRegistry) Option { return func(a *App) { a.conns = c } } +// WithMetrics 注入 Prometheus 注册表(投递耗时直方图)。 +func WithMetrics(m *metrics.Registry) Option { + return func(a *App) { a.met = m } +} + // New 创建消息服务实现。 func New(db *store.DB, lim Limits, hash auth.HashPool, opts ...Option) *App { if lim.RequestBurst <= 0 { diff --git a/internal/app/message/push.go b/internal/app/message/push.go index be11c47..aa9d557 100644 --- a/internal/app/message/push.go +++ b/internal/app/message/push.go @@ -229,11 +229,20 @@ WHERE seq = ? AND endpoint_id = ? AND state = 'pending' AND pushed_conn IS NULL` a.releaseLarge(it.seq, endpointID) } a.scheduleRepush(endpointID, time.Second) + } else { + a.observeDispatchToPush(it.sendAt, nowMs) } } return a.pushReceipts(ctx, endpointID, connID, nowMs) } +func (a *App) observeDispatchToPush(sendAtMs, pushedAtMs int64) { + if a.met == nil || pushedAtMs < sendAtMs { + return + } + a.met.DispatchToPushSeconds.Observe(float64(pushedAtMs-sendAtMs) / 1000.0) +} + func (a *App) rejectTooLarge(ctx context.Context, seq int64, endpointID, senderID string, nowMs int64) error { return a.db.Queue.Do(ctx, func(tx *sql.Tx) error { res, err := tx.Exec(` diff --git a/internal/broker/broker.go b/internal/broker/broker.go index f4989b3..0ecfdd5 100644 --- a/internal/broker/broker.go +++ b/internal/broker/broker.go @@ -12,6 +12,7 @@ import ( "time" "git.asio.asia/nixevol/NixMsg/internal/app/port" + "git.asio.asia/nixevol/NixMsg/internal/metrics" mqtt "github.com/mochi-mqtt/server/v2" "github.com/mochi-mqtt/server/v2/packets" ) @@ -68,15 +69,18 @@ type Options struct { Logger *slog.Logger // OnPublishDropped 可选;nil 时仅打 debug 日志。 OnPublishDropped PublishDroppedFunc + // Metrics 可选;会话建立/断开时更新 nixmsg_connections。 + Metrics *metrics.Registry } // Broker 内置 mochi,不自带监听端口。 type Broker struct { - server *mqtt.Server - auth Authenticator - uplink port.UplinkHandler - log *slog.Logger - onDrop PublishDroppedFunc + server *mqtt.Server + auth Authenticator + uplink port.UplinkHandler + log *slog.Logger + onDrop PublishDroppedFunc + metrics *metrics.Registry hook *nixHook @@ -105,6 +109,7 @@ type connState struct { handshook bool subscribedDown bool largeHeld int + metricsCounted bool mu sync.Mutex handshakeTimer *time.Timer @@ -151,6 +156,7 @@ func New(opts Options) (*Broker, error) { uplink: uplink, log: log, onDrop: opts.OnPublishDropped, + metrics: opts.Metrics, current: make(map[string]*connState), byClient: make(map[*mqtt.Client]*connState), queues: make(map[string]*uplinkQueue), diff --git a/internal/broker/hooks.go b/internal/broker/hooks.go index d503274..758da07 100644 --- a/internal/broker/hooks.go +++ b/internal/broker/hooks.go @@ -190,6 +190,7 @@ func (h *nixHook) OnSessionEstablished(cl *mqtt.Client, _ packets.Packet) { MaxPacketSize: st.maxPacketSize, } _ = h.b.uplink.OnSessionEstablished(context.Background(), info) + h.noteConnectionOpen(st) } func (h *nixHook) OnDisconnect(cl *mqtt.Client, err error, _ bool) { @@ -229,9 +230,27 @@ func (h *nixHook) OnDisconnect(cl *mqtt.Client, err error, _ bool) { } if sess, ok := h.b.uplink.(*Session); ok { sess.HandleDisconnect(context.Background(), info, reason, isCurrent) + h.noteConnectionClose(st) return } h.b.uplink.OnDisconnect(context.Background(), info, reason) + h.noteConnectionClose(st) +} + +func (h *nixHook) noteConnectionOpen(st *connState) { + if h.b.metrics == nil || st == nil || st.metricsCounted { + return + } + h.b.metrics.Connections.WithLabelValues(string(st.transport)).Inc() + st.metricsCounted = true +} + +func (h *nixHook) noteConnectionClose(st *connState) { + if h.b.metrics == nil || st == nil || !st.metricsCounted { + return + } + h.b.metrics.Connections.WithLabelValues(string(st.transport)).Dec() + st.metricsCounted = false } func (h *nixHook) OnQosComplete(cl *mqtt.Client, pk packets.Packet) { diff --git a/internal/broker/metrics_test.go b/internal/broker/metrics_test.go new file mode 100644 index 0000000..dde96fb --- /dev/null +++ b/internal/broker/metrics_test.go @@ -0,0 +1,118 @@ +package broker + +import ( + "io" + "net" + "net/http" + "net/http/httptest" + "strings" + "testing" + "time" + + "git.asio.asia/nixevol/NixMsg/internal/metrics" + dto "github.com/prometheus/client_model/go" +) + +func TestConnectionMetricsIncDec(t *testing.T) { + reg := metrics.New() + b, err := New(Options{Authenticator: AllowAuthenticator{}, Metrics: reg}) + if err != nil { + t.Fatal(err) + } + defer func() { _ = b.Close() }() + + if got := gaugeValue(t, reg, "nixmsg_connections", "tcp"); got != 0 { + t.Fatalf("before connect tcp=%v", got) + } + + r, w := net.Pipe() + done := make(chan struct{}) + go func() { + defer close(done) + _ = b.AttachTCP(r) + }() + connectAndSubscribe(t, w, "ep-metrics", 0) + + deadline := time.Now().Add(3 * time.Second) + for { + if _, ok := b.ConnInfoOf("ep-metrics"); ok { + break + } + if time.Now().After(deadline) { + t.Fatal("session not established") + } + time.Sleep(10 * time.Millisecond) + } + if got := gaugeValue(t, reg, "nixmsg_connections", "tcp"); got != 1 { + t.Fatalf("after connect tcp=%v want 1", got) + } + + _ = w.Close() + select { + case <-done: + case <-time.After(3 * time.Second): + t.Fatal("attach did not finish") + } + deadline = time.Now().Add(3 * time.Second) + for { + if got := gaugeValue(t, reg, "nixmsg_connections", "tcp"); got == 0 { + break + } + if time.Now().After(deadline) { + t.Fatalf("after disconnect tcp=%v want 0", gaugeValue(t, reg, "nixmsg_connections", "tcp")) + } + time.Sleep(10 * time.Millisecond) + } +} + +func TestWSConnectionMetrics(t *testing.T) { + reg := metrics.New() + b, err := New(Options{Authenticator: AllowAuthenticator{}, Metrics: reg}) + if err != nil { + t.Fatal(err) + } + defer func() { _ = b.Close() }() + + mux := http.NewServeMux() + mux.Handle("/mqtt", b.WSHandler(nil)) + srv := httptest.NewServer(mux) + defer srv.Close() + + // 仅验证 handler 暴露指标文本仍含初始标签;建连用 TCP 测即可。 + req := httptest.NewRequest(http.MethodGet, "/metrics", nil) + rec := httptest.NewRecorder() + reg.Handler().ServeHTTP(rec, req) + body, _ := io.ReadAll(rec.Body) + if !strings.Contains(string(body), `nixmsg_connections{transport="ws"} 0`) { + t.Fatalf("missing ws series: %s", body) + } +} + +func gaugeValue(t *testing.T, reg *metrics.Registry, name, transport string) float64 { + t.Helper() + mfs, err := reg.Gatherer().Gather() + if err != nil { + t.Fatal(err) + } + for _, mf := range mfs { + if mf.GetName() != name { + continue + } + for _, m := range mf.GetMetric() { + if matchLabel(m, "transport", transport) { + return m.GetGauge().GetValue() + } + } + } + t.Fatalf("metric %s transport=%s not found", name, transport) + return 0 +} + +func matchLabel(m *dto.Metric, key, val string) bool { + for _, lp := range m.GetLabel() { + if lp.GetName() == key && lp.GetValue() == val { + return true + } + } + return false +} diff --git a/internal/metrics/sample.go b/internal/metrics/sample.go new file mode 100644 index 0000000..1a3eaa3 --- /dev/null +++ b/internal/metrics/sample.go @@ -0,0 +1,36 @@ +package metrics + +import ( + "context" + "database/sql" +) + +// SampleStoreGauges 按库内真实计数刷新端总数、待投递、定时消息(无对应行则为 0)。 +func SampleStoreGauges(ctx context.Context, r *Registry, db *sql.DB) error { + if r == nil || db == nil { + return nil + } + var endpoints, pending, scheduled int + if err := db.QueryRowContext(ctx, `SELECT COUNT(*) FROM endpoints`).Scan(&endpoints); err != nil { + return err + } + if err := db.QueryRowContext(ctx, `SELECT COUNT(*) FROM deliveries WHERE state = 'pending'`).Scan(&pending); err != nil { + return err + } + if err := db.QueryRowContext(ctx, `SELECT COUNT(*) FROM messages WHERE state = 'scheduled'`).Scan(&scheduled); err != nil { + return err + } + r.EndpointsTotal.Set(float64(endpoints)) + r.DeliveriesPending.Set(float64(pending)) + r.MessagesScheduled.Set(float64(scheduled)) + return nil +} + +// SampleQueues 刷新写队列与密码哈希排队长度(传入当前真实长度,不做估算)。 +func SampleQueues(r *Registry, writeQueueLen, passwordHashQueueLen int) { + if r == nil { + return + } + r.WriteQueueLength.Set(float64(writeQueueLen)) + r.PasswordHashQueue.Set(float64(passwordHashQueueLen)) +} diff --git a/internal/metrics/sample_test.go b/internal/metrics/sample_test.go new file mode 100644 index 0000000..9ca54fd --- /dev/null +++ b/internal/metrics/sample_test.go @@ -0,0 +1,82 @@ +package metrics + +import ( + "context" + "path/filepath" + "testing" + + "git.asio.asia/nixevol/NixMsg/internal/store" +) + +func TestSampleStoreGaugesPendingNonZero(t *testing.T) { + t.Parallel() + dir := t.TempDir() + db, err := store.Open(filepath.Join(dir, "data"), "FULL") + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = db.Close() }) + + ctx := context.Background() + _, err = db.Write.ExecContext(ctx, ` +INSERT INTO endpoints(id, name, login_hash, enabled, source, created_at) +VALUES ('alice', 'A', 'x', 1, 'admin', 1)`) + if err != nil { + t.Fatal(err) + } + _, err = db.Write.ExecContext(ctx, ` +INSERT INTO messages(sender_id, id, dest_kind, dest_id, send_at, keep, ttl_seconds, receipt, state, reason, created_at, content_type, body_enc, meta) +VALUES ('alice', 'm1', 'endpoint', 'bob', 100, 1, 0, 0, 'dispatched', '', 1, 'text/plain', 'utf8', '{}')`) + if err != nil { + t.Fatal(err) + } + var seq int64 + if err := db.Write.QueryRowContext(ctx, `SELECT seq FROM messages WHERE id='m1'`).Scan(&seq); err != nil { + t.Fatal(err) + } + _, err = db.Write.ExecContext(ctx, ` +INSERT INTO deliveries(seq, endpoint_id, send_at, keep, state, reason, attempts, updated_at) +VALUES (?, 'bob', 100, 1, 'pending', '', 0, 1)`, seq) + if err != nil { + t.Fatal(err) + } + _, err = db.Write.ExecContext(ctx, ` +INSERT INTO messages(sender_id, id, dest_kind, dest_id, send_at, keep, ttl_seconds, receipt, state, reason, created_at, content_type, body_enc, meta) +VALUES ('alice', 'm2', 'endpoint', 'bob', 999999, 0, 0, 0, 'scheduled', '', 1, 'text/plain', 'utf8', '{}')`) + if err != nil { + t.Fatal(err) + } + + reg := New() + if err := SampleStoreGauges(ctx, reg, db.Read); err != nil { + t.Fatal(err) + } + SampleQueues(reg, 3, 1) + + mfs, err := reg.Gatherer().Gather() + if err != nil { + t.Fatal(err) + } + got := map[string]float64{} + for _, mf := range mfs { + switch mf.GetName() { + case "nixmsg_endpoints", "nixmsg_deliveries_pending", "nixmsg_messages_scheduled", + "nixmsg_write_queue_length", "nixmsg_password_hash_queue_length": + if len(mf.GetMetric()) > 0 { + got[mf.GetName()] = mf.GetMetric()[0].GetGauge().GetValue() + } + } + } + if got["nixmsg_endpoints"] != 1 { + t.Fatalf("endpoints=%v", got["nixmsg_endpoints"]) + } + if got["nixmsg_deliveries_pending"] != 1 { + t.Fatalf("pending=%v", got["nixmsg_deliveries_pending"]) + } + if got["nixmsg_messages_scheduled"] != 1 { + t.Fatalf("scheduled=%v", got["nixmsg_messages_scheduled"]) + } + if got["nixmsg_write_queue_length"] != 3 || got["nixmsg_password_hash_queue_length"] != 1 { + t.Fatalf("queues=%v", got) + } +} diff --git a/internal/store/queue.go b/internal/store/queue.go index a197194..567ff4c 100644 --- a/internal/store/queue.go +++ b/internal/store/queue.go @@ -43,6 +43,9 @@ type Queue struct { ready bool lastWriteErr error pending int + + // OnBatchCommit 可选;每次合并提交成功后回调耗时(秒级指标用)。 + OnBatchCommit func(d time.Duration) } // NewQueue 创建合并写入队列并启动写 goroutine。 @@ -129,6 +132,7 @@ func (q *Queue) loop() { } func (q *Queue) runBatch(batch []writeJob) { + started := time.Now() defer func() { q.mu.Lock() q.pending -= len(batch) @@ -229,6 +233,9 @@ func (q *Queue) runBatch(batch []writeJob) { } return } + if q.OnBatchCommit != nil { + q.OnBatchCommit(time.Since(started)) + } for _, o := range outcomes { if o.success { o.job.res <- nil