From 56ad3dad952d78d561ff37ed1e641bf46929c0ef Mon Sep 17 00:00:00 2001 From: andig Date: Mon, 27 Apr 2020 18:40:22 +0200 Subject: [PATCH] Opinionated refactor of SMA Energy Meter (#65) --- evcc.dist.yaml | 12 ++-- meter/config.go | 2 +- meter/sma.go | 48 +++++++--------- meter/sma/listener.go | 112 ++++++++++++++++--------------------- meter/sma/listener_test.go | 21 ++++--- 5 files changed, 87 insertions(+), 108 deletions(-) diff --git a/evcc.dist.yaml b/evcc.dist.yaml index 50a83c629..e11011a77 100644 --- a/evcc.dist.yaml +++ b/evcc.dist.yaml @@ -61,12 +61,12 @@ meters: type: script # use script cmd: /bin/sh -c "echo 0" # actual command timeout: 3s # kill script after 3 seconds -#- name: grid -# type: smameter # SMA Home Manager 2.0 or SMA Energy Meter 30 -# uri: 192.168.1.4 # IP Address of the device -#- name: pv -# type: smameter # SMA Home Manager 2.0 or SMA Energy Meter 30 -# uri: 192.168.1.4 # IP Address of the device +- name: sma-grid + type: sma # SMA Home Manager 2.0 or SMA Energy Meter 30 + uri: 192.168.1.4 # IP Address of the device +- name: sma-pv + type: sma # SMA Home Manager 2.0 or SMA Energy Meter 30 + uri: 192.168.1.4 # IP Address of the device chargers: - name: wallbe diff --git a/meter/config.go b/meter/config.go index 26c4a1d0a..568d9ae1d 100644 --- a/meter/config.go +++ b/meter/config.go @@ -13,7 +13,7 @@ func NewFromConfig(log *api.Logger, typ string, other map[string]interface{}) ap switch strings.ToLower(typ) { case "default", "configurable": c = NewConfigurableFromConfig(log, other) - case "smameter": + case "sma": c = NewSMAFromConfig(log, other) default: log.FATAL.Fatalf("invalid meter type '%s'", typ) diff --git a/meter/sma.go b/meter/sma.go index 1dcbec792..7d9002302 100644 --- a/meter/sma.go +++ b/meter/sma.go @@ -16,13 +16,13 @@ const ( // SMA supporting SMA Home Manager 2.0 and SMA Energy Meter 30 type SMA struct { - log *api.Logger - uri string - power float64 - lastUpdate time.Time - recv chan sma.TelegramData - mux sync.Mutex - once sync.Once + log *api.Logger + uri string + power float64 + updated time.Time + recv chan sma.Telegram + mux sync.Mutex + once sync.Once } // NewSMAFromConfig creates a SMA Meter from generic config @@ -42,7 +42,7 @@ func NewSMA(uri string) *SMA { sm := &SMA{ log: log, uri: uri, - recv: make(chan sma.TelegramData), + recv: make(chan sma.Telegram), } if sma.Instance == nil { @@ -61,11 +61,11 @@ func (sm *SMA) waitForInitialValue() { sm.mux.Lock() defer sm.mux.Unlock() - if sm.lastUpdate.IsZero() { + if sm.updated.IsZero() { sm.log.TRACE.Print("waiting for initial value") // wait for initial update - for sm.lastUpdate.IsZero() { + for sm.updated.IsZero() { sm.mux.Unlock() time.Sleep(waitTimeout) sm.mux.Lock() @@ -76,28 +76,22 @@ func (sm *SMA) waitForInitialValue() { // receive processes the channel message containing the multicast data func (sm *SMA) receive() { for msg := range sm.recv { - if msg.Data == nil { - continue - } - - var powerIn, powerOut float64 - var ok bool - - if powerIn, ok = msg.Data[sma.ObisImportPower]; !ok { - continue - } - - if powerOut, ok = msg.Data[sma.ObisExportPower]; !ok { + if msg.Values == nil { continue } sm.mux.Lock() - sm.lastUpdate = time.Now() - if powerOut > 0 { - sm.power = -powerOut + + if power, ok := msg.Values[sma.ObisExportPower]; ok { + sm.power = -power + sm.updated = time.Now() + } else if power, ok := msg.Values[sma.ObisImportPower]; ok { + sm.power = power + sm.updated = time.Now() } else { - sm.power = powerIn + sm.log.WARN.Println("missing obis for import/export power") } + sm.mux.Unlock() } } @@ -108,7 +102,7 @@ func (sm *SMA) CurrentPower() (float64, error) { sm.mux.Lock() defer sm.mux.Unlock() - if time.Since(sm.lastUpdate) > udpTimeout { + if time.Since(sm.updated) > udpTimeout { return 0, errors.New("recv timeout") } diff --git a/meter/sma/listener.go b/meter/sma/listener.go index 6c668ce4d..a3bf6e563 100644 --- a/meter/sma/listener.go +++ b/meter/sma/listener.go @@ -12,24 +12,28 @@ import ( ) const ( - multicastAddr = "239.12.255.254:9522" - udpBufferSize = 8192 - obisCodeLength = 4 + multicastAddr = "239.12.255.254:9522" + udpBufferSize = 8192 + + msgSerial = 20 // start of serial in preamble + msgPreamble = 28 // preamble size in bytes + msgCodeLength = 4 // length in bytes + ObisImportPower = "1:1.4.0" // Wirkleistung (W) ObisImportEnergy = "1:1.8.0" // Wirkarbeit (Ws) + ObisExportPower = "1:2.4.0" // Wirkleistung (W) ObisExportEnergy = "1:2.8.0" // Wirkarbeit (Ws) − ) -// obisCodeProp defines the properties needed to parse the SMA multicast telegram values -type obisCodeProp struct { +// obisDefinition defines the properties needed to parse the SMA multicast telegram values +type obisDefinition struct { length int // data size in bytes of the return value factor float64 // the factor to multiply the value by to get the proper value in the given unit } // list of Obis codes and their properties as defined in the SMA EMETER-Protokoll-TI-de-10.pdf document -var knownObisCodes = map[string]obisCodeProp{ - // Overal sums +var knownObisCodes = map[string]obisDefinition{ + // Overall sums ObisImportPower: {4, 0.1}, ObisImportEnergy: {8, 1}, // Wirkleistung (W)/-arbeit (Ws) + ObisExportPower: {4, 0.1}, ObisExportEnergy: {8, 1}, // Wirkleistung (W)/-arbeit (Ws) − "1:3.4.0": {4, 0.1}, "1:3.8.0": {8, 1}, // Blindleistung (W)/-arbeit (Ws) + @@ -68,13 +72,14 @@ var knownObisCodes = map[string]obisCodeProp{ "144:0.0.0": {4, 1}, // SW Version } +// Instance is the Listener singleton var Instance *Listener -// TelegramData defines the data structure of a SMA multicast data package -type TelegramData struct { +// Telegram defines the data structure of a SMA multicast data package +type Telegram struct { Addr string Serial string - Data map[string]float64 + Values map[string]float64 } // Listener for receiving SMA multicast data packages @@ -82,19 +87,19 @@ type Listener struct { mux sync.Mutex log *api.Logger conn *net.UDPConn - clients map[string]chan<- TelegramData + clients map[string]chan<- Telegram } // New creates a Listener func New(log *api.Logger, addr string) *Listener { // Parse the string address - laddr, err := net.ResolveUDPAddr("udp4", multicastAddr) + gaddr, err := net.ResolveUDPAddr("udp4", multicastAddr) if err != nil { log.FATAL.Fatalf("error resolving udp address: %s", err) } // Open up a connection - conn, err := net.ListenMulticastUDP("udp4", nil, laddr) + conn, err := net.ListenMulticastUDP("udp4", nil, gaddr) if err != nil { log.FATAL.Fatalf("error opening connecting: %s", err) } @@ -113,95 +118,76 @@ func New(log *api.Logger, addr string) *Listener { return l } -// processUDPData converts a SMA Multicast data package into TelegramData -func (l *Listener) processUDPData(src *net.UDPAddr, buffer []byte) (TelegramData, error) { - numBytes := len(buffer) +// processMessage converts a SMA multicast data package into Telegram +func (l *Listener) processMessage(src *net.UDPAddr, b []byte) (Telegram, error) { + numBytes := len(b) - if numBytes < 29 { - return TelegramData{}, errors.New("received data package is too small") + if numBytes <= msgPreamble { + return Telegram{}, errors.New("received data package is too small") } - obisCodeValues := make(map[string]float64) + obisValues := make(map[string]float64) - // yes this doesn't look nice, but keeping it until we found a better way to parse the data - // read obis code values, start at position 28, after initial static stuff - for i := 28; i < numBytes; i++ { - if i+obisCodeLength > numBytes-1 { - break + var obisDef obisDefinition + for i := msgPreamble; i < numBytes-msgCodeLength; i += msgCodeLength + obisDef.length { + // spec says value should be 1, but reading contains 0 + b0 := b[i+0] + if b0 == 0 { + b0 = 1 } - // create the string notation of the potential obis code - b := buffer[i : i+obisCodeLength] - - // Spec says value should be 1, but reading contains 0 - b3 := b[0] - if b3 == 0 { - b3 = 1 + code := fmt.Sprintf("%d:%d.%d.%d", b0, b[i+1], b[i+2], b[i+3]) + if obisDef, ok := knownObisCodes[code]; ok { + switch obisDef.length { + case 4: + obisValues[code] = obisDef.factor * float64(binary.BigEndian.Uint32(b[i+msgCodeLength:])) + case 8: + obisValues[code] = obisDef.factor * float64(binary.BigEndian.Uint64(b[i+msgCodeLength:])) + } } - - code := fmt.Sprintf("%d:%d.%d.%d", b3, b[1], b[2], b[3]) - - element, ok := knownObisCodes[code] - if !ok { - continue - } - - dataIndex := i + obisCodeLength - - var value float64 - switch element.length { - case 4: - value = float64(binary.BigEndian.Uint32(buffer[dataIndex : dataIndex+element.length+1])) - case 8: - value = float64(binary.BigEndian.Uint64(buffer[dataIndex : dataIndex+element.length+1])) - } - - obisCodeValues[code] = value * element.factor - - i = dataIndex + element.length - 1 } - serial := strconv.FormatUint(uint64(binary.BigEndian.Uint32(buffer[20:24])), 10) + serial := strconv.FormatUint(uint64(binary.BigEndian.Uint32(b[msgSerial:])), 10) - msg := TelegramData{ + msg := Telegram{ Addr: src.IP.String(), Serial: serial, - Data: obisCodeValues, + Values: obisValues, } return msg, nil } -// listen for Multicast data packages +// listen for multicast data packages func (l *Listener) listen() { buffer := make([]byte, udpBufferSize) - // Loop forever reading + for { - numBytes, src, err := l.conn.ReadFromUDP(buffer) + read, src, err := l.conn.ReadFromUDP(buffer) if err != nil { - l.log.WARN.Printf("readfromudp failed: %s", err) + l.log.WARN.Printf("udp read failed: %s", err) continue } - if msg, err := l.processUDPData(src, buffer[:numBytes-1]); err == nil { + if msg, err := l.processMessage(src, buffer[:read-1]); err == nil { l.send(msg) } } } // Subscribe adds a client address and message channel -func (l *Listener) Subscribe(addr string, c chan<- TelegramData) { +func (l *Listener) Subscribe(addr string, c chan<- Telegram) { l.mux.Lock() defer l.mux.Unlock() if l.clients == nil { - l.clients = make(map[string]chan<- TelegramData) + l.clients = make(map[string]chan<- Telegram) } l.clients[addr] = c } -func (l *Listener) send(msg TelegramData) { +func (l *Listener) send(msg Telegram) { l.mux.Lock() defer l.mux.Unlock() diff --git a/meter/sma/listener_test.go b/meter/sma/listener_test.go index efc05133e..1d92af471 100644 --- a/meter/sma/listener_test.go +++ b/meter/sma/listener_test.go @@ -6,14 +6,14 @@ import ( "testing" ) -func TestListenerProcessUDPData(t *testing.T) { +func TestListenerProcessMessage(t *testing.T) { tests := []struct { name string ip net.IP port int response []byte wantErr bool - want TelegramData + want Telegram }{ { "SMA Home Manager - success", @@ -60,10 +60,10 @@ func TestListenerProcessUDPData(t *testing.T) { 0x00, 0x00, 0x00, 0x1e, 0x90, 0x00, 0x00, 0x00, 0x02, 0x03, 0x05, 0x52, 0x00, 0x00, 0x00, 0x00, }, false, - TelegramData{ + Telegram{ Addr: "192.168.1.4", Serial: "0", - Data: map[string]float64{ + Values: map[string]float64{ "1:1.4.0": 0, "1:1.8.0": 6.89131908e+09, "1:2.4.0": 37.9, "1:2.8.0": 2.371258944e+10, "1:3.4.0": 0, "1:3.8.0": 6.30376488e+09, @@ -144,10 +144,10 @@ func TestListenerProcessUDPData(t *testing.T) { 0x02, 0x00, 0x12, 0x52, 0x00, 0x00, 0x00, 0x00, }, false, - TelegramData{ + Telegram{ Addr: "192.168.1.4", Serial: "0", - Data: map[string]float64{ + Values: map[string]float64{ "1:1.4.0": 0, "1:1.8.0": 1.2385008e+08, "1:2.4.0": 222, "1:2.8.0": 4.137912288e+10, "1:3.4.0": 0, "1:3.8.0": 1.66667796e+09, @@ -189,18 +189,17 @@ func TestListenerProcessUDPData(t *testing.T) { l := &Listener{} buffer := tc.response - numBytes := len(buffer) + read := len(buffer) src := &net.UDPAddr{IP: tc.ip, Port: tc.port} - got, err := l.processUDPData(src, buffer[:numBytes-1]) + got, err := l.processMessage(src, buffer[:read-1]) if (err != nil) != tc.wantErr { - t.Errorf("Listener.processUDPData() error = %v, wantErr %v", err, tc.wantErr) + t.Errorf("Listener.processMessage() error = %v, wantErr %v", err, tc.wantErr) return } if !reflect.DeepEqual(got, tc.want) { - t.Errorf("Listener.processUDPData() = %v, want %v", got, tc.want) + t.Errorf("Listener.processMessage() got %v, want %v", got, tc.want) } - }) } }