From 706ed627faa08925b8bab82fd168ac953e1cb93f Mon Sep 17 00:00:00 2001 From: andig Date: Tue, 14 Apr 2020 17:53:11 +0200 Subject: [PATCH] Support writing through mqtt --- provider/config.go | 34 ++++++++++++++++++------------- provider/helper.go | 51 +++++++++++++++++++++++++--------------------- provider/mqtt.go | 40 ++++++++++++++++++++++++++++++++++++ 3 files changed, 88 insertions(+), 37 deletions(-) diff --git a/provider/config.go b/provider/config.go index ad2f4c35b..b7b7441c5 100644 --- a/provider/config.go +++ b/provider/config.go @@ -20,7 +20,7 @@ type Config struct { // mqttConfig is the specific mqtt getter/setter configuration type mqttConfig struct { - Topic string + Topic, Payload string // Payload only applies to setters Multiplier float64 Timeout time.Duration } @@ -140,22 +140,12 @@ func NewBoolGetterFromConfig(log *api.Logger, config Config) (res BoolGetter) { return } -// NewBoolSetterFromConfig creates a BoolSetter from config -func NewBoolSetterFromConfig(log *api.Logger, param string, config Config) (res BoolSetter) { - switch strings.ToLower(config.Type) { - case "script": - pc := scriptFromConfig(log, config.Other) - exec := NewScriptProvider(pc.Timeout) - res = exec.BoolSetter(param, pc.Cmd) - default: - log.FATAL.Fatalf("invalid setter type %s", config.Type) - } - return -} - // NewIntSetterFromConfig creates a IntSetter from config func NewIntSetterFromConfig(log *api.Logger, param string, config Config) (res IntSetter) { switch strings.ToLower(config.Type) { + case "mqtt": + pc := mqttFromConfig(log, config.Other) + res = MQTT.IntSetter(param, pc.Topic, pc.Payload) case "script": pc := scriptFromConfig(log, config.Other) exec := NewScriptProvider(pc.Timeout) @@ -165,3 +155,19 @@ func NewIntSetterFromConfig(log *api.Logger, param string, config Config) (res I } return } + +// NewBoolSetterFromConfig creates a BoolSetter from config +func NewBoolSetterFromConfig(log *api.Logger, param string, config Config) (res BoolSetter) { + switch strings.ToLower(config.Type) { + case "mqtt": + pc := mqttFromConfig(log, config.Other) + res = MQTT.BoolSetter(param, pc.Topic, pc.Payload) + case "script": + pc := scriptFromConfig(log, config.Other) + exec := NewScriptProvider(pc.Timeout) + res = exec.BoolSetter(param, pc.Cmd) + default: + log.FATAL.Fatalf("invalid setter type %s", config.Type) + } + return +} diff --git a/provider/helper.go b/provider/helper.go index 7358cce55..8c015ec84 100644 --- a/provider/helper.go +++ b/provider/helper.go @@ -14,32 +14,37 @@ func truish(s string) bool { return s == "1" || strings.ToLower(s) == "true" || strings.ToLower(s) == "on" } -// replaceFormatted replaces all occurrances of ${key} with val from the kv map. -// All keys of kv must exist inside the string to apply replacements to +// formatValue will apply specific formatting in addition to standard sprintf +func formatValue(format string, val interface{}) string { + switch val := val.(type) { + case bool: + if format == "%d" { + if val { + return "1" + } + return "0" + } + } + + if format == "" { + format = "%v" + } + + return fmt.Sprintf(format, val) +} + +// replaceFormatted replaces all occurrences of ${key} with formatted val from the kv map func replaceFormatted(s string, kv map[string]interface{}) (string, error) { - matches := re.FindAllStringSubmatch(s, -1) - - for len(matches) > 0 { - for _, m := range matches { - key := m[1] - val, ok := kv[key] - if !ok { - return "", errors.New("could not find match for " + m[0]) - } - - // apply format - format := m[3] - if format != "" { - val = fmt.Sprintf(format, val) - } - - // update string - literalMatch := m[0] - s = strings.ReplaceAll(s, literalMatch, fmt.Sprintf("%v", val)) + for m := re.FindStringSubmatch(s); m != nil; m = re.FindStringSubmatch(s) { + // find key and replacement value + val, ok := kv[m[1]] + if !ok { + return "", errors.New("could find value for: " + m[0]) } - // update matches - matches = re.FindAllStringSubmatch(s, -1) + // update all literal matches + new := formatValue(m[3], val) + s = strings.ReplaceAll(s, m[0], new) } return s, nil diff --git a/provider/mqtt.go b/provider/mqtt.go index dd5622c70..0fbafec86 100644 --- a/provider/mqtt.go +++ b/provider/mqtt.go @@ -148,6 +148,46 @@ func (m *MqttClient) BoolGetter(topic string, timeout time.Duration) BoolGetter return h.boolGetter } +// IntSetter publishes topic with parameter replaced by int value +func (m *MqttClient) IntSetter(param, topic, message string) IntSetter { + return func(v int64) error { + payload, err := replaceFormatted(message, map[string]interface{}{ + param: v, + }) + if err != nil { + return err + } + + mlog.TRACE.Printf("send %s: '%s'", topic, payload) + token := m.Client.Publish(topic, m.qos, false, payload) + if token.WaitTimeout(publishTimeout) { + return token.Error() + } + + return fmt.Errorf("%s send timeout", topic) + } +} + +// BoolSetter invokes script with parameter replaced by bool value +func (m *MqttClient) BoolSetter(param, topic, message string) BoolSetter { + return func(v bool) error { + payload, err := replaceFormatted(message, map[string]interface{}{ + param: v, + }) + if err != nil { + return err + } + + mlog.TRACE.Printf("send %s: '%s'", topic, payload) + token := m.Client.Publish(topic, m.qos, false, payload) + if token.WaitTimeout(publishTimeout) { + return token.Error() + } + + return fmt.Errorf("%s send timeout", topic) + } +} + // WaitForToken synchronously waits until token operation completed func (m *MqttClient) WaitForToken(token mqtt.Token) { if token.WaitTimeout(publishTimeout) {