From 4dcff007d705429c8d582df3e4c78d414251e5f0 Mon Sep 17 00:00:00 2001 From: andig Date: Sun, 7 Jun 2026 11:35:52 +0200 Subject: [PATCH] aa55 udp: dedupe concurrent block reads with single flight (#30589) --- plugin/aa55/aa55_test.go | 59 ++++++++++++++++++++++++++++++++++++++++ plugin/aa55/aa55udp.go | 36 +++++++++++++----------- plugin/aa55/cache.go | 33 ++++++++++++++++++++-- 3 files changed, 110 insertions(+), 18 deletions(-) diff --git a/plugin/aa55/aa55_test.go b/plugin/aa55/aa55_test.go index 04f3301a2..217e03566 100644 --- a/plugin/aa55/aa55_test.go +++ b/plugin/aa55/aa55_test.go @@ -4,6 +4,8 @@ import ( "encoding/binary" "encoding/hex" "math" + "sync" + "sync/atomic" "testing" "github.com/stretchr/testify/assert" @@ -320,6 +322,63 @@ func TestCache_PutGet(t *testing.T) { 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 d99b82374..a93b2b0fe 100644 --- a/plugin/aa55/aa55udp.go +++ b/plugin/aa55/aa55udp.go @@ -127,30 +127,34 @@ func (p *AA55UDP) query() (float64, error) { return v * p.scale, nil } -// fetch returns the response payload, using caching for block-read mode. +// fetch returns the response payload. In block-read mode the shared cache +// serves and de-duplicates requests. func (p *AA55UDP) fetch() ([]byte, error) { - if p.cacheKey != nil { - if payload, ok := cache.get(p.cacheKey); ok { - p.log.TRACE.Printf("cache hit for %s pdu=%x", p.conn.RemoteAddr(), p.pdu) - return payload, nil - } + // Register mode: single targeted read, no caching. + if p.cacheKey == nil { + return p.exchange() } + payload, ok, err := cache.fetch(p.cacheKey, p.exchange) + if err != nil { + return nil, err + } + if ok { + p.log.TRACE.Printf("cache hit for %s pdu=%x", p.conn.RemoteAddr(), p.pdu) + } + return payload, nil +} + +// exchange performs one request/response round trip and returns the response +// payload with the AA55 header stripped. It is the cache-miss path shared by +// single flight in block-read mode. +func (p *AA55UDP) exchange() ([]byte, error) { packet := append(p.pdu, modbusCRC16(p.pdu)...) raw, err := p.sendRecv(packet) if err != nil { return nil, err } - - payload, err := stripHeader(raw) - if err != nil { - return nil, err - } - - if p.cacheKey != nil { - cache.put(p.cacheKey, payload) - } - return payload, nil + return stripHeader(raw) } // sendRecv sends packet over p.conn and returns the raw response bytes. diff --git a/plugin/aa55/cache.go b/plugin/aa55/cache.go index aeb38e457..c9f1c6788 100644 --- a/plugin/aa55/cache.go +++ b/plugin/aa55/cache.go @@ -3,6 +3,8 @@ package aa55 import ( "sync" "time" + + "golang.org/x/sync/singleflight" ) const cacheTTL = 2 * time.Second @@ -23,14 +25,41 @@ type cacheEntry struct { } type responseCache struct { - mu sync.Mutex - data map[string]cacheEntry + 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.