allow bridge mq cost msg (#162)

This commit is contained in:
ZhangJian He
2022-05-20 21:27:48 +08:00
committed by GitHub
parent 92758c8c85
commit 5dc2114daf
7 changed files with 21 additions and 14 deletions
+4 -2
View File
@@ -5,11 +5,13 @@ import (
"go.uber.org/zap" "go.uber.org/zap"
) )
func (b *Broker) Publish(e *bridge.Elements) { func (b *Broker) Publish(e *bridge.Elements) bool {
if b.bridgeMQ != nil { if b.bridgeMQ != nil {
err := b.bridgeMQ.Publish(e) cost, err := b.bridgeMQ.Publish(e)
if err != nil { if err != nil {
log.Error("send message to mq error.", zap.Error(err)) log.Error("send message to mq error.", zap.Error(err))
} }
return cost
} }
return false
} }
+5 -1
View File
@@ -414,7 +414,7 @@ func (c *client) processClientPublish(packet *packets.PublishPacket) {
} }
//publish to bridge mq //publish to bridge mq
c.broker.Publish(&bridge.Elements{ cost := c.broker.Publish(&bridge.Elements{
ClientID: c.info.clientID, ClientID: c.info.clientID,
Username: c.info.username, Username: c.info.username,
Action: bridge.Publish, Action: bridge.Publish,
@@ -423,6 +423,10 @@ func (c *client) processClientPublish(packet *packets.PublishPacket) {
Topic: topic, Topic: topic,
}) })
if cost {
return
}
switch packet.Qos { switch packet.Qos {
case QosAtMostOnce: case QosAtMostOnce:
c.ProcessPublishMessage(packet) c.ProcessPublishMessage(packet)
+2 -2
View File
@@ -15,7 +15,7 @@ func (c *client) SendInfo() {
} }
url := c.info.localIP + ":" + c.broker.config.Cluster.Port url := c.info.localIP + ":" + c.broker.config.Cluster.Port
infoMsg := NewInfo(c.broker.id, url, false) infoMsg := NewInfo(c.broker.id, url)
err := c.WriterPacket(infoMsg) err := c.WriterPacket(infoMsg)
if err != nil { if err != nil {
log.Error("send info message error, ", zap.Error(err)) log.Error("send info message error, ", zap.Error(err))
@@ -60,7 +60,7 @@ func (c *client) SendConnect() {
log.Info("send connect success") log.Info("send connect success")
} }
func NewInfo(sid, url string, isforword bool) *packets.PublishPacket { func NewInfo(sid, url string) *packets.PublishPacket {
pub := packets.NewControlPacket(packets.Publish).(*packets.PublishPacket) pub := packets.NewControlPacket(packets.Publish).(*packets.PublishPacket)
pub.Qos = 0 pub.Qos = 0
pub.TopicName = BrokerInfoTopic pub.TopicName = BrokerInfoTopic
+2 -1
View File
@@ -37,7 +37,8 @@ const (
) )
type BridgeMQ interface { type BridgeMQ interface {
Publish(e *Elements) error // Publish return true to cost the message
Publish(e *Elements) (bool, error)
} }
func NewBridgeMQ(name string) BridgeMQ { func NewBridgeMQ(name string) BridgeMQ {
+3 -3
View File
@@ -353,7 +353,7 @@ func (c *csvLog) logFilePrune() error {
// Publish implements the bridge interface - it accepts an Element then checks to see if that element is a // Publish implements the bridge interface - it accepts an Element then checks to see if that element is a
// message published to the admin topic for the plugin // message published to the admin topic for the plugin
// //
func (c *csvLog) Publish(e *Elements) error { func (c *csvLog) Publish(e *Elements) (bool, error) {
// A short-lived lock on c allows us to // A short-lived lock on c allows us to
// get the Command topic then release the lock // get the Command topic then release the lock
// This then allows us to process the command - which may // This then allows us to process the command - which may
@@ -372,7 +372,7 @@ func (c *csvLog) Publish(e *Elements) error {
// If the outfile is set to "{NULL}" we don't do anything with the message - we just return nil // If the outfile is set to "{NULL}" we don't do anything with the message - we just return nil
// This feature is here to allow CSVLOG to be enabled/disabled at runtime // This feature is here to allow CSVLOG to be enabled/disabled at runtime
if OutFile == "{NULL}" { if OutFile == "{NULL}" {
return nil return false, nil
} }
if e.Topic == CommandTopic { if e.Topic == CommandTopic {
@@ -410,5 +410,5 @@ func (c *csvLog) Publish(e *Elements) error {
// Push the message into the channel and return // Push the message into the channel and return
// the channel is buffered and is read by a goroutine so this should block for the shortest possible time // the channel is buffered and is read by a goroutine so this should block for the shortest possible time
c.msgchan <- e c.msgchan <- e
return nil return false, nil
} }
+3 -3
View File
@@ -63,7 +63,7 @@ func (k *kafka) connect() {
} }
//Publish publish to kafka //Publish publish to kafka
func (k *kafka) Publish(e *Elements) error { func (k *kafka) Publish(e *Elements) (bool, error) {
config := k.kafkaConfig config := k.kafkaConfig
key := e.ClientID key := e.ClientID
topics := make(map[string]bool) topics := make(map[string]bool)
@@ -96,10 +96,10 @@ func (k *kafka) Publish(e *Elements) error {
topics[config.DisconnectTopic] = true topics[config.DisconnectTopic] = true
} }
default: default:
return errors.New("error action: " + e.Action) return false, errors.New("error action: " + e.Action)
} }
return k.publish(topics, key, e) return false, k.publish(topics, key, e)
} }
+2 -2
View File
@@ -2,6 +2,6 @@ package bridge
type mockMQ struct{} type mockMQ struct{}
func (m *mockMQ) Publish(e *Elements) error { func (m *mockMQ) Publish(e *Elements) (bool, error) {
return nil return false, nil
} }