Refactor mqtt broker connection handling (#2003)
This commit is contained in:
parent
5a929ff8b0
commit
6ad297a9f0
3 changed files with 7 additions and 14 deletions
|
|
@ -121,13 +121,9 @@ func configureDatabase(conf server.InfluxConfig, loadPoints []loadpoint.API, in
|
|||
// setup mqtt
|
||||
func configureMQTT(conf mqttConfig) error {
|
||||
log := util.NewLogger("mqtt")
|
||||
clientID := conf.ClientID
|
||||
if clientID == "" {
|
||||
clientID = mqtt.ClientID()
|
||||
}
|
||||
|
||||
var err error
|
||||
mqtt.Instance, err = mqtt.RegisteredClient(log, conf.Broker, conf.User, conf.Password, clientID, 1, func(options *paho.ClientOptions) {
|
||||
mqtt.Instance, err = mqtt.RegisteredClient(log, conf.Broker, conf.User, conf.Password, conf.ClientID, 1, func(options *paho.ClientOptions) {
|
||||
topic := fmt.Sprintf("%s/status", conf.RootTopic())
|
||||
options.SetWill(topic, "offline", 1, true)
|
||||
})
|
||||
|
|
|
|||
|
|
@ -29,10 +29,14 @@ var registry clientRegistry = make(map[string]*Client)
|
|||
|
||||
// RegisteredClient reuses an registered Mqtt publisher or creates a new one
|
||||
func RegisteredClient(log *util.Logger, broker, user, password, clientID string, qos byte, opts ...Option) (*Client, error) {
|
||||
key := fmt.Sprintf("%s.%s", broker, log.Name())
|
||||
key := fmt.Sprintf("%s.%s:%s", broker, user, password)
|
||||
client, err := registry.Get(key)
|
||||
|
||||
if err != nil {
|
||||
if clientID == "" {
|
||||
clientID = ClientID()
|
||||
}
|
||||
|
||||
if client, err = NewClient(log, broker, user, password, clientID, qos, opts...); err == nil {
|
||||
registry.Add(key, client)
|
||||
}
|
||||
|
|
@ -48,7 +52,7 @@ func RegisteredClientOrDefault(log *util.Logger, cc Config) (*Client, error) {
|
|||
client := Instance
|
||||
|
||||
if cc.Broker != "" {
|
||||
client, err = RegisteredClient(log, cc.Broker, cc.User, cc.Password, ClientID(), 1)
|
||||
client, err = RegisteredClient(log, cc.Broker, cc.User, cc.Password, cc.ClientID, 1)
|
||||
}
|
||||
|
||||
if client == nil && err == nil {
|
||||
|
|
|
|||
|
|
@ -30,7 +30,6 @@ var LogAreaPadding = 6
|
|||
// Logger wraps a jww notepad to avoid leaking implementation detail
|
||||
type Logger struct {
|
||||
*jww.Notepad
|
||||
name string
|
||||
*Redactor
|
||||
}
|
||||
|
||||
|
|
@ -55,7 +54,6 @@ func NewLogger(area string) *Logger {
|
|||
logger := &Logger{
|
||||
Notepad: notepad,
|
||||
Redactor: redactor,
|
||||
name: area,
|
||||
}
|
||||
|
||||
loggers[area] = logger
|
||||
|
|
@ -63,11 +61,6 @@ func NewLogger(area string) *Logger {
|
|||
return logger
|
||||
}
|
||||
|
||||
// Name returns the loggers name
|
||||
func (l *Logger) Name() string {
|
||||
return l.name
|
||||
}
|
||||
|
||||
// Redact adds items for redaction
|
||||
func (l *Logger) Redact(items ...string) *Logger {
|
||||
l.Redactor.Redact(items...)
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue