Modbus: add shared block reading (#30846)
This commit is contained in:
parent
efc57ed702
commit
18609ee874
15 changed files with 428 additions and 262 deletions
|
|
@ -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
|
||||
// ---------------------------------------------------------------------------
|
||||
|
|
|
|||
|
|
@ -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
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -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),
|
||||
}
|
||||
}
|
||||
Loading…
Add table
Add a link
Reference in a new issue