diff --git a/internal/mqttforward/dedup_queue.go b/internal/mqttforward/dedup_queue.go index 91816a7..8a80cd3 100644 --- a/internal/mqttforward/dedup_queue.go +++ b/internal/mqttforward/dedup_queue.go @@ -53,6 +53,12 @@ func (dq *DedupQueue) Stop() { close(dq.stopCh) } +func (dq *DedupQueue) Len() int { + dq.mu.Lock() + defer dq.mu.Unlock() + return len(dq.entries) +} + func (dq *DedupQueue) cleanup() { now := time.Now() dq.mu.Lock() diff --git a/internal/web/mqtt_status.go b/internal/web/mqtt_status.go index 5b750c8..3f9848b 100644 --- a/internal/web/mqtt_status.go +++ b/internal/web/mqtt_status.go @@ -30,6 +30,7 @@ type MQTTRuntimeStatus struct { Stats *mqttforwardpkg.Stats ClientStats *mqttforwardpkg.ClientStats DBQueue *storepkg.WriteQueue + DedupQueue *mqttforwardpkg.DedupQueue } // AdminMQTTStatus 是 admin 路由 GET /admin/mqtt-status 返回的 JSON 视图。 @@ -50,6 +51,7 @@ type AdminMQTTStatus struct { MessagesSent int64 `json:"messages_sent"` MessagesDropped int64 `json:"messages_dropped"` DBWriteQueueLength int `json:"db_write_queue_length"` + DedupQueueLength int `json:"dedup_queue_len"` Retained int64 `json:"retained"` Inflight int64 `json:"inflight"` InflightDropped int64 `json:"inflight_dropped"` @@ -71,7 +73,7 @@ type AdminMQTTClient struct { // Status 实现 MQTTStatusProvider。 func (m MQTTRuntimeStatus) Status() AdminMQTTStatus { if m.Server == nil || m.Server.Info == nil { - return AdminMQTTStatus{Running: false, Address: m.Address, TLS: m.TLS, DBWriteQueueLength: m.DBQueue.Len()} + return AdminMQTTStatus{Running: false, Address: m.Address, TLS: m.TLS, DBWriteQueueLength: m.DBQueue.Len(), DedupQueueLength: m.dedupQueueLen()} } info := m.Server.Info.Clone() status := AdminMQTTStatus{ @@ -91,6 +93,7 @@ func (m MQTTRuntimeStatus) Status() AdminMQTTStatus { MessagesSent: m.Stats.Forwarded(), MessagesDropped: m.Stats.Dropped(), DBWriteQueueLength: m.DBQueue.Len(), + DedupQueueLength: m.dedupQueueLen(), Retained: info.Retained, Inflight: info.Inflight, InflightDropped: info.InflightDropped, @@ -136,6 +139,13 @@ func mqttClientInfo(c *mqtt.Client) mqttClientInfoView { } } +func (m MQTTRuntimeStatus) dedupQueueLen() int { + if m.DedupQueue == nil { + return 0 + } + return m.DedupQueue.Len() +} + // DisconnectClient 实现 MQTTStatusProvider:发送 Disconnect 报文并关闭连接。 // 使用 ErrAdministrativeAction 作为断开理由,便于日志区分。 func (m MQTTRuntimeStatus) DisconnectClient(clientID string) bool { diff --git a/main.go b/main.go index b129e96..6d1eea5 100644 --- a/main.go +++ b/main.go @@ -474,7 +474,7 @@ func run(cfg *configpkg.Config) error { if err != nil { return err } - mqttStatus := webpkg.MQTTRuntimeStatus{Server: server, Address: mqttAddr, TLS: cfg.MQTT.TLS.Enabled, Stats: messageStats, ClientStats: clientStats, DBQueue: dbQueue} + mqttStatus := webpkg.MQTTRuntimeStatus{Server: server, Address: mqttAddr, TLS: cfg.MQTT.TLS.Enabled, Stats: messageStats, ClientStats: clientStats, DBQueue: dbQueue, DedupQueue: mqttHook.dedupQueue} handler := webpkg.NewRouter(cfg.Web, cfg.ConsoleLog.Web, store, sessions, mqttStatus, blocking, forwardManager, settings, botSender, aiService) webAddresses := []string{} if cfg.Web.PortEnabled { diff --git a/meshmap_frontend/src/components/AdminDashboard.vue b/meshmap_frontend/src/components/AdminDashboard.vue index f63b78b..3fc909b 100644 --- a/meshmap_frontend/src/components/AdminDashboard.vue +++ b/meshmap_frontend/src/components/AdminDashboard.vue @@ -200,6 +200,7 @@ onBeforeUnmount(() => {