aa55 udp: dedupe concurrent block reads with single flight (#30589)
This commit is contained in:
parent
b4f7a42306
commit
4dcff007d7
3 changed files with 110 additions and 18 deletions
|
|
@ -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
|
||||
// ---------------------------------------------------------------------------
|
||||
|
|
|
|||
|
|
@ -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.
|
||||
|
|
|
|||
|
|
@ -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.
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue