From d97e8d18e59e12c37d08ebf2066083024e40876a Mon Sep 17 00:00:00 2001 From: andig Date: Fri, 31 Dec 2021 16:11:10 +0100 Subject: [PATCH] Refactor data processing into pipeline component (#2169) --- meter/discovergy.go | 10 +- provider/http.go | 238 ++++------------------------------ provider/mqtt.go | 80 ++++-------- provider/mqtt_handler.go | 35 ++--- provider/pipeline/pipeline.go | 231 +++++++++++++++++++++++++++++++++ 5 files changed, 297 insertions(+), 297 deletions(-) create mode 100644 provider/pipeline/pipeline.go diff --git a/meter/discovergy.go b/meter/discovergy.go index a7fcae4f4..c566c1b13 100644 --- a/meter/discovergy.go +++ b/meter/discovergy.go @@ -6,6 +6,7 @@ import ( "github.com/evcc-io/evcc/api" "github.com/evcc-io/evcc/provider" + "github.com/evcc-io/evcc/provider/pipeline" "github.com/evcc-io/evcc/util" "github.com/evcc-io/evcc/util/request" "github.com/evcc-io/evcc/util/transport" @@ -71,13 +72,16 @@ func NewDiscovergyFromConfig(other map[string]interface{}) (api.Meter, error) { uri := fmt.Sprintf("%s/last_reading?meterId=%s", discovergyAPI, meterID) power, err := provider.NewHTTP(log, http.MethodGet, uri, false, 0.001*cc.Scale, 0).WithAuth("basic", cc.User, cc.Password) - if err == nil { - _, err = power.WithJq(".values.power") - } if err != nil { return nil, err } + pipe, err := new(pipeline.Pipeline).WithJq(".values.power") + if err != nil { + return nil, err + } + power = power.WithPipeline(pipe) + return NewConfigurable(power.FloatGetter()) } diff --git a/provider/http.go b/provider/http.go index 4e5b29833..c01b18262 100644 --- a/provider/http.go +++ b/provider/http.go @@ -1,26 +1,18 @@ package provider import ( - "bytes" - "encoding/hex" "fmt" "io" "math" - "regexp" "strconv" "strings" "time" - xj "github.com/basgys/goxml2json" - "github.com/evcc-io/evcc/provider/javascript" + "github.com/evcc-io/evcc/provider/pipeline" "github.com/evcc-io/evcc/util" - "github.com/evcc-io/evcc/util/jq" "github.com/evcc-io/evcc/util/request" "github.com/evcc-io/evcc/util/transport" - "github.com/itchyny/gojq" "github.com/jpfielding/go-http-digest/pkg/digest" - "github.com/robertkrimen/otto" - "github.com/volkszaehler/mbmd/meters/rs485" ) // HTTP implements HTTP request provider @@ -29,15 +21,10 @@ type HTTP struct { url, method string headers map[string]string body string - re *regexp.Regexp - jq *gojq.Query - unpack string - decode string - vm *otto.Otto - script string scale float64 cache time.Duration updated time.Time + pipeline *pipeline.Pipeline val []byte // Cached http response value err error // Cached http response error } @@ -54,20 +41,15 @@ type Auth struct { // NewHTTPProviderFromConfig creates a HTTP provider func NewHTTPProviderFromConfig(other map[string]interface{}) (IntProvider, error) { cc := struct { - URI, Method string - Headers map[string]string - Body string - Regex string - Jq string - Unpack string - Decode string - VM string - Script string - Scale float64 - Insecure bool - Auth Auth - Timeout time.Duration - Cache time.Duration + URI, Method string + Headers map[string]string + Body string + pipeline.Settings `mapstructure:",squash"` + Scale float64 + Insecure bool + Auth Auth + Timeout time.Duration + Cache time.Duration }{ Headers: make(map[string]string), Scale: 1, @@ -85,32 +67,14 @@ func NewHTTPProviderFromConfig(other map[string]interface{}) (IntProvider, error cc.Insecure, cc.Scale, cc.Cache, - ).WithHeaders(cc.Headers).WithBody(cc.Body) - http.Client.Timeout = cc.Timeout + ). + WithHeaders(cc.Headers). + WithBody(cc.Body) - var err error - if err == nil && cc.Regex != "" { - _, err = http.WithRegex(cc.Regex) - } - - if err == nil && cc.Jq != "" { - _, err = http.WithJq(cc.Jq) - } - - if err == nil && cc.Unpack != "" { - _, err = http.WithUnpack(cc.Unpack) - } - - if err == nil && cc.Decode != "" { - _, err = http.WithDecode(cc.Decode) - } - - if err == nil && cc.Script != "" { - _, err = http.WithScript(cc.VM, cc.Script) - } - - if err == nil && cc.Auth.Type != "" { - _, err = http.WithAuth(cc.Auth.Type, cc.Auth.User, cc.Auth.Password) + pipe, err := pipeline.New(cc.Settings) + if err == nil { + http = http.WithPipeline(pipe) + http.Client.Timeout = cc.Timeout } return http, err @@ -151,52 +115,10 @@ func (p *HTTP) WithHeaders(headers map[string]string) *HTTP { return p } -// WithRegex adds a regex query applied to the mqtt listener payload -func (p *HTTP) WithRegex(regex string) (*HTTP, error) { - re, err := regexp.Compile(regex) - if err != nil { - return nil, fmt.Errorf("invalid regex '%s': %w", re, err) - } - - p.re = re - - return p, nil -} - -// WithJq adds a jq query applied to the mqtt listener payload -func (p *HTTP) WithJq(jq string) (*HTTP, error) { - op, err := gojq.Parse(jq) - if err != nil { - return nil, fmt.Errorf("invalid jq query '%s': %w", jq, err) - } - - p.jq = op - - return p, nil -} - -// WithUnpack adds data unpacking -func (p *HTTP) WithUnpack(unpack string) (*HTTP, error) { - p.unpack = strings.ToLower(unpack) - - return p, nil -} - -// WithDecode adds data decoding -func (p *HTTP) WithDecode(decode string) (*HTTP, error) { - p.decode = strings.ToLower(decode) - - return p, nil -} - -// WithScript adds a javascript script to process the response -func (p *HTTP) WithScript(vm, script string) (*HTTP, error) { - regvm := javascript.RegisteredVM(strings.ToLower(vm)) - - p.vm = regvm - p.script = script - - return p, nil +// WithPipeline adds a processing pipeline +func (p *HTTP) WithPipeline(pipeline *pipeline.Pipeline) *HTTP { + p.pipeline = pipeline + return p } // WithAuth adds authorized transport @@ -237,72 +159,6 @@ func (p *HTTP) request(body ...string) ([]byte, error) { return p.val, p.err } -// transform XML into JSON with attribute names getting 'attr' prefix -func (p *HTTP) transformXML(value []byte) []byte { - // only do a simple check, as some devices e.g. Kostal Piko MP plus don't seem to send proper XML - if !bytes.HasPrefix(value, []byte("<")) { - return value - } - - xmlReader := bytes.NewReader(value) - - // Decode XML document - root := new(xj.Node) - if err := xj.NewDecoder(xmlReader).DecodeWithCustomPrefixes(root, "", "attr"); err != nil { - return value - } - - // Then encode it in JSON - json := new(bytes.Buffer) - if err := xj.NewEncoder(json).Encode(root); err != nil { - return value - } - - return json.Bytes() -} - -func (p *HTTP) unpackValue(value []byte) (string, error) { - switch p.unpack { - case "hex": - b, err := hex.DecodeString(string(value)) - if err != nil { - return "", err - } - return string(b), nil - } - - return "", fmt.Errorf("invalid unpack: %s", p.unpack) -} - -// decode a hex string to a proper value -// TODO reuse similar code from Modbus -func (p *HTTP) decodeValue(value []byte) (interface{}, error) { - switch p.decode { - case "float32", "ieee754": - return rs485.RTUIeee754ToFloat64(value), nil - case "float32s", "ieee754s": - return rs485.RTUIeee754ToFloat64Swapped(value), nil - case "float64": - return rs485.RTUUint64ToFloat64(value), nil - case "uint16": - return rs485.RTUUint16ToFloat64(value), nil - case "uint32": - return rs485.RTUUint32ToFloat64(value), nil - case "uint32s": - return rs485.RTUUint32ToFloat64Swapped(value), nil - case "uint64": - return rs485.RTUUint64ToFloat64(value), nil - case "int16": - return rs485.RTUInt16ToFloat64(value), nil - case "int32": - return rs485.RTUInt32ToFloat64(value), nil - case "int32s": - return rs485.RTUInt32ToFloat64Swapped(value), nil - } - - return nil, fmt.Errorf("invalid decoding: %s", p.decode) -} - // FloatGetter parses float from request func (p *HTTP) FloatGetter() func() (float64, error) { g := p.StringGetter() @@ -336,55 +192,9 @@ func (p *HTTP) IntGetter() func() (int64, error) { func (p *HTTP) StringGetter() func() (string, error) { return func() (string, error) { b, err := p.request(p.body) - if err != nil { - return string(b), err - } - b = p.transformXML(b) - - if p.re != nil { - m := p.re.FindSubmatch(b) - if len(m) > 1 { - b = m[1] // first submatch - } - } - - if p.jq != nil { - v, err := jq.Query(p.jq, b) - if err != nil { - return string(b), err - } - b = []byte(fmt.Sprintf("%v", v)) - } - - if p.unpack != "" { - v, err := p.unpackValue(b) - if err != nil { - return string(b), err - } - b = []byte(fmt.Sprintf("%v", v)) - } - - if p.decode != "" { - v, err := p.decodeValue(b) - if err != nil { - return string(b), err - } - b = []byte(fmt.Sprintf("%v", v)) - } - - if p.vm != nil { - err := p.vm.Set("val", string(b)) - if err != nil { - return string(b), err - } - - v, err := p.vm.Eval(p.script) - if err != nil { - return string(b), err - } - - return v.ToString() + if err == nil && p.pipeline != nil { + b, err = p.pipeline.Process(b) } return string(b), err diff --git a/provider/mqtt.go b/provider/mqtt.go index 291b1d6f0..b77d908db 100644 --- a/provider/mqtt.go +++ b/provider/mqtt.go @@ -1,25 +1,22 @@ package provider import ( - "fmt" - "regexp" "time" "github.com/evcc-io/evcc/provider/mqtt" + "github.com/evcc-io/evcc/provider/pipeline" "github.com/evcc-io/evcc/util" - "github.com/itchyny/gojq" ) // Mqtt provider type Mqtt struct { - log *util.Logger - client *mqtt.Client - topic string - payload string - scale float64 - timeout time.Duration - re *regexp.Regexp - jq *gojq.Query + log *util.Logger + client *mqtt.Client + topic string + payload string + scale float64 + timeout time.Duration + pipeline *pipeline.Pipeline } func init() { @@ -29,12 +26,11 @@ func init() { // NewMqttFromConfig creates Mqtt provider func NewMqttFromConfig(other map[string]interface{}) (IntProvider, error) { cc := struct { - mqtt.Config `mapstructure:",squash"` - Topic, Payload string // Payload only applies to setters - Scale float64 - Timeout time.Duration - Regex string - Jq string + mqtt.Config `mapstructure:",squash"` + Topic, Payload string // Payload only applies to setters + Scale float64 + Timeout time.Duration + pipeline.Settings `mapstructure:",squash"` }{ Scale: 1, } @@ -52,19 +48,12 @@ func NewMqttFromConfig(other map[string]interface{}) (IntProvider, error) { m := NewMqtt(log, client, cc.Topic, cc.Scale, cc.Timeout).WithPayload(cc.Payload) - if cc.Regex != "" { - if m, err = m.WithRegex(cc.Regex); err != nil { - return nil, err - } + pipe, err := pipeline.New(cc.Settings) + if err == nil { + m = m.WithPipeline(pipe) } - if cc.Jq != "" { - if m, err = m.WithJq(cc.Jq); err != nil { - return nil, err - } - } - - return m, nil + return m, err } // NewMqtt creates mqtt provider for given topic @@ -86,28 +75,10 @@ func (m *Mqtt) WithPayload(payload string) *Mqtt { return m } -// WithRegex adds a regex query applied to the mqtt listener payload -func (m *Mqtt) WithRegex(regex string) (*Mqtt, error) { - re, err := regexp.Compile(regex) - if err != nil { - return m, fmt.Errorf("invalid regex '%s': %w", re, err) - } - - m.re = re - - return m, nil -} - -// WithJq adds a jq query applied to the mqtt listener payload -func (m *Mqtt) WithJq(jq string) (*Mqtt, error) { - op, err := gojq.Parse(jq) - if err != nil { - return m, fmt.Errorf("invalid jq query '%s': %w", jq, err) - } - - m.jq = op - - return m, nil +// WithPipeline adds a processing pipeline +func (p *Mqtt) WithPipeline(pipeline *pipeline.Pipeline) *Mqtt { + p.pipeline = pipeline + return p } var _ FloatProvider = (*Mqtt)(nil) @@ -117,11 +88,10 @@ var _ FloatProvider = (*Mqtt)(nil) // if initial value is not received within `timeout` or max. 10s if timeout is not given. func (m *Mqtt) newReceiver() *msgHandler { h := &msgHandler{ - topic: m.topic, - scale: m.scale, - mux: util.NewWaiter(m.timeout, func() { m.log.DEBUG.Printf("%s wait for initial value", m.topic) }), - re: m.re, - jq: m.jq, + topic: m.topic, + scale: m.scale, + mux: util.NewWaiter(m.timeout, func() { m.log.DEBUG.Printf("%s wait for initial value", m.topic) }), + pipeline: m.pipeline, } m.client.Listen(m.topic, h.receive) diff --git a/provider/mqtt_handler.go b/provider/mqtt_handler.go index b5f3bb120..4da973893 100644 --- a/provider/mqtt_handler.go +++ b/provider/mqtt_handler.go @@ -3,22 +3,19 @@ package provider import ( "fmt" "math" - "regexp" "strconv" "time" + "github.com/evcc-io/evcc/provider/pipeline" "github.com/evcc-io/evcc/util" - "github.com/evcc-io/evcc/util/jq" - "github.com/itchyny/gojq" ) type msgHandler struct { - mux *util.Waiter - scale float64 - topic string - payload string - re *regexp.Regexp - jq *gojq.Query + mux *util.Waiter + scale float64 + topic string + pipeline *pipeline.Pipeline + payload string } func (h *msgHandler) receive(payload string) { @@ -38,24 +35,12 @@ func (h *msgHandler) hasValue() (string, error) { return "", fmt.Errorf("%s outdated: %v", h.topic, late.Truncate(time.Second)) } - var err error - payload := h.payload - - if h.re != nil { - m := h.re.FindStringSubmatch(payload) - if len(m) > 1 { - payload = m[1] // first submatch - } + if h.pipeline != nil { + b, err := h.pipeline.Process([]byte(h.payload)) + return string(b), err } - if h.jq != nil { - var val interface{} - if val, err = jq.Query(h.jq, []byte(payload)); err == nil { - payload = fmt.Sprintf("%v", val) - } - } - - return payload, err + return h.payload, nil } func (h *msgHandler) floatGetter() (float64, error) { diff --git a/provider/pipeline/pipeline.go b/provider/pipeline/pipeline.go new file mode 100644 index 000000000..8d4756090 --- /dev/null +++ b/provider/pipeline/pipeline.go @@ -0,0 +1,231 @@ +package pipeline + +import ( + "bytes" + "encoding/hex" + "fmt" + "regexp" + "strings" + + xj "github.com/basgys/goxml2json" + "github.com/evcc-io/evcc/provider/javascript" + "github.com/evcc-io/evcc/util/jq" + "github.com/itchyny/gojq" + "github.com/robertkrimen/otto" + "github.com/volkszaehler/mbmd/meters/rs485" +) + +type Pipeline struct { + re *regexp.Regexp + jq *gojq.Query + unpack string + decode string + vm *otto.Otto + script string +} + +type Settings struct { + Regex string + Jq string + Unpack string + Decode string + VM string + Script string +} + +func New(cc Settings) (*Pipeline, error) { + p := new(Pipeline) + + var err error + if err == nil && cc.Regex != "" { + _, err = p.WithRegex(cc.Regex) + } + + if err == nil && cc.Jq != "" { + _, err = p.WithJq(cc.Jq) + } + + if err == nil && cc.Unpack != "" { + _, err = p.WithUnpack(cc.Unpack) + } + + if err == nil && cc.Decode != "" { + _, err = p.WithDecode(cc.Decode) + } + + if err == nil && cc.Script != "" { + _, err = p.WithScript(cc.VM, cc.Script) + } + + return p, err +} + +// WithRegex adds a regex query applied to the mqtt listener payload +func (p *Pipeline) WithRegex(regex string) (*Pipeline, error) { + re, err := regexp.Compile(regex) + if err != nil { + return nil, fmt.Errorf("invalid regex '%s': %w", re, err) + } + + p.re = re + + return p, nil +} + +// WithJq adds a jq query applied to the mqtt listener payload +func (p *Pipeline) WithJq(jq string) (*Pipeline, error) { + op, err := gojq.Parse(jq) + if err != nil { + return nil, fmt.Errorf("invalid jq query '%s': %w", jq, err) + } + + p.jq = op + + return p, nil +} + +// WithUnpack adds data unpacking +func (p *Pipeline) WithUnpack(unpack string) (*Pipeline, error) { + p.unpack = strings.ToLower(unpack) + + return p, nil +} + +// WithDecode adds data decoding +func (p *Pipeline) WithDecode(decode string) (*Pipeline, error) { + p.decode = strings.ToLower(decode) + + return p, nil +} + +// WithScript adds a javascript script to process the response +func (p *Pipeline) WithScript(vm, script string) (*Pipeline, error) { + regvm := javascript.RegisteredVM(strings.ToLower(vm)) + + p.vm = regvm + p.script = script + + return p, nil +} + +// transform XML into JSON with attribute names getting 'attr' prefix +func (p *Pipeline) transformXML(value []byte) []byte { + value = bytes.TrimSpace(value) + + // only do a simple check, as some devices e.g. Kostal Piko MP plus don't seem to send proper XML + if !bytes.HasPrefix(value, []byte("<")) { + return value + } + + in := bytes.NewReader(value) + + // Decode XML document + root := new(xj.Node) + if err := xj.NewDecoder(in).DecodeWithCustomPrefixes(root, "", "attr"); err != nil { + return value + } + + // Then encode it in JSON + out := new(bytes.Buffer) + if err := xj.NewEncoder(out).Encode(root); err != nil { + return value + } + + return out.Bytes() +} + +func (p *Pipeline) unpackValue(value []byte) (string, error) { + switch p.unpack { + case "hex": + b, err := hex.DecodeString(string(value)) + if err != nil { + return "", err + } + return string(b), nil + } + + return "", fmt.Errorf("invalid unpack: %s", p.unpack) +} + +// decode a hex string to a proper value +// TODO reuse similar code from Modbus +func (p *Pipeline) decodeValue(value []byte) (interface{}, error) { + switch p.decode { + case "float32", "ieee754": + return rs485.RTUIeee754ToFloat64(value), nil + case "float32s", "ieee754s": + return rs485.RTUIeee754ToFloat64Swapped(value), nil + case "float64": + return rs485.RTUUint64ToFloat64(value), nil + case "uint16": + return rs485.RTUUint16ToFloat64(value), nil + case "uint32": + return rs485.RTUUint32ToFloat64(value), nil + case "uint32s": + return rs485.RTUUint32ToFloat64Swapped(value), nil + case "uint64": + return rs485.RTUUint64ToFloat64(value), nil + case "int16": + return rs485.RTUInt16ToFloat64(value), nil + case "int32": + return rs485.RTUInt32ToFloat64(value), nil + case "int32s": + return rs485.RTUInt32ToFloat64Swapped(value), nil + } + + return nil, fmt.Errorf("invalid decoding: %s", p.decode) +} + +func (p *Pipeline) Process(in []byte) ([]byte, error) { + b := p.transformXML(in) + + if p.re != nil { + m := p.re.FindSubmatch(b) + if len(m) > 1 { + b = m[1] // first submatch + } + } + + if p.jq != nil { + v, err := jq.Query(p.jq, b) + if err != nil { + return b, err + } + b = []byte(fmt.Sprintf("%v", v)) + } + + if p.unpack != "" { + v, err := p.unpackValue(b) + if err != nil { + return b, err + } + b = []byte(fmt.Sprintf("%v", v)) + } + + if p.decode != "" { + v, err := p.decodeValue(b) + if err != nil { + return b, err + } + b = []byte(fmt.Sprintf("%v", v)) + } + + if p.vm != nil { + if err := p.vm.Set("val", string(b)); err != nil { + return b, err + } + + v, err := p.vm.Eval(p.script) + if err != nil { + return b, err + } + + s, err := v.ToString() + b = []byte(s) + if err != nil { + return b, err + } + } + + return b, nil +}