diff --git a/plugin/aa55/aa55_test.go b/plugin/aa55/aa55_test.go index 217e03566..72826c67d 100644 --- a/plugin/aa55/aa55_test.go +++ b/plugin/aa55/aa55_test.go @@ -4,8 +4,6 @@ import ( "encoding/binary" "encoding/hex" "math" - "sync" - "sync/atomic" "testing" "github.com/stretchr/testify/assert" @@ -304,81 +302,6 @@ func TestET_SoC_GW25K(t *testing.T) { assertBlockOffset(t, capGW25kETBattery, 14, "uint16be", 1.0, 100.0) } -// --------------------------------------------------------------------------- -// Cache -// --------------------------------------------------------------------------- - -func TestCache_GetMiss(t *testing.T) { - c := newResponseCache() - _, ok := c.get([]byte("nope")) - assert.False(t, ok) -} - -func TestCache_PutGet(t *testing.T) { - c := newResponseCache() - c.put([]byte("k"), []byte{1, 2, 3}) - got, ok := c.get([]byte("k")) - require.True(t, ok) - assert.Equal(t, []byte{1, 2, 3}, got) -} - -// TestCache_FetchSingleFlight verifies that concurrent fetches for the same key -// collapse into a single load (one UDP exchange), all observe the same payload, -// and the cache is left warm for subsequent reads. -func TestCache_FetchSingleFlight(t *testing.T) { - c := newResponseCache() - key := []byte("10.0.0.1:8899/f703891c007d") - want := []byte{0xde, 0xad, 0xbe, 0xef} - - var calls atomic.Int32 - var once sync.Once - entered := make(chan struct{}) - release := make(chan struct{}) - load := func() ([]byte, error) { - calls.Add(1) - once.Do(func() { close(entered) }) - <-release // hold the flight open while followers pile up - return want, nil - } - - const n = 8 - var wg sync.WaitGroup - payloads := make([][]byte, n) - - wg.Go(func() { - got, _, err := c.fetch(key, load) - assert.NoError(t, err) - payloads[0] = got - }) - - <-entered // ensure the flight is open before followers join - - for i := 1; i < n; i++ { - wg.Go(func() { - got, _, err := c.fetch(key, load) - assert.NoError(t, err) - payloads[i] = got - }) - } - - close(release) - wg.Wait() - - assert.Equal(t, int32(1), calls.Load(), "concurrent reads of the same block must share one exchange") - for i := range n { - assert.Equal(t, want, payloads[i]) - } - - // cache is now warm: a subsequent read hits without loading again - got, ok, err := c.fetch(key, func() ([]byte, error) { - t.Fatal("must not load on warm cache") - return nil, nil - }) - require.NoError(t, err) - assert.True(t, ok) - assert.Equal(t, want, got) -} - // --------------------------------------------------------------------------- // Helpers // --------------------------------------------------------------------------- diff --git a/plugin/aa55/aa55udp.go b/plugin/aa55/aa55udp.go index 3bd329c3d..a0dbfca36 100644 --- a/plugin/aa55/aa55udp.go +++ b/plugin/aa55/aa55udp.go @@ -7,8 +7,17 @@ import ( "time" "github.com/evcc-io/evcc/util" + "github.com/evcc-io/evcc/util/modbus" ) +// cacheTTL serves all sources within one poll cycle (well under 1s) while +// forcing a fresh read on the next cycle. +const cacheTTL = 2 * time.Second + +// cache de-duplicates block reads across all AA55UDP instances so multiple +// sources covering the same (host, block) share one UDP exchange per cycle. +var cache = modbus.NewCache(cacheTTL) + // AA55UDP is the GoodWe AA55-over-UDP source plugin transport. // // Two read modes are supported, both built from logical parameters @@ -26,19 +35,10 @@ type AA55UDP struct { offset int // byte offset into the response payload (0 for register reads) decode string // int32be | uint32be | uint32nan | int16be | uint16be | float32be scale float64 - cacheKey []byte // precomputed cache key (remoteAddr/pdu); nil disables caching + cacheKey string // precomputed cache key (remoteAddr/pdu); empty disables caching delay time.Duration // minimum gap between sends to the inverter (0 disables) } -// Block describes the enclosing register block to fetch in block-read mode. -// When set, one UDP exchange reads Count registers starting at Register, and -// each source extracts its own target register at the computed offset, sharing -// the response via the cache. -type Block struct { - Register uint16 - Count uint16 -} - // readConfig holds the resolved read mode configuration. type readConfig struct { pdu []byte @@ -49,7 +49,7 @@ type readConfig struct { // buildReadConfig resolves the read mode from the target register (register, // count, id) and the optional enclosing block. In both modes the PDU is built // on the Go side; the template only supplies logical parameters. -func buildReadConfig(id int, register, count uint16, block *Block) (readConfig, error) { +func buildReadConfig(id int, register, count uint16, block *modbus.Block) (readConfig, error) { if id < 0 || id > 255 { return readConfig{}, fmt.Errorf("id must be 0-255, got %d", id) } @@ -63,13 +63,12 @@ func buildReadConfig(id int, register, count uint16, block *Block) (readConfig, if block.Count == 0 { return readConfig{}, errors.New("block count must be ≥ 1") } - // The target register must fit entirely within the block. - if register < block.Register || uint32(register)+uint32(count) > uint32(block.Register)+uint32(block.Count) { + if !block.Contains(register, count) { return readConfig{}, fmt.Errorf("register %d+%d does not fit in block %d+%d", register, count, block.Register, block.Count) } return readConfig{ pdu: buildPDU(byte(id), block.Register, block.Count), - offset: int(register-block.Register) * 2, + offset: block.ByteOffset(register), useCache: true, }, nil } @@ -81,7 +80,7 @@ func buildReadConfig(id int, register, count uint16, block *Block) (readConfig, // New constructs an AA55UDP from a high-level configuration. It validates // decode, resolves the read mode (register vs block), and wraps the conn. // The caller is responsible for dialling conn. -func New(log *util.Logger, conn *net.UDPConn, id int, register, count uint16, block *Block, decode string, scale float64, delay time.Duration) (*AA55UDP, error) { +func New(log *util.Logger, conn *net.UDPConn, id int, register, count uint16, block *modbus.Block, decode string, scale float64, delay time.Duration) (*AA55UDP, error) { if err := validateDecode(decode); err != nil { return nil, err } @@ -99,7 +98,7 @@ func New(log *util.Logger, conn *net.UDPConn, id int, register, count uint16, bl offset: cfg.offset, } if cfg.useCache { - ap.cacheKey = []byte(conn.RemoteAddr().String() + "/" + string(cfg.pdu)) + ap.cacheKey = conn.RemoteAddr().String() + "/" + string(cfg.pdu) } return ap, nil } @@ -133,11 +132,11 @@ func (p *AA55UDP) query() (float64, error) { // serves and de-duplicates requests. func (p *AA55UDP) fetch() ([]byte, error) { // Register mode: single targeted read, no caching. - if p.cacheKey == nil { + if p.cacheKey == "" { return p.exchange() } - payload, ok, err := cache.fetch(p.cacheKey, p.exchange) + payload, ok, err := cache.Fetch(p.cacheKey, p.exchange) if err != nil { return nil, err } diff --git a/plugin/aa55/aa55udp_test.go b/plugin/aa55/aa55udp_test.go index c067b3826..58b1cd098 100644 --- a/plugin/aa55/aa55udp_test.go +++ b/plugin/aa55/aa55udp_test.go @@ -6,6 +6,7 @@ import ( "testing" "github.com/evcc-io/evcc/util" + "github.com/evcc-io/evcc/util/modbus" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" ) @@ -113,7 +114,7 @@ func TestBuildReadConfig_RegisterMode(t *testing.T) { // register/count and the target register's offset is computed within it. // ET grid (0x8943) within block READ 125 @ 0x891C → offset (35139-35100)*2 = 78. func TestBuildReadConfig_BlockMode(t *testing.T) { - cfg, err := buildReadConfig(0xF7, 0x8943, 2, &Block{Register: 0x891C, Count: 125}) + cfg, err := buildReadConfig(0xF7, 0x8943, 2, &modbus.Block{Register: 0x891C, Count: 125}) require.NoError(t, err) assert.Equal(t, []byte{0xf7, 0x03, 0x89, 0x1c, 0x00, 0x7d}, cfg.pdu) assert.Equal(t, 78, cfg.offset) @@ -124,11 +125,11 @@ func TestBuildReadConfig_BlockMode(t *testing.T) { // that does not fit entirely within the configured block. func TestBuildReadConfig_RejectsRegisterOutsideBlock(t *testing.T) { // before block start - _, err := buildReadConfig(0xF7, 0x8900, 2, &Block{Register: 0x891C, Count: 125}) + _, err := buildReadConfig(0xF7, 0x8900, 2, &modbus.Block{Register: 0x891C, Count: 125}) require.Error(t, err) // past block end (0x891C+125 = 0x8999; 0x8998+2 overruns) - _, err = buildReadConfig(0xF7, 0x8998, 2, &Block{Register: 0x891C, Count: 125}) + _, err = buildReadConfig(0xF7, 0x8998, 2, &modbus.Block{Register: 0x891C, Count: 125}) require.Error(t, err) } diff --git a/plugin/aa55/cache.go b/plugin/aa55/cache.go deleted file mode 100644 index c9f1c6788..000000000 --- a/plugin/aa55/cache.go +++ /dev/null @@ -1,90 +0,0 @@ -package aa55 - -import ( - "sync" - "time" - - "golang.org/x/sync/singleflight" -) - -const cacheTTL = 2 * time.Second - -// cache is the package-level response cache shared across all AA55UDP plugin -// instances. Sharing at package level ensures that multiple source blocks for -// the same (host, pdu) pair — e.g. the four Ppv string registers all using -// READ 125 @ 0x891C — share one UDP exchange per poll cycle. -// -// TTL is 2 s: long enough to serve all source blocks within one evcc poll -// cycle (which completes in well under 1 s), short enough that the next cycle -// always fetches fresh data. -var cache = newResponseCache() - -type cacheEntry struct { - payload []byte - expiresAt time.Time -} - -type responseCache struct { - mu sync.Mutex - data map[string]cacheEntry - flight singleflight.Group -} - -func newResponseCache() *responseCache { - return &responseCache{data: make(map[string]cacheEntry)} -} - -// fetch returns the cached payload for key if it is fresh. On a miss, load is -// invoked exactly once across all concurrent callers sharing the same key. -func (c *responseCache) fetch(key []byte, load func() ([]byte, error)) ([]byte, bool, error) { - if payload, ok := c.get(key); ok { - return payload, true, nil - } - - payload, err, _ := c.flight.Do(string(key), func() (any, error) { - // re-check under the flight: a prior flight may have populated the - // cache between our miss above and acquiring the call. - if payload, ok := c.get(key); ok { - return payload, nil - } - payload, err := load() - if err != nil { - return nil, err - } - c.put(key, payload) - return payload, nil - }) - if err != nil { - return nil, false, err - } - return payload.([]byte), false, nil -} - -// get returns the cached payload if it exists and is fresh, or (nil, false) -// otherwise. Expired entries are deleted on access. The map lookup -// m[string(key)] is alloc-free — the Go compiler elides the conversion. -func (c *responseCache) get(key []byte) ([]byte, bool) { - c.mu.Lock() - defer c.mu.Unlock() - - entry, ok := c.data[string(key)] - if !ok { - return nil, false - } - if time.Now().After(entry.expiresAt) { - delete(c.data, string(key)) - return nil, false - } - return entry.payload, true -} - -// put inserts or overwrites a payload in the cache. -func (c *responseCache) put(key, payload []byte) { - c.mu.Lock() - defer c.mu.Unlock() - - c.data[string(key)] = cacheEntry{ - payload: payload, - expiresAt: time.Now().Add(cacheTTL), - } -} diff --git a/plugin/aa55udp.go b/plugin/aa55udp.go index 75d0aedfd..7b2320ce7 100644 --- a/plugin/aa55udp.go +++ b/plugin/aa55udp.go @@ -9,6 +9,7 @@ import ( "github.com/evcc-io/evcc/plugin/aa55" "github.com/evcc-io/evcc/util" + "github.com/evcc-io/evcc/util/modbus" "github.com/evcc-io/evcc/util/request" ) @@ -16,16 +17,17 @@ func init() { registry.AddCtx("aa55udp", NewAA55UDPFromConfig) } -// NewAA55UDPFromConfig creates a GoodWe AA55-over-UDP plugin. +// NewAA55UDPFromConfig creates a GoodWe AA55-over-UDP plugin. The register is +// configured as a modbus.Register; its count is derived from the decode width. // // Register read mode (single register): // // source: aa55udp // host: 192.168.1.26 -// id: 127 # 0x7F for DT/DNS/ES/EM (default); 247 (0xF7) for ET/EH/BT/BH -// register: 30127 -// count: 2 # 1 = 16-bit, 2 = 32-bit -// decode: int32be +// id: 127 # 0x7F for DT/DNS/ES/EM (default); 247 (0xF7) for ET/EH/BT/BH +// register: +// address: 30127 +// decode: int32be # int32be | uint32be | uint32nan | int16be | uint16be | float32be // scale: 1.0 // // Block read mode (fetch an enclosing block once, extract the target register): @@ -33,27 +35,24 @@ func init() { // source: aa55udp // host: 192.168.1.26 # inverter IP; port 8899 is always used // id: 247 # 0x7F (default) for DT/DNS/ES/EM; 247 (0xF7) for ET/EH/BT/BH -// register: 35139 # target register -// count: 2 # 1 = 16-bit, 2 = 32-bit +// register: +// address: 35139 # target register +// decode: int32be // block: # enclosing block fetched in a single UDP exchange // register: 35100 # block start register // count: 125 # block length (registers) -// decode: int32be # int32be | uint32be | uint32nan | int16be | uint16be | float32be // scale: 1.0 # optional multiplier (default 1.0) // delay: 100ms # optional min gap between sends to one inverter (0 disables) func NewAA55UDPFromConfig(ctx context.Context, other map[string]any) (Plugin, error) { cc := struct { Host string Id int - Register uint16 - Count uint16 - Block *aa55.Block - Decode string + Register modbus.Register + Block *modbus.Block Scale float64 Delay time.Duration }{ Id: int(aa55.InverterAddr), - Count: 2, Scale: 1.0, } @@ -61,6 +60,11 @@ func NewAA55UDPFromConfig(ctx context.Context, other map[string]any) (Plugin, er return nil, err } + count, err := cc.Register.Length() + if err != nil { + return nil, err + } + raddr, err := net.ResolveUDPAddr("udp4", net.JoinHostPort(cc.Host, "8899")) if err != nil { return nil, err @@ -72,7 +76,7 @@ func NewAA55UDPFromConfig(ctx context.Context, other map[string]any) (Plugin, er return nil, err } - res, err := aa55.New(util.NewLogger("aa55udp"), conn, cc.Id, cc.Register, cc.Count, cc.Block, cc.Decode, cc.Scale, cc.Delay) + res, err := aa55.New(util.NewLogger("aa55udp"), conn, cc.Id, cc.Register.Address, count, cc.Block, cc.Register.Decode, cc.Scale, cc.Delay) if err != nil { return nil, fmt.Errorf("aa55udp: %w", err) } diff --git a/plugin/aa55udp_test.go b/plugin/aa55udp_test.go index a89163270..3f5cc8a5e 100644 --- a/plugin/aa55udp_test.go +++ b/plugin/aa55udp_test.go @@ -7,16 +7,14 @@ import ( "github.com/stretchr/testify/require" ) -// TestAA55UDPFromConfig_Block verifies that the nested block map decodes and a -// target register that fits within the block is accepted (block-read mode). +// TestAA55UDPFromConfig_Block verifies the modbus.Register and nested block map +// decode, with the count (2 registers) derived from the int32be decode width. func TestAA55UDPFromConfig_Block(t *testing.T) { _, err := NewAA55UDPFromConfig(context.Background(), map[string]any{ "host": "127.0.0.1", "id": 247, - "register": 35139, - "count": 2, + "register": map[string]any{"address": 35139, "decode": "int32be"}, "block": map[string]any{"register": 35100, "count": 125}, - "decode": "int32be", }) require.NoError(t, err) } @@ -27,10 +25,8 @@ func TestAA55UDPFromConfig_BlockRejectsOutOfRange(t *testing.T) { _, err := NewAA55UDPFromConfig(context.Background(), map[string]any{ "host": "127.0.0.1", "id": 247, - "register": 36017, // outside READ 125 @ 35100 - "count": 2, + "register": map[string]any{"address": 36017, "decode": "float32be"}, // outside READ 125 @ 35100 "block": map[string]any{"register": 35100, "count": 125}, - "decode": "float32be", }) require.Error(t, err) } @@ -39,9 +35,7 @@ func TestAA55UDPFromConfig_BlockRejectsOutOfRange(t *testing.T) { func TestAA55UDPFromConfig_Register(t *testing.T) { _, err := NewAA55UDPFromConfig(context.Background(), map[string]any{ "host": "127.0.0.1", - "register": 30127, - "count": 2, - "decode": "int32be", + "register": map[string]any{"address": 30127, "decode": "int32be"}, }) require.NoError(t, err) } diff --git a/plugin/modbus.go b/plugin/modbus.go index 91bbe41d5..9eae8a9e4 100644 --- a/plugin/modbus.go +++ b/plugin/modbus.go @@ -3,6 +3,7 @@ package plugin import ( "bytes" "context" + "errors" "fmt" "math" "strings" @@ -13,11 +14,20 @@ import ( gridx "github.com/grid-x/modbus" ) +// modbusBlockTTL dedups block reads within a poll cycle, short enough to force +// a fresh read on the next cycle. +const modbusBlockTTL = time.Second + +// modbusBlockCache lets sources covering the same (device, block) share one +// read per poll cycle. +var modbusBlockCache = modbus.NewCache(modbusBlockTTL) + // Modbus implements modbus RTU and TCP access type Modbus struct { log *util.Logger conn *modbus.Connection reg modbus.Register + block modbus.Block scale float64 } @@ -30,6 +40,7 @@ func NewModbusFromConfig(ctx context.Context, other map[string]any) (Plugin, err cc := struct { modbus.Settings `mapstructure:",squash"` Register modbus.Register + Block modbus.Block Scale float64 Delay time.Duration ConnectDelay time.Duration @@ -66,16 +77,27 @@ func NewModbusFromConfig(ctx context.Context, other map[string]any) (Plugin, err return nil, err } + if cc.Block.Count > 0 { + fc, err := cc.Register.FuncCode() + if err != nil { + return nil, err + } + if fc != gridx.FuncCodeReadHoldingRegisters && fc != gridx.FuncCodeReadInputRegisters { + return nil, errors.New("block read requires holding or input register type") + } + } + mb := &Modbus{ log: log, conn: conn, reg: cc.Register, + block: cc.Block, scale: cc.Scale, } return mb, nil } -func (m *Modbus) readBytes(op modbus.RegisterOperation) ([]byte, error) { +func (m *Modbus) read(op modbus.RegisterOperation) ([]byte, error) { switch op.FuncCode { case gridx.FuncCodeReadHoldingRegisters: return m.conn.ReadHoldingRegisters(op.Addr, op.Length) @@ -91,6 +113,29 @@ func (m *Modbus) readBytes(op modbus.RegisterOperation) ([]byte, error) { } } +// readBytes returns the bytes for op. In block mode it fetches the enclosing +// block once (shared via modbusBlockCache) and extracts op at its offset. +func (m *Modbus) readBytes(op modbus.RegisterOperation) ([]byte, error) { + if m.block.Count == 0 { + return m.read(op) + } + + blockOp := modbus.RegisterOperation{FuncCode: op.FuncCode, Addr: m.block.Register, Length: m.block.Count} + key := fmt.Sprintf("%s/%d/%d/%d", m.conn.Addr(), op.FuncCode, m.block.Register, m.block.Count) + + payload, hit, err := modbusBlockCache.Fetch(key, func() ([]byte, error) { + return m.read(blockOp) + }) + if err != nil { + return nil, err + } + if hit { + m.log.TRACE.Printf("block cache hit %s", key) + } + + return m.block.Extract(op, payload) +} + var _ FloatGetter = (*Modbus)(nil) // FloatGetter implements func() (float64, error) diff --git a/plugin/modbus_test.go b/plugin/modbus_test.go new file mode 100644 index 000000000..31b5b8b9b --- /dev/null +++ b/plugin/modbus_test.go @@ -0,0 +1,53 @@ +package plugin + +import ( + "encoding/binary" + "testing" + + "github.com/evcc-io/evcc/util/modbus" + gridx "github.com/grid-x/modbus" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// huaweiGridBlock mirrors the shared grid block of the huawei-sun2000-hybrid +// template: 16 registers (32 bytes) starting at 37107. +var huaweiGridBlock = modbus.Block{Register: 37107, Count: 16} + +func op(addr, length uint16) modbus.RegisterOperation { + return modbus.RegisterOperation{FuncCode: gridx.FuncCodeReadHoldingRegisters, Addr: addr, Length: length} +} + +// TestExtractBlock verifies that each Huawei grid register is sliced from the +// block payload at the correct offset and decodes to the expected value. +func TestExtractBlock(t *testing.T) { + payload := make([]byte, 2*huaweiGridBlock.Count) + // power @ 37113 (offset 12), int32 = 1234 W + binary.BigEndian.PutUint32(payload[12:], uint32(int32(1234))) + // imported energy @ 37121 (offset 28), uint32 = 567890 + binary.BigEndian.PutUint32(payload[28:], 567890) + + power, err := huaweiGridBlock.Extract(op(37113, 2), payload) + require.NoError(t, err) + powerDec, err := (modbus.Register{Type: "holding", Decode: "int32"}).DecodeFunc() + require.NoError(t, err) + assert.Equal(t, 1234.0, powerDec(power)) + + energy, err := huaweiGridBlock.Extract(op(37121, 2), payload) + require.NoError(t, err) + energyDec, err := (modbus.Register{Type: "holding", Decode: "uint32"}).DecodeFunc() + require.NoError(t, err) + assert.Equal(t, 567890.0, energyDec(energy)) +} + +func TestExtractBlockBounds(t *testing.T) { + payload := make([]byte, 2*huaweiGridBlock.Count) + + // register outside the block is rejected + _, err := huaweiGridBlock.Extract(op(37200, 2), payload) + require.Error(t, err) + + // short payload is rejected + _, err = huaweiGridBlock.Extract(op(37121, 2), payload[:10]) + require.Error(t, err) +} diff --git a/templates/definition/meter/goodwe-wifi-dt.yaml b/templates/definition/meter/goodwe-wifi-dt.yaml index c04730439..9785b2205 100644 --- a/templates/definition/meter/goodwe-wifi-dt.yaml +++ b/templates/definition/meter/goodwe-wifi-dt.yaml @@ -14,13 +14,13 @@ render: | power: source: aa55udp host: {{ .host }} - register: 30127 - count: 2 - decode: int32be + register: + address: 30127 + decode: int32be energy: source: aa55udp host: {{ .host }} - register: 30145 - count: 2 - decode: uint32be + register: + address: 30145 + decode: uint32be scale: 0.1 diff --git a/templates/definition/meter/goodwe-wifi-es.yaml b/templates/definition/meter/goodwe-wifi-es.yaml index e7eb69805..9cab17415 100644 --- a/templates/definition/meter/goodwe-wifi-es.yaml +++ b/templates/definition/meter/goodwe-wifi-es.yaml @@ -17,29 +17,29 @@ render: | power: source: aa55udp host: {{ or .host .uri }} - register: 29964 - count: 2 - decode: int32be + register: + address: 29964 + decode: int32be {{- end }} {{- if eq .usage "pv" }} power: source: aa55udp host: {{ or .host .uri }} - register: 29958 - count: 2 - decode: int32be + register: + address: 29958 + decode: int32be {{- end }} {{- if eq .usage "battery" }} power: source: aa55udp host: {{ or .host .uri }} - register: 29970 - count: 2 - decode: int32be + register: + address: 29970 + decode: int32be soc: source: aa55udp host: {{ or .host .uri }} - register: 29966 - count: 1 - decode: uint16be + register: + address: 29966 + decode: uint16be {{- end }} diff --git a/templates/definition/meter/goodwe-wifi-et.yaml b/templates/definition/meter/goodwe-wifi-et.yaml index add5295ad..99b5fca3f 100644 --- a/templates/definition/meter/goodwe-wifi-et.yaml +++ b/templates/definition/meter/goodwe-wifi-et.yaml @@ -25,17 +25,17 @@ render: | block: # READ 125 @ 0x891C register: 35100 count: 125 - register: 35139 # grid power - count: 2 - decode: int32be + register: + address: 35139 # grid power + decode: int32be delay: {{ .delay }} energy: source: aa55udp host: {{ or .host .uri }} id: 247 - register: 36017 - count: 2 - decode: float32be + register: + address: 36017 + decode: float32be scale: 0.001 delay: {{ .delay }} {{- end }} @@ -49,9 +49,9 @@ render: | block: # READ 125 @ 0x891C register: 35100 count: 125 - register: 35105 # Ppv1 - count: 2 - decode: uint32nan + register: + address: 35105 # Ppv1 + decode: uint32nan delay: {{ .delay }} - source: aa55udp host: {{ or .host .uri }} @@ -59,9 +59,9 @@ render: | block: register: 35100 count: 125 - register: 35109 # Ppv2 - count: 2 - decode: uint32nan + register: + address: 35109 # Ppv2 + decode: uint32nan delay: {{ .delay }} - source: aa55udp host: {{ or .host .uri }} @@ -69,9 +69,9 @@ render: | block: register: 35100 count: 125 - register: 35113 # Ppv3 - count: 2 - decode: uint32nan + register: + address: 35113 # Ppv3 + decode: uint32nan delay: {{ .delay }} - source: aa55udp host: {{ or .host .uri }} @@ -79,9 +79,9 @@ render: | block: register: 35100 count: 125 - register: 35117 # Ppv4 - count: 2 - decode: uint32nan + register: + address: 35117 # Ppv4 + decode: uint32nan delay: {{ .delay }} energy: source: aa55udp @@ -90,9 +90,9 @@ render: | block: register: 35100 count: 125 - register: 35191 # pv energy - count: 2 - decode: uint32be + register: + address: 35191 # pv energy + decode: uint32be scale: 0.1 delay: {{ .delay }} {{- end }} @@ -104,9 +104,9 @@ render: | block: # READ 125 @ 0x891C register: 35100 count: 125 - register: 35182 # battery power - count: 2 - decode: int32be + register: + address: 35182 # battery power + decode: int32be delay: {{ .delay }} energy: source: aa55udp @@ -115,9 +115,9 @@ render: | block: register: 35100 count: 125 - register: 35209 # battery discharge energy - count: 2 - decode: uint32be + register: + address: 35209 # battery discharge energy + decode: uint32be scale: 0.1 delay: {{ .delay }} soc: @@ -127,8 +127,8 @@ render: | block: # READ 13 @ 0x9088 register: 37000 count: 13 - register: 37007 # soc - count: 1 - decode: uint16be + register: + address: 37007 # soc + decode: uint16be delay: {{ .delay }} {{- end }} diff --git a/templates/definition/meter/huawei-sun2000-hybrid.yaml b/templates/definition/meter/huawei-sun2000-hybrid.yaml index 0910b3055..2cb9ee9a3 100644 --- a/templates/definition/meter/huawei-sun2000-hybrid.yaml +++ b/templates/definition/meter/huawei-sun2000-hybrid.yaml @@ -82,6 +82,9 @@ render: | address: 37113 # Grid import/export power type: holding decode: int32nan + block: # 37107-37122 + register: 37107 + count: 16 scale: -1 energy: source: modbus @@ -91,6 +94,9 @@ render: | address: 37121 # Active energy imported from the grid type: holding decode: uint32nan + block: # 37107-37122 + register: 37107 + count: 16 scale: 0.01 returnenergy: source: modbus @@ -100,6 +106,9 @@ render: | address: 37119 # Active energy exported to the grid type: holding decode: uint32nan + block: # 37107-37122 + register: 37107 + count: 16 scale: 0.01 currents: - source: modbus @@ -109,6 +118,9 @@ render: | address: 37107 # Huawei phase A grid current type: holding decode: int32nan + block: # 37107-37122 + register: 37107 + count: 16 scale: -0.01 - source: modbus {{- include "modbus" . | indent 2 }} @@ -117,6 +129,9 @@ render: | address: 37109 # Huawei phase B grid current type: holding decode: int32nan + block: # 37107-37122 + register: 37107 + count: 16 scale: -0.01 - source: modbus {{- include "modbus" . | indent 2 }} @@ -125,6 +140,9 @@ render: | address: 37111 # Huawei phase C grid current type: holding decode: int32nan + block: # 37107-37122 + register: 37107 + count: 16 scale: -0.01 {{- end }} {{- if eq .usage "pv" }} diff --git a/util/modbus/block.go b/util/modbus/block.go new file mode 100644 index 000000000..8f54f26ce --- /dev/null +++ b/util/modbus/block.go @@ -0,0 +1,38 @@ +package modbus + +import "fmt" + +// Block describes an enclosing register block: Count 16-bit registers starting +// at Register, fetched once per poll cycle and shared via a Cache. +type Block struct { + Register uint16 + Count uint16 +} + +// Contains reports whether the count registers starting at register fit +// entirely within the block. +func (b Block) Contains(register, count uint16) bool { + return register >= b.Register && + uint32(register)+uint32(count) <= uint32(b.Register)+uint32(b.Count) +} + +// ByteOffset returns the byte offset of register within the block payload +// (each register occupies two bytes). +func (b Block) ByteOffset(register uint16) int { + return int(register-b.Register) * 2 +} + +// Extract returns the bytes for op sliced out of a block payload. +func (b Block) Extract(op RegisterOperation, payload []byte) ([]byte, error) { + if !b.Contains(op.Addr, op.Length) { + return nil, fmt.Errorf("register %d+%d does not fit in block %d+%d", op.Addr, op.Length, b.Register, b.Count) + } + + offset := b.ByteOffset(op.Addr) + end := offset + int(op.Length)*2 + if len(payload) < end { + return nil, fmt.Errorf("block payload too short (len=%d, need=%d)", len(payload), end) + } + + return payload[offset:end], nil +} diff --git a/util/modbus/block_test.go b/util/modbus/block_test.go new file mode 100644 index 000000000..3bd9cf0dd --- /dev/null +++ b/util/modbus/block_test.go @@ -0,0 +1,100 @@ +package modbus + +import ( + "sync" + "sync/atomic" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +const testTTL = 2 * time.Second + +func TestBlockContains(t *testing.T) { + b := Block{Register: 37107, Count: 16} + + assert.True(t, b.Contains(37107, 2), "first register") + assert.True(t, b.Contains(37121, 2), "last register exactly at block end") + assert.False(t, b.Contains(37106, 2), "before block start") + assert.False(t, b.Contains(37122, 2), "overruns block end") +} + +func TestBlockByteOffset(t *testing.T) { + b := Block{Register: 37107, Count: 16} + + assert.Equal(t, 0, b.ByteOffset(37107)) + assert.Equal(t, 12, b.ByteOffset(37113)) + assert.Equal(t, 28, b.ByteOffset(37121)) +} + +func TestCacheGetMiss(t *testing.T) { + c := NewCache(testTTL) + _, ok := c.get("nope") + assert.False(t, ok) +} + +func TestCachePutGet(t *testing.T) { + c := NewCache(testTTL) + c.put("k", []byte{1, 2, 3}) + got, ok := c.get("k") + require.True(t, ok) + assert.Equal(t, []byte{1, 2, 3}, got) +} + +// TestCacheFetchSingleFlight verifies that concurrent fetches for the same key +// collapse into a single load and leave the cache warm. +func TestCacheFetchSingleFlight(t *testing.T) { + c := NewCache(testTTL) + key := "10.0.0.1::1/3/37107/16" + want := []byte{0xde, 0xad, 0xbe, 0xef} + + var calls atomic.Int32 + var once sync.Once + entered := make(chan struct{}) + release := make(chan struct{}) + load := func() ([]byte, error) { + calls.Add(1) + once.Do(func() { close(entered) }) + <-release // hold the flight open while followers pile up + return want, nil + } + + const n = 8 + var wg sync.WaitGroup + payloads := make([][]byte, n) + + wg.Go(func() { + got, _, err := c.Fetch(key, load) + assert.NoError(t, err) + payloads[0] = got + }) + + <-entered // ensure the flight is open before followers join + + for i := 1; i < n; i++ { + wg.Go(func() { + got, _, err := c.Fetch(key, load) + assert.NoError(t, err) + payloads[i] = got + }) + } + + close(release) + wg.Wait() + + assert.Equal(t, int32(1), calls.Load(), "concurrent reads of the same block must share one exchange") + for i := range n { + assert.Equal(t, want, payloads[i]) + } + + // cache is now warm: a subsequent read hits without loading again + got, ok, err := c.Fetch(key, func() ([]byte, error) { + t.Fatal("must not load on warm cache") + return nil, nil + }) + require.NoError(t, err) + assert.True(t, ok) + assert.Equal(t, want, got) +} diff --git a/util/modbus/blockcache.go b/util/modbus/blockcache.go new file mode 100644 index 000000000..1422c1943 --- /dev/null +++ b/util/modbus/blockcache.go @@ -0,0 +1,81 @@ +package modbus + +import ( + "sync" + "time" + + "golang.org/x/sync/singleflight" +) + +// Cache is a TTL response cache with single-flight de-duplication, sharing one +// device exchange per key. The TTL is caller-chosen (depends on poll cadence). +type Cache struct { + ttl time.Duration + mu sync.Mutex + data map[string]cacheEntry + flight singleflight.Group +} + +type cacheEntry struct { + payload []byte + expiresAt time.Time +} + +// NewCache returns a Cache that holds entries for ttl. +func NewCache(ttl time.Duration) *Cache { + return &Cache{ttl: ttl, data: make(map[string]cacheEntry)} +} + +// Fetch returns the cached payload for key if it is fresh. On a miss, load is +// invoked exactly once across all concurrent callers sharing the same key. +func (c *Cache) Fetch(key string, load func() ([]byte, error)) ([]byte, bool, error) { + if payload, ok := c.get(key); ok { + return payload, true, nil + } + + payload, err, _ := c.flight.Do(key, func() (any, error) { + // re-check under the flight: a prior flight may have populated the + // cache between our miss above and acquiring the call. + if payload, ok := c.get(key); ok { + return payload, nil + } + payload, err := load() + if err != nil { + return nil, err + } + c.put(key, payload) + return payload, nil + }) + if err != nil { + return nil, false, err + } + return payload.([]byte), false, nil +} + +// get returns the cached payload if it exists and is fresh. Expired entries are +// deleted on access. +func (c *Cache) get(key string) ([]byte, bool) { + c.mu.Lock() + defer c.mu.Unlock() + + e, ok := c.data[key] + if !ok { + return nil, false + } + if time.Now().After(e.expiresAt) { + delete(c.data, key) + return nil, false + } + return e.payload, true +} + +// put inserts or overwrites a payload in the cache. +func (c *Cache) put(key string, payload []byte) { + c.mu.Lock() + defer c.mu.Unlock() + + c.data[key] = cacheEntry{ + payload: payload, + expiresAt: time.Now().Add(c.ttl), + } +}