From 4782e84c159c4e060cd4ebfcc4e15b0be3ee1696 Mon Sep 17 00:00:00 2001 From: kevin Date: Tue, 30 Jun 2026 12:10:29 +0800 Subject: [PATCH] =?UTF-8?q?=E6=B7=BB=E5=8A=A0=20MQTT=20=E6=B6=88=E6=81=AF?= =?UTF-8?q?=E5=8E=BB=E9=87=8D=E9=98=9F=E5=88=97=EF=BC=9AHook=20=E5=B1=82?= =?UTF-8?q?=E5=9F=BA=E4=BA=8E=20payload+topic=20hash=20=E5=8E=BB=E9=87=8D?= =?UTF-8?q?=EF=BC=8CTTL=2015=20=E7=A7=92=EF=BC=8C=E5=AE=9A=E6=97=B6?= =?UTF-8?q?=E6=B8=85=E7=90=86?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- internal/mqttforward/dedup_queue.go | 72 +++++++++++++++++++++++++++++ main.go | 13 +++++- 2 files changed, 84 insertions(+), 1 deletion(-) create mode 100644 internal/mqttforward/dedup_queue.go diff --git a/internal/mqttforward/dedup_queue.go b/internal/mqttforward/dedup_queue.go new file mode 100644 index 0000000..91816a7 --- /dev/null +++ b/internal/mqttforward/dedup_queue.go @@ -0,0 +1,72 @@ +package mqttforward + +import ( + "crypto/sha256" + "encoding/hex" + "sync" + "time" +) + +const dedupTTL = 15 * time.Second + +type DedupQueue struct { + mu sync.Mutex + entries map[string]time.Time + stopCh chan struct{} +} + +func NewDedupQueue() *DedupQueue { + return &DedupQueue{ + entries: make(map[string]time.Time), + stopCh: make(chan struct{}), + } +} + +func (dq *DedupQueue) TryForward(topic string, payload []byte) bool { + hash := dedupHash(topic, payload) + now := time.Now() + dq.mu.Lock() + defer dq.mu.Unlock() + if expiry, ok := dq.entries[hash]; ok && now.Before(expiry) { + return false + } + dq.entries[hash] = now.Add(dedupTTL) + return true +} + +func (dq *DedupQueue) Start() { + go func() { + ticker := time.NewTicker(dedupTTL) + defer ticker.Stop() + for { + select { + case <-ticker.C: + dq.cleanup() + case <-dq.stopCh: + return + } + } + }() +} + +func (dq *DedupQueue) Stop() { + close(dq.stopCh) +} + +func (dq *DedupQueue) cleanup() { + now := time.Now() + dq.mu.Lock() + defer dq.mu.Unlock() + for hash, expiry := range dq.entries { + if now.After(expiry) { + delete(dq.entries, hash) + } + } +} + +func dedupHash(topic string, payload []byte) string { + h := sha256.New() + h.Write([]byte(topic)) + h.Write(payload) + return hex.EncodeToString(h.Sum(nil)) +} diff --git a/main.go b/main.go index ddf85d2..b129e96 100644 --- a/main.go +++ b/main.go @@ -57,8 +57,9 @@ type meshtasticFilterHook struct { settings *rspkg.Cache pkiResolver func(toNodeNum, fromNodeNum uint32) ([]byte, []byte, bool) autoAcker func(record map[string]any) - consoleLog bool // 控制台是否打印 MQTT 连接/订阅事件 + consoleLog bool // 控制台是否打印 MQTT 连接/订阅事件 packetConsoleLog bool // 控制台是否打印 Meshtastic 数据包 + dedupQueue *mqttforwardpkg.DedupQueue } // ID 返回用于识别 Meshtastic payload 过滤器的 hook 名称。 @@ -225,6 +226,10 @@ func (h *meshtasticFilterHook) OnPublish(cl *mqtt.Client, pk packets.Packet) (pa h.rejectPublish(cl, pk, record) return pk, packets.ErrRejectPacket } + if h.dedupQueue != nil && !h.dedupQueue.TryForward(pk.TopicName, pk.Payload) { + h.stats.IncDropped() + return pk, packets.ErrRejectPacket + } h.stats.IncForwarded() h.dbQueue.EnqueueRecord(record, mqttClientInfoFromClient(cl)) @@ -521,6 +526,9 @@ func run(cfg *configpkg.Config) error { if err := server.Close(); err != nil && runErr == nil { runErr = err } + if mqttHook.dedupQueue != nil { + mqttHook.dedupQueue.Stop() + } return runErr } @@ -529,6 +537,8 @@ func startMQTTServer(cfg *configpkg.Config, store *storepkg.Store, dbQueue *stor if err := server.AddHook(new(mqttauth.AllowHook), nil); err != nil { return nil, nil, "", err } + dedupQueue := mqttforwardpkg.NewDedupQueue() + dedupQueue.Start() hook := &meshtasticFilterHook{ server: server, key: cfg.Key, @@ -540,6 +550,7 @@ func startMQTTServer(cfg *configpkg.Config, store *storepkg.Store, dbQueue *stor pkiResolver: botpkg.NewPKIKeyResolver(store), consoleLog: cfg.ConsoleLog.MQTT, packetConsoleLog: cfg.ConsoleLog.Meshtastic, + dedupQueue: dedupQueue, } if err := server.AddHook(hook, nil); err != nil { return nil, nil, "", err