Sunspec: serialize access to a shared device tree (#33014)
This commit is contained in:
parent
9f70a18501
commit
1b909edd9b
3 changed files with 73 additions and 8 deletions
|
|
@ -19,7 +19,7 @@ import (
|
|||
type ModbusSunspec struct {
|
||||
log *util.Logger
|
||||
conn *modbus.Connection
|
||||
device *sunsdev.SunSpec
|
||||
device *sunspecDevice
|
||||
op modbus.SunSpecOperation
|
||||
scale float64
|
||||
mask uint64
|
||||
|
|
@ -86,7 +86,7 @@ func NewModbusSunspecFromConfig(ctx context.Context, other map[string]any) (Plug
|
|||
device := sunspecSubDevices.Get(conn, cc.SubDevice)
|
||||
if device == nil {
|
||||
// silence KOSTAL implementation errors
|
||||
device = sunsdev.NewDevice("sunspec", cc.SubDevice)
|
||||
device = &sunspecDevice{SunSpec: sunsdev.NewDevice("sunspec", cc.SubDevice)}
|
||||
if err := device.InitializeWithTree(devices); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
|
@ -111,6 +111,9 @@ func NewModbusSunspecFromConfig(ctx context.Context, other map[string]any) (Plug
|
|||
mask: mask,
|
||||
}
|
||||
|
||||
device.mu.Lock()
|
||||
defer device.mu.Unlock()
|
||||
|
||||
for _, op := range ops {
|
||||
if _, _, err := device.QueryPointAny(conn, op.Model, op.Block, op.Point); err == nil {
|
||||
mb.op = op
|
||||
|
|
@ -131,6 +134,9 @@ func recoverToError(err *error) {
|
|||
func (m *ModbusSunspec) floatGetter() (f float64, err error) {
|
||||
defer recoverToError(&err)
|
||||
|
||||
m.device.mu.Lock()
|
||||
defer m.device.mu.Unlock()
|
||||
|
||||
res, err := m.device.QueryPoint(
|
||||
m.conn,
|
||||
m.op.Model,
|
||||
|
|
@ -173,6 +179,9 @@ func (m *ModbusSunspec) BoolGetter() (func() (bool, error), error) {
|
|||
return func() (res bool, err error) {
|
||||
defer recoverToError(&err)
|
||||
|
||||
m.device.mu.Lock()
|
||||
defer m.device.mu.Unlock()
|
||||
|
||||
_, point, err := m.blockPoint()
|
||||
if err != nil {
|
||||
return false, err
|
||||
|
|
@ -222,6 +231,8 @@ func sunspecBool(val int64, mask uint64) bool {
|
|||
return val != 0
|
||||
}
|
||||
|
||||
// blockPoint reads the block and returns its point. The device lock must be held
|
||||
// by the caller until the point value has been consumed or written.
|
||||
func (m *ModbusSunspec) blockPoint() (block sunspec.Block, point sunspec.Point, err error) {
|
||||
defer recoverToError(&err)
|
||||
|
||||
|
|
@ -242,7 +253,10 @@ var _ FloatSetter = (*Modbus)(nil)
|
|||
|
||||
// FloatSetter executes configured modbus write operation and implements FloatSetter
|
||||
func (m *ModbusSunspec) FloatSetter(_ string) (func(float64) error, error) {
|
||||
m.device.mu.Lock()
|
||||
block, point, err := m.blockPoint()
|
||||
m.device.mu.Unlock()
|
||||
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
|
@ -252,6 +266,10 @@ func (m *ModbusSunspec) FloatSetter(_ string) (func(float64) error, error) {
|
|||
return func(val float64) (err error) {
|
||||
defer recoverToError(&err)
|
||||
|
||||
// setting the point and writing it must not interleave with a concurrent read
|
||||
m.device.mu.Lock()
|
||||
defer m.device.mu.Unlock()
|
||||
|
||||
val = val * m.scale
|
||||
switch typ {
|
||||
case typelabel.Float32:
|
||||
|
|
@ -268,7 +286,10 @@ var _ IntSetter = (*Modbus)(nil)
|
|||
|
||||
// IntSetter executes configured modbus write operation and implements IntSetter
|
||||
func (m *ModbusSunspec) IntSetter(_ string) (func(int64) error, error) {
|
||||
m.device.mu.Lock()
|
||||
block, point, err := m.blockPoint()
|
||||
m.device.mu.Unlock()
|
||||
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
|
@ -278,6 +299,10 @@ func (m *ModbusSunspec) IntSetter(_ string) (func(int64) error, error) {
|
|||
return func(val int64) (err error) {
|
||||
defer recoverToError(&err)
|
||||
|
||||
// setting the point and writing it must not interleave with a concurrent read
|
||||
m.device.mu.Lock()
|
||||
defer m.device.mu.Unlock()
|
||||
|
||||
val = int64(float64(val) * m.scale)
|
||||
|
||||
// SetValue is used to include the scale factor when writing
|
||||
|
|
|
|||
|
|
@ -2,12 +2,21 @@ package plugin
|
|||
|
||||
import (
|
||||
"strconv"
|
||||
"sync"
|
||||
|
||||
gosunspec "github.com/andig/gosunspec"
|
||||
"github.com/evcc-io/evcc/util/modbus"
|
||||
"github.com/volkszaehler/mbmd/meters/sunspec"
|
||||
)
|
||||
|
||||
// sunspecDevice serializes access to a device tree shared by all values of one
|
||||
// subdevice. gosunspec decodes into the shared model on every block read, hence
|
||||
// a read and the point access following it must not interleave with another one.
|
||||
type sunspecDevice struct {
|
||||
mu sync.Mutex
|
||||
*sunspec.SunSpec
|
||||
}
|
||||
|
||||
var sunspecDevices = sunspecDeviceCache{
|
||||
data: make(map[string][]gosunspec.Device),
|
||||
}
|
||||
|
|
@ -26,15 +35,15 @@ func (c *sunspecDeviceCache) Put(conn *modbus.Connection, devices []gosunspec.De
|
|||
}
|
||||
|
||||
var sunspecSubDevices = sunspecSubDeviceCache{
|
||||
data: make(map[string][]*sunspec.SunSpec),
|
||||
data: make(map[string][]*sunspecDevice),
|
||||
}
|
||||
|
||||
// sunspecSubDeviceCache is a cache for a sunspec devices's models
|
||||
type sunspecSubDeviceCache struct {
|
||||
data map[string][]*sunspec.SunSpec
|
||||
data map[string][]*sunspecDevice
|
||||
}
|
||||
|
||||
func (c *sunspecSubDeviceCache) Get(conn *modbus.Connection, subDevice int) *sunspec.SunSpec {
|
||||
func (c *sunspecSubDeviceCache) Get(conn *modbus.Connection, subDevice int) *sunspecDevice {
|
||||
addr := sunspecSubdeviceAddr(conn, subDevice)
|
||||
for _, dev := range c.data[addr] {
|
||||
if dev.Descriptor().SubDevice == subDevice {
|
||||
|
|
@ -45,7 +54,7 @@ func (c *sunspecSubDeviceCache) Get(conn *modbus.Connection, subDevice int) *sun
|
|||
return nil
|
||||
}
|
||||
|
||||
func (c *sunspecSubDeviceCache) Put(conn *modbus.Connection, subDevice int, dev *sunspec.SunSpec) {
|
||||
func (c *sunspecSubDeviceCache) Put(conn *modbus.Connection, subDevice int, dev *sunspecDevice) {
|
||||
addr := sunspecSubdeviceAddr(conn, subDevice)
|
||||
c.data[addr] = append(c.data[addr], dev)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,12 +1,14 @@
|
|||
package plugin
|
||||
|
||||
import (
|
||||
"sync"
|
||||
"testing"
|
||||
|
||||
sunspec "github.com/andig/gosunspec"
|
||||
"github.com/andig/gosunspec/memory"
|
||||
"github.com/andig/gosunspec/models/model704"
|
||||
"github.com/evcc-io/evcc/util/modbus"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
sunsdev "github.com/volkszaehler/mbmd/meters/sunspec"
|
||||
)
|
||||
|
|
@ -33,7 +35,7 @@ func TestSunspecBool(t *testing.T) {
|
|||
|
||||
// newSunspecTestDevice builds an in-memory SunSpec model 704 (DER AC Controls)
|
||||
// device; mbmd never touches the modbus.Client argument, so no connection is needed.
|
||||
func newSunspecTestDevice(t *testing.T) (*sunsdev.SunSpec, sunspec.Block) {
|
||||
func newSunspecTestDevice(t *testing.T) (*sunspecDevice, sunspec.Block) {
|
||||
t.Helper()
|
||||
|
||||
slab, err := memory.NewSlabBuilder().AddModel(model704.ModelID).Build()
|
||||
|
|
@ -47,7 +49,7 @@ func newSunspecTestDevice(t *testing.T) (*sunsdev.SunSpec, sunspec.Block) {
|
|||
|
||||
block := devices[0].MustModel(sunspec.ModelId(model704.ModelID)).MustBlock(0)
|
||||
|
||||
dev := sunsdev.NewDevice("test")
|
||||
dev := &sunspecDevice{SunSpec: sunsdev.NewDevice("test")}
|
||||
require.NoError(t, dev.InitializeWithTree(devices))
|
||||
|
||||
return dev, block
|
||||
|
|
@ -111,3 +113,32 @@ func TestSunspecBoolGetterInt(t *testing.T) {
|
|||
require.NoError(t, err)
|
||||
require.True(t, got)
|
||||
}
|
||||
|
||||
// TestSunspecConcurrentSharedDevice reads different points of one shared device
|
||||
// tree concurrently, as the config page's parallel capability probes do.
|
||||
func TestSunspecConcurrentSharedDevice(t *testing.T) {
|
||||
dev, block := newSunspecTestDevice(t)
|
||||
block.MustPoint(model704.WMaxLimPctEna).SetEnum16(1)
|
||||
block.MustPoint(model704.WMaxLimPct).SetUint16(50)
|
||||
require.NoError(t, block.Write(model704.WMaxLimPctEna))
|
||||
require.NoError(t, block.Write(model704.WMaxLimPct))
|
||||
|
||||
var wg sync.WaitGroup
|
||||
for range 20 {
|
||||
for _, point := range []string{model704.WMaxLimPctEna, model704.WMaxLimPct} {
|
||||
wg.Go(func() {
|
||||
mb := &ModbusSunspec{
|
||||
device: dev,
|
||||
op: modbus.SunSpecOperation{Model: model704.ModelID, Point: point},
|
||||
}
|
||||
|
||||
g, err := mb.BoolGetter()
|
||||
assert.NoError(t, err) // require fails the test goroutine only
|
||||
|
||||
_, err = g()
|
||||
assert.NoError(t, err)
|
||||
})
|
||||
}
|
||||
}
|
||||
wg.Wait()
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue