Refactor data processing into pipeline component (#2169)
This commit is contained in:
parent
bac4e07980
commit
d97e8d18e5
5 changed files with 297 additions and 297 deletions
|
|
@ -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())
|
||||
}
|
||||
|
||||
|
|
|
|||
238
provider/http.go
238
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
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
|
|
|||
|
|
@ -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) {
|
||||
|
|
|
|||
231
provider/pipeline/pipeline.go
Normal file
231
provider/pipeline/pipeline.go
Normal file
|
|
@ -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
|
||||
}
|
||||
Loading…
Add table
Add a link
Reference in a new issue