Modbus: add UDP and allow concurrent access (#13676)

This commit is contained in:
andig 2024-08-02 21:32:45 +02:00 • committed by GitHub
parent 67d98707f8
commit 03f62e5593
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
32 changed files with 329 additions and 305 deletions

View file

@ -65,7 +65,7 @@ func NewABBFromConfig(other map[string]interface{}) (api.Charger, error) {
return nil, err
}
return NewABB(cc.URI, cc.Device, cc.Comset, cc.Baudrate, modbus.ProtocolFromRTU(cc.RTU), cc.ID)
return NewABB(cc.URI, cc.Device, cc.Comset, cc.Baudrate, cc.Protocol(), cc.ID)
}
// NewABB creates ABB charger

View file

@ -55,7 +55,7 @@ func NewAlphatecFromConfig(other map[string]interface{}) (api.Charger, error) {
return nil, err
}
return NewAlphatec(cc.URI, cc.Device, cc.Comset, cc.Baudrate, modbus.ProtocolFromRTU(cc.RTU), cc.ID)
return NewAlphatec(cc.URI, cc.Device, cc.Comset, cc.Baudrate, cc.Protocol(), cc.ID)
}
// NewAlphatec creates Alphatec charger

View file

@ -76,7 +76,7 @@ func NewDeltaFromConfig(other map[string]interface{}) (api.Charger, error) {
return nil, err
}
return NewDelta(cc.URI, cc.Device, cc.Comset, cc.Baudrate, modbus.ProtocolFromRTU(cc.RTU), cc.ID, cc.Connector)
return NewDelta(cc.URI, cc.Device, cc.Comset, cc.Baudrate, cc.Settings.Protocol(), cc.ID, cc.Connector)
}
// NewDelta creates Delta charger

View file

@ -39,7 +39,7 @@ func NewEvseDINFromConfig(other map[string]interface{}) (api.Charger, error) {
return nil, err
}
return NewEvseDIN(cc.URI, cc.Device, cc.Comset, cc.Baudrate, modbus.ProtocolFromRTU(cc.RTU), cc.ID)
return NewEvseDIN(cc.URI, cc.Device, cc.Comset, cc.Baudrate, cc.Protocol(), cc.ID)
}
// NewEvseDIN creates EVSE DIN charger

View file

@ -70,7 +70,7 @@ func NewHeidelbergECFromConfig(other map[string]interface{}) (api.Charger, error
return nil, err
}
return NewHeidelbergEC(cc.URI, cc.Device, cc.Comset, cc.Baudrate, modbus.ProtocolFromRTU(cc.RTU), cc.ID)
return NewHeidelbergEC(cc.URI, cc.Device, cc.Comset, cc.Baudrate, cc.Protocol(), cc.ID)
}
// NewHeidelbergEC creates HeidelbergEC charger

View file

@ -83,7 +83,7 @@ func NewMennekesCompactFromConfig(other map[string]interface{}) (api.Charger, er
return nil, err
}
return NewMennekesCompact(cc.URI, cc.Device, cc.Comset, cc.Baudrate, modbus.ProtocolFromRTU(cc.RTU), cc.ID, cc.Timeout)
return NewMennekesCompact(cc.URI, cc.Device, cc.Comset, cc.Baudrate, cc.Settings.Protocol(), cc.ID, cc.Timeout)
}
// NewMennekesCompact creates Mennekes charger

View file

@ -39,7 +39,7 @@ func NewOboFromConfig(other map[string]interface{}) (api.Charger, error) {
return nil, err
}
return NewObo(cc.URI, cc.Device, cc.Comset, cc.Baudrate, modbus.ProtocolFromRTU(cc.RTU), cc.ID)
return NewObo(cc.URI, cc.Device, cc.Comset, cc.Baudrate, cc.Protocol(), cc.ID)
}
// NewObo creates OBO Bettermann charger

View file

@ -34,7 +34,7 @@ func NewPhoenixEVSerFromConfig(other map[string]interface{}) (api.Charger, error
return nil, err
}
return NewPhoenixEVSer(cc.URI, cc.Device, cc.Comset, cc.Baudrate, modbus.ProtocolFromRTU(cc.RTU), cc.ID)
return NewPhoenixEVSer(cc.URI, cc.Device, cc.Comset, cc.Baudrate, cc.Protocol(), cc.ID)
}
// NewPhoenixEVSer creates a Phoenix charger

View file

@ -66,7 +66,7 @@ func NewPrachtAlphaFromConfig(other map[string]interface{}) (api.Charger, error)
return nil, err
}
return NewPrachtAlpha(cc.URI, cc.Device, cc.Comset, cc.Baudrate, modbus.ProtocolFromRTU(cc.RTU), cc.ID, cc.Timeout, cc.Connector)
return NewPrachtAlpha(cc.URI, cc.Device, cc.Comset, cc.Baudrate, cc.Settings.Protocol(), cc.ID, cc.Timeout, cc.Connector)
}
// NewPrachtAlpha creates PrachtAlpha charger

View file

@ -62,7 +62,7 @@ func NewPulsaresFromConfig(other map[string]interface{}) (api.Charger, error) {
return nil, err
}
wb, err := NewPulsares(cc.URI, cc.Device, cc.Comset, cc.Baudrate, modbus.ProtocolFromRTU(cc.RTU), cc.ID)
wb, err := NewPulsares(cc.URI, cc.Device, cc.Comset, cc.Baudrate, cc.Protocol(), cc.ID)
if err != nil {
return nil, err
}

View file

@ -77,7 +77,7 @@ func NewsmartEVSEFromConfig(other map[string]interface{}) (api.Charger, error) {
return nil, err
}
return NewsmartEVSE(cc.URI, cc.Device, cc.Comset, cc.Baudrate, modbus.ProtocolFromRTU(cc.RTU), cc.ID)
return NewsmartEVSE(cc.URI, cc.Device, cc.Comset, cc.Baudrate, cc.Protocol(), cc.ID)
}
// NewsmartEVSE creates a new charger

View file

@ -70,7 +70,7 @@ func NewSolaxFromConfig(other map[string]interface{}) (api.Charger, error) {
return nil, err
}
return NewSolax(cc.URI, cc.Device, cc.Comset, cc.Baudrate, modbus.ProtocolFromRTU(cc.RTU), cc.ID)
return NewSolax(cc.URI, cc.Device, cc.Comset, cc.Baudrate, cc.Protocol(), cc.ID)
}
// NewSolax creates Solax charger

View file

@ -75,7 +75,7 @@ func NewSungrowFromConfig(other map[string]interface{}) (api.Charger, error) {
return nil, err
}
return NewSungrow(cc.URI, cc.Device, cc.Comset, cc.Baudrate, modbus.ProtocolFromRTU(cc.RTU), cc.ID)
return NewSungrow(cc.URI, cc.Device, cc.Comset, cc.Baudrate, cc.Protocol(), cc.ID)
}
// NewSungrow creates Sungrow charger

View file

@ -204,6 +204,8 @@ NEXT:
}
func configureMeters(static []config.Named, names ...string) error {
g, _ := errgroup.WithContext(context.Background())
for i, cc := range static {
if cc.Name == "" {
return fmt.Errorf("cannot create meter %d: missing name", i+1)
@ -217,14 +219,18 @@ func configureMeters(static []config.Named, names ...string) error {
log.WARN.Printf("create meter %d: %v", i+1, err)
}
instance, err := meter.NewFromConfig(cc.Type, cc.Other)
if err != nil {
return &DeviceError{cc.Name, fmt.Errorf("cannot create meter '%s': %w", cc.Name, err)}
}
g.Go(func() error {
instance, err := meter.NewFromConfig(cc.Type, cc.Other)
if err != nil {
return &DeviceError{cc.Name, fmt.Errorf("cannot create meter '%s': %w", cc.Name, err)}
}
if err := config.Meters().Add(config.NewStaticDevice(cc, instance)); err != nil {
return &DeviceError{cc.Name, err}
}
if err := config.Meters().Add(config.NewStaticDevice(cc, instance)); err != nil {
return &DeviceError{cc.Name, err}
}
return nil
})
}
// append devices from database
@ -234,25 +240,29 @@ func configureMeters(static []config.Named, names ...string) error {
}
for _, conf := range configurable {
cc := conf.Named()
g.Go(func() error {
cc := conf.Named()
if len(names) > 0 && !slices.Contains(names, cc.Name) {
return nil
}
// TODO add fake devices
instance, err := meter.NewFromConfig(cc.Type, cc.Other)
if err != nil {
return &DeviceError{cc.Name, fmt.Errorf("cannot create meter '%s': %w", cc.Name, err)}
}
if err := config.Meters().Add(config.NewConfigurableDevice(conf, instance)); err != nil {
return &DeviceError{cc.Name, err}
}
if len(names) > 0 && !slices.Contains(names, cc.Name) {
return nil
}
// TOTO add fake devices
instance, err := meter.NewFromConfig(cc.Type, cc.Other)
if err != nil {
return &DeviceError{cc.Name, fmt.Errorf("cannot create meter '%s': %w", cc.Name, err)}
}
if err := config.Meters().Add(config.NewConfigurableDevice(conf, instance)); err != nil {
return &DeviceError{cc.Name, err}
}
})
}
return nil
return g.Wait()
}
func configureChargers(static []config.Named, names ...string) error {
@ -299,7 +309,7 @@ func configureChargers(static []config.Named, names ...string) error {
return nil
}
// TOTO add fake devices
// TODO add fake devices
instance, err := charger.NewFromConfig(cc.Type, cc.Other)
if err != nil {

4
go.mod
View file

@ -89,7 +89,7 @@ require (
github.com/teslamotors/vehicle-command v0.0.2
github.com/traefik/yaegi v0.16.1
github.com/tv42/httpunix v0.0.0-20191220191345-2ba4b9c3382c
github.com/volkszaehler/mbmd v0.0.0-20240611142726-33463eb0324e
github.com/volkszaehler/mbmd v0.0.0-20240727104742-3191c0dbfb9e
github.com/writeas/go-strip-markdown/v2 v2.1.1
gitlab.com/bboehmke/sunny v0.16.0
go.uber.org/mock v0.4.0
@ -197,6 +197,8 @@ require (
replace gopkg.in/yaml.v3 => github.com/andig/yaml v0.0.0-20240531135838-1ff5761ab467
replace github.com/grid-x/modbus => github.com/evcc-io/modbus v0.0.0-20240503125516-9fd99fe0e438
replace github.com/enbility/spine-go => github.com/enbility/spine-go v0.0.0-20240726200332-a983de1e34b8
replace github.com/enbility/ship-go => github.com/enbility/ship-go v0.0.0-20240731093131-37b1302bca66

10
go.sum
View file

@ -132,6 +132,8 @@ github.com/enbility/zeroconf/v2 v2.0.0-20240210101930-d0004078577b/go.mod h1:Bjz
github.com/envoyproxy/go-control-plane v0.6.9/go.mod h1:SBwIajubJHhxtWwsL9s8ss4safvEdbitLhGGK48rN6g=
github.com/envoyproxy/go-control-plane v0.9.1-0.20191026205805-5f8ba28d4473/go.mod h1:YTl/9mNaCwkRvm6d1a2C3ymFceY/DCBVvsKhRF0iEA4=
github.com/envoyproxy/protoc-gen-validate v0.1.0/go.mod h1:iSmxcyjqTsJpI2R4NaDN7+kN2VEUnK/pcBlmesArF7c=
github.com/evcc-io/modbus v0.0.0-20240503125516-9fd99fe0e438 h1:I14MN2UauarS2Lo+prnNXEj31MxJS1mEfLUrEnOUBsw=
github.com/evcc-io/modbus v0.0.0-20240503125516-9fd99fe0e438/go.mod h1:WpbUAyptAAi0VAriSRopZa6uhiJOJCTz7KFvgGtNRXc=
github.com/evcc-io/ocpp-go v0.0.0-20240730071053-d69e53b0fce9 h1:FLv1vmLnfc8DanI5U1qOTe1Zr0OgZ/tuOvHALtZ2sOI=
github.com/evcc-io/ocpp-go v0.0.0-20240730071053-d69e53b0fce9/go.mod h1:ZynYDWGw6CslG3vyPuucLsy6AyE+h3XXYlr39jhNiQY=
github.com/evcc-io/tesla-proxy-client v0.0.0-20240221194046-4168b3759701 h1:3JplY3KS6KMDVDNAU+3+KWmSWmoHIU34qwuIpW6SiHk=
@ -276,10 +278,6 @@ github.com/gregdel/pushover v1.3.1 h1:4bMLITOZ15+Zpi6qqoGqOPuVHCwSUvMCgVnN5Xhilf
github.com/gregdel/pushover v1.3.1/go.mod h1:EcaO66Nn1StkpEm1iKtBTV3d2A16SoMsVER1PthX7to=
github.com/gregjones/httpcache v0.0.0-20190611155906-901d90724c79 h1:+ngKgrYPPJrOjhax5N+uePQ0Fh1Z7PheYoUI/0nzkPA=
github.com/gregjones/httpcache v0.0.0-20190611155906-901d90724c79/go.mod h1:FecbI9+v66THATjSRHfNgh1IVFe/9kFxbXtjV0ctIMA=
github.com/grid-x/modbus v0.0.0-20210714071042-7af2b65ec03b/go.mod h1:YaK0rKJenZ74vZFcSSLlAQqtG74PMI68eDjpDCDDmTw=
github.com/grid-x/modbus v0.0.0-20240503115206-582f2ab60a18 h1:8V5xRtdD70kGC4/IHqFq+kcBSWr4k6nscAUgWwJ6A5k=
github.com/grid-x/modbus v0.0.0-20240503115206-582f2ab60a18/go.mod h1:WpbUAyptAAi0VAriSRopZa6uhiJOJCTz7KFvgGtNRXc=
github.com/grid-x/serial v0.0.0-20191104121038-e24bc9bf6f08/go.mod h1:kdOd86/VGFWRrtkNwf1MPk0u1gIjc4Y7R2j7nhwc7Rk=
github.com/grid-x/serial v0.0.0-20211107191517-583c7356b3aa h1:Rsn6ARgNkXrsXJIzhkE4vQr5Gbx2LvtEMv4BJOK4LyU=
github.com/grid-x/serial v0.0.0-20211107191517-583c7356b3aa/go.mod h1:kdOd86/VGFWRrtkNwf1MPk0u1gIjc4Y7R2j7nhwc7Rk=
github.com/grpc-ecosystem/go-grpc-middleware v1.0.1-0.20190118093823-f849b5445de4/go.mod h1:FiyG127CGDf3tlThmgyCl78X/SZQqEOJBCDaAfeWzPs=
@ -662,8 +660,8 @@ github.com/vmihailenco/msgpack/v5 v5.4.1 h1:cQriyiUvjTwOHg8QZaPihLWeRAAVoCpE00IU
github.com/vmihailenco/msgpack/v5 v5.4.1/go.mod h1:GaZTsDaehaPpQVyxrf5mtQlH+pc21PIudVV/E3rRQok=
github.com/vmihailenco/tagparser/v2 v2.0.0 h1:y09buUbR+b5aycVFQs/g70pqKVZNBmxwAhO7/IwNM9g=
github.com/vmihailenco/tagparser/v2 v2.0.0/go.mod h1:Wri+At7QHww0WTrCBeu4J6bNtoV6mEfg5OIWRZA9qds=
github.com/volkszaehler/mbmd v0.0.0-20240611142726-33463eb0324e h1:1iAo0qalenDcyo9ySgRWjDuMTpfyA/rBDKklMMNahZk=
github.com/volkszaehler/mbmd v0.0.0-20240611142726-33463eb0324e/go.mod h1:p1nUKfszvbZ+mMMsMi1vaKGA/eBAXC34VSKCk5slyG4=
github.com/volkszaehler/mbmd v0.0.0-20240727104742-3191c0dbfb9e h1:d0KoYuIPXTI/vn85axsoD7w616oUalAzmDLZwv9yB9Q=
github.com/volkszaehler/mbmd v0.0.0-20240727104742-3191c0dbfb9e/go.mod h1:MA3vKWI4KozcpWvD9dwN6CL+WqwdK9XGXT5Iv49IsUM=
github.com/writeas/go-strip-markdown/v2 v2.1.1 h1:hAxUM21Uhznf/FnbVGiJciqzska6iLei22Ijc3q2e28=
github.com/writeas/go-strip-markdown/v2 v2.1.1/go.mod h1:UvvgPJgn1vvN8nWuE5e7v/+qmDu3BSVnKAB6Gl7hFzA=
github.com/xiang90/probing v0.0.0-20190116061207-43a291ad63a2/go.mod h1:UETIi67q53MR2AWcXfiuqkDkRtnGDLqkBTpCHuJHxtU=

View file

@ -57,20 +57,19 @@ func NewModbusMbmdFromConfig(other map[string]interface{}) (api.Meter, error) {
cc.RTU = &b
}
conn, err := modbus.NewConnection(cc.URI, cc.Device, cc.Comset, cc.Baudrate, modbus.ProtocolFromRTU(cc.RTU), cc.ID)
modbus.Lock()
defer modbus.Unlock()
conn, err := modbus.NewConnection(cc.URI, cc.Device, cc.Comset, cc.Baudrate, cc.Settings.Protocol(), cc.ID)
if err != nil {
return nil, err
}
// set non-default delay
if cc.Delay > 0 {
conn.Delay(cc.Delay)
}
// set non-default timeout
if cc.Timeout > 0 {
conn.Timeout(cc.Timeout)
}
conn.Timeout(cc.Timeout)
// set non-default delay
conn.Delay(cc.Delay)
log := util.NewLogger("modbus")
conn.Logger(log.TRACE)

View file

@ -13,9 +13,10 @@ var acceptable = []string{
"missing token",
"mqtt not configured",
"not a SunSpec device",
"missing credentials", // sockets
"power: timeout", // sockets
"missing password", // Powerwall
"connect: connection refused", // sockets
"missing credentials", // sockets
"power: timeout", // sockets
"missing password", // Powerwall
"connect: no route to host",
"connect: connection refused",
"connect: network is unreachable",

View file

@ -41,25 +41,22 @@ func NewModbusFromConfig(other map[string]interface{}) (Provider, error) {
return nil, err
}
conn, err := modbus.NewConnection(cc.URI, cc.Device, cc.Comset, cc.Baudrate, modbus.ProtocolFromRTU(cc.RTU), cc.ID)
modbus.Lock()
defer modbus.Unlock()
conn, err := modbus.NewConnection(cc.URI, cc.Device, cc.Comset, cc.Baudrate, cc.Settings.Protocol(), cc.ID)
if err != nil {
return nil, err
}
// set non-default timeout
if cc.Timeout > 0 {
conn.Timeout(cc.Timeout)
}
conn.Timeout(cc.Timeout)
// set non-default delay
if cc.Delay > 0 {
conn.Delay(cc.Delay)
}
conn.Delay(cc.Delay)
// set non-default connect delay
if cc.ConnectDelay > 0 {
conn.ConnectDelay(cc.ConnectDelay)
}
conn.ConnectDelay(cc.ConnectDelay)
log := util.NewLogger("modbus")
conn.Logger(log.TRACE)

View file

@ -44,25 +44,22 @@ func NewModbusSunspecFromConfig(other map[string]interface{}) (Provider, error)
return nil, err
}
modbus.Lock()
defer modbus.Unlock()
conn, err := modbus.NewConnection(cc.URI, cc.Device, cc.Comset, cc.Baudrate, modbus.Tcp, cc.ID)
if err != nil {
return nil, err
}
// set non-default timeout
if cc.Timeout > 0 {
conn.Timeout(cc.Timeout)
}
conn.Timeout(cc.Timeout)
// set non-default delay
if cc.Delay > 0 {
conn.Delay(cc.Delay)
}
conn.Delay(cc.Delay)
// set non-default connect delay
if cc.ConnectDelay > 0 {
conn.ConnectDelay(cc.ConnectDelay)
}
conn.ConnectDelay(cc.ConnectDelay)
log := util.NewLogger("sunspec")
conn.Logger(log.TRACE)

View file

@ -98,7 +98,7 @@ LOOP:
func (h *handler) HandleDiscreteInputs(req *mbserver.DiscreteInputsRequest) ([]bool, error) {
h.log.TRACE.Printf("read discrete: id %d addr %d qty %d", req.UnitId, req.Addr, req.Quantity)
b, err := h.conn.ReadDiscreteInputsWithSlave(req.UnitId, req.Addr, req.Quantity)
b, err := h.conn.Clone(req.UnitId).ReadDiscreteInputs(req.Addr, req.Quantity)
return h.bytesToBoolResult("read discrete", req.Quantity, b, err)
}
@ -120,24 +120,24 @@ func (h *handler) HandleCoils(req *mbserver.CoilsRequest) ([]bool, error) {
u = 0xFF00
}
b, err := h.conn.WriteSingleCoilWithSlave(req.UnitId, req.Addr, u)
b, err := h.conn.Clone(req.UnitId).WriteSingleCoil(req.Addr, u)
return h.bytesToBoolResult("write coil", req.Quantity, b, err)
}
h.log.TRACE.Printf("write coils: id %d addr %d qty %d val %v", req.UnitId, req.Addr, req.Quantity, req.Args)
args := coilsToBytes(req.Args)
b, err := h.conn.WriteMultipleCoilsWithSlave(req.UnitId, req.Addr, req.Quantity, args)
b, err := h.conn.Clone(req.UnitId).WriteMultipleCoils(req.Addr, req.Quantity, args)
return h.bytesToBoolResult("write coils", req.Quantity, b, err)
}
h.log.TRACE.Printf("read coils: id %d addr %d qty %d", req.UnitId, req.Addr, req.Quantity)
b, err := h.conn.ReadCoilsWithSlave(req.UnitId, req.Addr, req.Quantity)
b, err := h.conn.Clone(req.UnitId).ReadCoils(req.Addr, req.Quantity)
return h.bytesToBoolResult("read coils", req.Quantity, b, err)
}
func (h *handler) HandleInputRegisters(req *mbserver.InputRegistersRequest) ([]uint16, error) {
h.log.TRACE.Printf("read input: id %d addr %d qty %d", req.UnitId, req.Addr, req.Quantity)
b, err := h.conn.ReadInputRegistersWithSlave(req.UnitId, req.Addr, req.Quantity)
b, err := h.conn.Clone(req.UnitId).ReadInputRegisters(req.Addr, req.Quantity)
return h.exceptionToUint16AndError("read input", b, err)
}
@ -154,16 +154,16 @@ func (h *handler) HandleHoldingRegisters(req *mbserver.HoldingRegistersRequest)
if req.WriteFuncCode == gridx.FuncCodeWriteSingleRegister {
h.log.TRACE.Printf("write holding: id %d addr %d val %04x", req.UnitId, req.Addr, req.Args[0])
b, err := h.conn.WriteSingleRegisterWithSlave(req.UnitId, req.Addr, req.Args[0])
b, err := h.conn.Clone(req.UnitId).WriteSingleRegister(req.Addr, req.Args[0])
return h.exceptionToUint16AndError("write holding", b, err)
}
h.log.TRACE.Printf("write holdings: id %d addr %d qty %d val %0x", req.UnitId, req.Addr, req.Quantity, asBytes(req.Args))
b, err := h.conn.WriteMultipleRegistersWithSlave(req.UnitId, req.Addr, req.Quantity, asBytes(req.Args))
b, err := h.conn.Clone(req.UnitId).WriteMultipleRegisters(req.Addr, req.Quantity, asBytes(req.Args))
return h.exceptionToUint16AndError("write multiple holding", b, err)
}
h.log.TRACE.Printf("read holdings: id %d addr %d qty %d", req.UnitId, req.Addr, req.Quantity)
b, err := h.conn.ReadHoldingRegistersWithSlave(req.UnitId, req.Addr, req.Quantity)
b, err := h.conn.Clone(req.UnitId).ReadHoldingRegisters(req.Addr, req.Quantity)
return h.exceptionToUint16AndError("read holding", b, err)
}

View file

@ -12,7 +12,7 @@ import (
)
func StartProxy(port int, config modbus.Settings, readOnly ReadOnlyMode) error {
conn, err := modbus.NewConnection(config.URI, config.Device, config.Comset, config.Baudrate, modbus.ProtocolFromRTU(config.RTU), config.ID)
conn, err := modbus.NewConnection(config.URI, config.Device, config.Comset, config.Baudrate, config.Protocol(), config.ID)
if err != nil {
return err
}

View file

@ -27,21 +27,21 @@ func TestConcurrentRead(t *testing.T) {
require.NoError(t, srv.Start(l))
defer func() { _ = srv.Stop() }()
// client
conn, err := modbus.NewConnection(l.Addr().String(), "", "", 0, modbus.Tcp, 1)
require.NoError(t, err)
var wg sync.WaitGroup
for i := 1; i <= 10; i++ {
wg.Add(1)
go func(id int) {
// client
conn, err := modbus.NewConnection(l.Addr().String(), "", "", 0, modbus.Tcp, uint8(id))
require.NoError(t, err)
for i := 0; i < 50; i++ {
addr := uint16(rand.Int31n(200) + 1)
qty := uint16(rand.Int31n(32) + 1)
b, err := conn.ReadInputRegistersWithSlave(uint8(id), addr, qty)
b, err := conn.ReadInputRegisters(addr, qty)
require.NoError(t, err)
if err == nil {
@ -94,24 +94,24 @@ func TestReadCoils(t *testing.T) {
require.NoError(t, err)
{ // read
b, err := conn.ReadCoilsWithSlave(1, 1, 1)
b, err := conn.ReadCoils(1, 1)
require.NoError(t, err)
assert.Equal(t, []byte{0x01}, b)
b, err = conn.ReadCoilsWithSlave(1, 1, 2)
b, err = conn.ReadCoils(1, 2)
require.NoError(t, err)
assert.Equal(t, []byte{0x03}, b)
b, err = conn.ReadCoilsWithSlave(1, 1, 9)
b, err = conn.ReadCoils(1, 9)
require.NoError(t, err)
assert.Equal(t, []byte{0xFF, 0x01}, b)
}
{ // write
b, err := conn.WriteSingleCoilWithSlave(1, 1, 0xFF00)
b, err := conn.WriteSingleCoil(1, 0xFF00)
require.NoError(t, err)
assert.Equal(t, []byte{0xFF, 0x00}, b)
b, err = conn.WriteMultipleCoilsWithSlave(1, 1, 9, []byte{0xFF, 0x01})
b, err = conn.WriteMultipleCoils(1, 9, []byte{0xFF, 0x01})
require.NoError(t, err)
assert.Equal(t, []byte{0x00, 0x09}, b)
}

View file

@ -8,7 +8,7 @@ params:
- name: usage
choice: ["grid", "pv", "battery"]
- name: modbus
choice: ["rs485", "tcpip"]
choice: ["rs485", "tcpip", "udp"]
baudrate: 9600
id: 247
- name: capacity

113
util/modbus/connection.go Normal file
View file

@ -0,0 +1,113 @@
package modbus
import (
"sync"
"time"
"github.com/grid-x/modbus"
"github.com/volkszaehler/mbmd/meters"
)
// Connection is a logical modbus connection per slave ID sharing a physical connection
type Connection struct {
meters.Connection
mu sync.Mutex
logger *logger
logical meters.Logger
delay time.Duration
}
func (c *Connection) Delay(delay time.Duration) {
c.delay = delay
}
func (c *Connection) Clone(slaveID uint8) *Connection {
return &Connection{
Connection: c.Connection.Clone(slaveID),
logger: c.logger,
}
}
// TODO resolve conflicts
func (c *Connection) ConnectDelay(delay time.Duration) {
if delay > 0 {
c.Connection.ConnectDelay(delay)
}
}
// TODO resolve conflicts
func (c *Connection) Timeout(timeout time.Duration) {
if timeout > 0 {
_ = c.Connection.Timeout(timeout)
}
}
func (c *Connection) Logger(l modbus.Logger) {
c.mu.Lock()
defer c.mu.Unlock()
c.logical = l
}
func (c *Connection) prepare() {
c.mu.Lock()
defer c.mu.Unlock()
time.Sleep(c.delay)
c.logger.Logger(c.logical)
}
func (c *Connection) ReadCoils(address, quantity uint16) ([]byte, error) {
c.prepare()
return c.ModbusClient().ReadCoils(address, quantity)
}
func (c *Connection) WriteSingleCoil(address, value uint16) ([]byte, error) {
c.prepare()
return c.ModbusClient().WriteSingleCoil(address, value)
}
func (c *Connection) ReadInputRegisters(address, quantity uint16) ([]byte, error) {
c.prepare()
return c.ModbusClient().ReadInputRegisters(address, quantity)
}
func (c *Connection) ReadHoldingRegisters(address, quantity uint16) ([]byte, error) {
c.prepare()
return c.ModbusClient().ReadHoldingRegisters(address, quantity)
}
func (c *Connection) WriteSingleRegister(address, value uint16) ([]byte, error) {
c.prepare()
return c.ModbusClient().WriteSingleRegister(address, value)
}
func (c *Connection) WriteMultipleRegisters(address, quantity uint16, value []byte) ([]byte, error) {
c.prepare()
return c.ModbusClient().WriteMultipleRegisters(address, quantity, value)
}
func (c *Connection) ReadDiscreteInputs(address, quantity uint16) (results []byte, err error) {
c.prepare()
return c.ModbusClient().ReadDiscreteInputs(address, quantity)
}
func (c *Connection) WriteMultipleCoils(address, quantity uint16, value []byte) (results []byte, err error) {
c.prepare()
return c.ModbusClient().WriteMultipleCoils(address, quantity, value)
}
func (c *Connection) ReadWriteMultipleRegisters(readAddress, readQuantity, writeAddress, writeQuantity uint16, value []byte) (results []byte, err error) {
c.prepare()
return c.ModbusClient().ReadWriteMultipleRegisters(readAddress, readQuantity, writeAddress, writeQuantity, value)
}
func (c *Connection) MaskWriteRegister(address, andMask, orMask uint16) (results []byte, err error) {
c.prepare()
return c.ModbusClient().MaskWriteRegister(address, andMask, orMask)
}
func (c *Connection) ReadFIFOQueue(address uint16) (results []byte, err error) {
c.prepare()
return c.ModbusClient().ReadFIFOQueue(address)
}

29
util/modbus/log.go Normal file
View file

@ -0,0 +1,29 @@
package modbus
import (
"sync"
"github.com/grid-x/modbus"
"github.com/volkszaehler/mbmd/meters"
)
type logger struct {
mu sync.RWMutex
logger meters.Logger
}
func (l *logger) Logger(logger modbus.Logger) {
l.mu.Lock()
defer l.mu.Unlock()
l.logger = logger
}
func (l *logger) Printf(format string, v ...interface{}) {
l.mu.RLock()
defer l.mu.RUnlock()
if l.logger != nil {
l.logger.Printf(format, v...)
}
}

View file

@ -5,7 +5,6 @@ import (
"fmt"
"strings"
"sync"
"time"
"github.com/evcc-io/evcc/util"
"github.com/volkszaehler/mbmd/meters"
@ -19,6 +18,7 @@ const (
Tcp Protocol = iota
Rtu
Ascii
Udp
CoilOn uint16 = 0xFF00
)
@ -38,9 +38,22 @@ type Settings struct {
SubDevice int
URI, Device, Comset string
Baudrate int
UDP bool
RTU *bool // indicates RTU over TCP if true
}
// Protocol identifies the wire format from the RTU setting
func (s Settings) Protocol() Protocol {
switch {
case s.UDP:
return Udp
case s.RTU != nil && *s.RTU:
return Rtu
default:
return Tcp
}
}
func (s *Settings) String() string {
if s.URI != "" {
return s.URI
@ -48,180 +61,6 @@ func (s *Settings) String() string {
return s.Device
}
// Connection decorates a meters.Connection with transparent slave id and error handling
type Connection struct {
slaveID uint8
mu sync.Mutex
conn meters.Connection
delay time.Duration
}
func (mb *Connection) prepare(slaveID uint8) {
mb.conn.Slave(slaveID)
if mb.delay > 0 {
time.Sleep(mb.delay)
}
}
func (mb *Connection) handle(res []byte, err error) ([]byte, error) {
if err != nil {
mb.conn.Close()
}
return res, err
}
// Delay sets delay so use between subsequent modbus operations
func (mb *Connection) Delay(delay time.Duration) {
mb.delay = delay
}
// ConnectDelay sets the initial delay after connecting before starting communication
func (mb *Connection) ConnectDelay(delay time.Duration) {
mb.conn.ConnectDelay(delay)
}
// Logger sets logger implementation
func (mb *Connection) Logger(logger meters.Logger) {
mb.conn.Logger(logger)
}
// Timeout sets the connection timeout (not idle timeout)
func (mb *Connection) Timeout(timeout time.Duration) {
mb.conn.Timeout(timeout)
}
// ReadCoils wraps the underlying implementation
func (mb *Connection) ReadCoilsWithSlave(slaveID uint8, address, quantity uint16) ([]byte, error) {
mb.mu.Lock()
defer mb.mu.Unlock()
mb.prepare(slaveID)
return mb.handle(mb.conn.ModbusClient().ReadCoils(address, quantity))
}
// WriteSingleCoil wraps the underlying implementation
func (mb *Connection) WriteSingleCoilWithSlave(slaveID uint8, address, value uint16) ([]byte, error) {
mb.mu.Lock()
defer mb.mu.Unlock()
mb.prepare(slaveID)
return mb.handle(mb.conn.ModbusClient().WriteSingleCoil(address, value))
}
// ReadInputRegisters wraps the underlying implementation
func (mb *Connection) ReadInputRegistersWithSlave(slaveID uint8, address, quantity uint16) ([]byte, error) {
mb.mu.Lock()
defer mb.mu.Unlock()
mb.prepare(slaveID)
return mb.handle(mb.conn.ModbusClient().ReadInputRegisters(address, quantity))
}
// ReadHoldingRegisters wraps the underlying implementation
func (mb *Connection) ReadHoldingRegistersWithSlave(slaveID uint8, address, quantity uint16) ([]byte, error) {
mb.mu.Lock()
defer mb.mu.Unlock()
mb.prepare(slaveID)
return mb.handle(mb.conn.ModbusClient().ReadHoldingRegisters(address, quantity))
}
// WriteSingleRegister wraps the underlying implementation
func (mb *Connection) WriteSingleRegisterWithSlave(slaveID uint8, address, value uint16) ([]byte, error) {
mb.mu.Lock()
defer mb.mu.Unlock()
mb.prepare(slaveID)
return mb.handle(mb.conn.ModbusClient().WriteSingleRegister(address, value))
}
// WriteMultipleRegisters wraps the underlying implementation
func (mb *Connection) WriteMultipleRegistersWithSlave(slaveID uint8, address, quantity uint16, value []byte) ([]byte, error) {
mb.mu.Lock()
defer mb.mu.Unlock()
mb.prepare(slaveID)
return mb.handle(mb.conn.ModbusClient().WriteMultipleRegisters(address, quantity, value))
}
// ReadDiscreteInputs wraps the underlying implementation
func (mb *Connection) ReadDiscreteInputsWithSlave(slaveID uint8, address, quantity uint16) (results []byte, err error) {
mb.mu.Lock()
defer mb.mu.Unlock()
mb.prepare(slaveID)
return mb.handle(mb.conn.ModbusClient().ReadDiscreteInputs(address, quantity))
}
// WriteMultipleCoils wraps the underlying implementation
func (mb *Connection) WriteMultipleCoilsWithSlave(slaveID uint8, address, quantity uint16, value []byte) (results []byte, err error) {
mb.mu.Lock()
defer mb.mu.Unlock()
mb.prepare(slaveID)
return mb.handle(mb.conn.ModbusClient().WriteMultipleCoils(address, quantity, value))
}
// ReadWriteMultipleRegisters wraps the underlying implementation
func (mb *Connection) ReadWriteMultipleRegistersWithSlave(slaveID uint8, readAddress, readQuantity, writeAddress, writeQuantity uint16, value []byte) (results []byte, err error) {
mb.mu.Lock()
defer mb.mu.Unlock()
mb.prepare(slaveID)
return mb.handle(mb.conn.ModbusClient().ReadWriteMultipleRegisters(readAddress, readQuantity, writeAddress, writeQuantity, value))
}
// MaskWriteRegister wraps the underlying implementation
func (mb *Connection) MaskWriteRegisterWithSlave(slaveID uint8, address, andMask, orMask uint16) (results []byte, err error) {
mb.mu.Lock()
defer mb.mu.Unlock()
mb.prepare(slaveID)
return mb.handle(mb.conn.ModbusClient().MaskWriteRegister(address, andMask, orMask))
}
// ReadFIFOQueue wraps the underlying implementation
func (mb *Connection) ReadFIFOQueueWithSlave(slaveID uint8, address uint16) (results []byte, err error) {
mb.mu.Lock()
defer mb.mu.Unlock()
mb.prepare(slaveID)
return mb.handle(mb.conn.ModbusClient().ReadFIFOQueue(address))
}
func (mb *Connection) ReadCoils(address, quantity uint16) ([]byte, error) {
return mb.ReadCoilsWithSlave(mb.slaveID, address, quantity)
}
func (mb *Connection) WriteSingleCoil(address, value uint16) ([]byte, error) {
return mb.WriteSingleCoilWithSlave(mb.slaveID, address, value)
}
func (mb *Connection) ReadInputRegisters(address, quantity uint16) ([]byte, error) {
return mb.ReadInputRegistersWithSlave(mb.slaveID, address, quantity)
}
func (mb *Connection) ReadHoldingRegisters(address, quantity uint16) ([]byte, error) {
return mb.ReadHoldingRegistersWithSlave(mb.slaveID, address, quantity)
}
func (mb *Connection) WriteSingleRegister(address, value uint16) ([]byte, error) {
return mb.WriteSingleRegisterWithSlave(mb.slaveID, address, value)
}
func (mb *Connection) WriteMultipleRegisters(address, quantity uint16, value []byte) ([]byte, error) {
return mb.WriteMultipleRegistersWithSlave(mb.slaveID, address, quantity, value)
}
func (mb *Connection) ReadDiscreteInputs(address, quantity uint16) (results []byte, err error) {
return mb.ReadDiscreteInputsWithSlave(mb.slaveID, address, quantity)
}
func (mb *Connection) WriteMultipleCoils(address, quantity uint16, value []byte) (results []byte, err error) {
return mb.WriteMultipleCoilsWithSlave(mb.slaveID, address, quantity, value)
}
func (mb *Connection) ReadWriteMultipleRegisters(readAddress, readQuantity, writeAddress, writeQuantity uint16, value []byte) (results []byte, err error) {
return mb.ReadWriteMultipleRegistersWithSlave(mb.slaveID, readAddress, readQuantity, writeAddress, writeQuantity, value)
}
func (mb *Connection) MaskWriteRegister(address, andMask, orMask uint16) (results []byte, err error) {
return mb.MaskWriteRegisterWithSlave(mb.slaveID, address, andMask, orMask)
}
func (mb *Connection) ReadFIFOQueue(address uint16) (results []byte, err error) {
return mb.ReadFIFOQueueWithSlave(mb.slaveID, address)
}
var (
connections = make(map[string]meters.Connection)
mu sync.Mutex
@ -240,65 +79,69 @@ func registeredConnection(key string, newConn meters.Connection) meters.Connecti
return newConn
}
// ProtocolFromRTU identifies the wire format from the RTU setting
func ProtocolFromRTU(rtu *bool) Protocol {
if rtu != nil && *rtu {
return Rtu
}
return Tcp
}
// NewConnection creates physical modbus device from config
func NewConnection(uri, device, comset string, baudrate int, proto Protocol, slaveID uint8) (*Connection, error) {
var conn meters.Connection
if device != "" && uri != "" {
return nil, errors.New("invalid modbus configuration: can only have either uri or device")
conn, err := physicalConnection(proto, Settings{
URI: uri,
Device: device,
Comset: comset,
Baudrate: baudrate,
})
if err != nil {
return nil, err
}
if device != "" {
switch strings.ToUpper(comset) {
res := &Connection{
Connection: conn.Clone(slaveID),
logger: new(logger),
}
return res, nil
}
func physicalConnection(proto Protocol, cfg Settings) (meters.Connection, error) {
var conn meters.Connection
if (cfg.Device != "") == (cfg.URI != "") {
return nil, errors.New("invalid modbus configuration: must have either uri or device")
}
if cfg.Device != "" {
switch strings.ToUpper(cfg.Comset) {
case "8N1", "8E1", "8N2":
case "80":
comset = "8E1"
cfg.Comset = "8E1"
default:
return nil, fmt.Errorf("invalid comset: %s", comset)
return nil, fmt.Errorf("invalid comset: %s", cfg.Comset)
}
if baudrate == 0 {
if cfg.Baudrate == 0 {
return nil, errors.New("invalid modbus configuration: need baudrate and comset")
}
if proto == Ascii {
conn = registeredConnection(device, meters.NewASCII(device, baudrate, comset))
conn = registeredConnection(cfg.Device, meters.NewASCII(cfg.Device, cfg.Baudrate, cfg.Comset))
} else {
conn = registeredConnection(device, meters.NewRTU(device, baudrate, comset))
conn = registeredConnection(cfg.Device, meters.NewRTU(cfg.Device, cfg.Baudrate, cfg.Comset))
}
}
if uri != "" {
uri = util.DefaultPort(uri, 502)
if cfg.URI != "" {
cfg.URI = util.DefaultPort(cfg.URI, 502)
switch proto {
case Udp:
conn = registeredConnection(cfg.URI, meters.NewRTUOverUDP(cfg.URI))
case Rtu:
conn = registeredConnection(uri, meters.NewRTUOverTCP(uri))
conn = registeredConnection(cfg.URI, meters.NewRTUOverTCP(cfg.URI))
case Ascii:
conn = registeredConnection(uri, meters.NewASCIIOverTCP(uri))
conn = registeredConnection(cfg.URI, meters.NewASCIIOverTCP(cfg.URI))
default:
conn = registeredConnection(uri, meters.NewTCP(uri))
conn = registeredConnection(cfg.URI, meters.NewTCP(cfg.URI))
}
}
if conn == nil {
return nil, errors.New("invalid modbus configuration: need either uri or device")
}
slaveConn := &Connection{
slaveID: slaveID,
conn: conn,
}
return slaveConn, nil
return conn, nil
}
// NewDevice creates physical modbus device from config

13
util/modbus/mutex.go Normal file
View file

@ -0,0 +1,13 @@
package modbus
import "sync"
var mu2 sync.Mutex
func Lock() {
mu2.Lock()
}
func Unlock() {
mu2.Unlock()
}

View file

@ -435,6 +435,7 @@ modbus:
interfaces:
rs485: ["rs485serial", "rs485tcpip"]
tcpip: ["tcpip"]
udp: ["udp"]
types:
rs485serial:
description:
@ -476,6 +477,18 @@ modbus:
- reference: true
name: port
default: 502
udp:
description:
generic: UDP
params:
- reference: true
referencename: modbusid
name: id
- reference: true
name: host
- reference: true
name: port
default: 502
devicegroups:
generic:

View file

@ -13,6 +13,11 @@ rtu: true
# Modbus TCP
uri: {{ .host }}:{{ .port }}
rtu: false
{{- else if or (eq .modbus "udp") .udp }}
# Modbus UDP
uri: {{ .host }}:{{ if (ne .port "502") }}{{ .port }}{{ else }}8899{{ end }}
udp: true
rtu: true
{{- else }}
# configuration error - should not happen
modbusConnectionTypeNotDefined: {{ .modbus }}

View file

@ -48,6 +48,8 @@ func TestClass(t *testing.T, class Class, instantiate func(t *testing.T, values
// we only test one modbus setup
if slices.Contains(modbusChoices, ModbusChoiceTCPIP) {
values[ModbusKeyTCPIP] = true
} else if slices.Contains(modbusChoices, ModbusChoiceUDP) {
values[ModbusKeyUDP] = true
} else {
values[ModbusKeyRS485TCPIP] = true
}

View file

@ -17,9 +17,11 @@ const (
ModbusChoiceRS485 = "rs485"
ModbusChoiceTCPIP = "tcpip"
ModbusChoiceUDP = "udp"
ModbusKeyRS485Serial = "rs485serial"
ModbusKeyRS485TCPIP = "rs485tcpip"
ModbusKeyTCPIP = "tcpip"
ModbusKeyUDP = "udp"
ModbusParamNameId = "id"
ModbusParamNameDevice = "device"
@ -37,7 +39,7 @@ const (
RenderModeInstance
)
var ValidModbusChoices = []string{ModbusChoiceRS485, ModbusChoiceTCPIP}
var ValidModbusChoices = []string{ModbusChoiceRS485, ModbusChoiceTCPIP, ModbusChoiceUDP}
const (
CapabilityISO151182 = "iso151182" // ISO 15118-2 support
@ -62,7 +64,7 @@ var predefinedTemplateProperties = []string{
"type", "template", "name",
ModbusParamNameId, ModbusParamNameDevice, ModbusParamNameBaudrate, ModbusParamNameComset,
ModbusParamNameURI, ModbusParamNameHost, ModbusParamNamePort, ModbusParamNameRTU,
ModbusKeyTCPIP, ModbusKeyRS485Serial, ModbusKeyRS485TCPIP,
ModbusKeyTCPIP, ModbusKeyUDP, ModbusKeyRS485Serial, ModbusKeyRS485TCPIP,
}
// TextLanguage contains language-specific texts