281 lines
8 KiB
Go
281 lines
8 KiB
Go
package aa55
|
|
|
|
import (
|
|
"errors"
|
|
"fmt"
|
|
"net"
|
|
"time"
|
|
|
|
"github.com/evcc-io/evcc/util"
|
|
"github.com/evcc-io/evcc/util/modbus"
|
|
gridx "github.com/grid-x/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
|
|
|
|
// readTimeout bounds the wait for a response
|
|
const readTimeout = 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
|
|
// (id, register, count) on the Go side:
|
|
//
|
|
// Register read: single register, value at offset 0 of response payload.
|
|
// Block read: enclosing block fetched once, target register extracted at
|
|
// its computed offset. Multiple sources sharing the same
|
|
// (host, block) share one UDP exchange per poll cycle via the
|
|
// response cache.
|
|
type AA55UDP struct {
|
|
log *util.Logger
|
|
conn *net.UDPConn
|
|
pdu []byte // 6-byte PDU body, no CRC
|
|
offset int // byte offset into the response payload (0 for register reads)
|
|
length int // value length in bytes
|
|
decode func([]byte) float64 // modbus register decoder
|
|
scale float64
|
|
cacheKey string // precomputed cache key (remoteAddr/pdu); empty disables caching
|
|
delay time.Duration // minimum gap between sends to the inverter (0 disables)
|
|
|
|
// write mode (set by NewSetter)
|
|
id byte
|
|
reg modbus.Register
|
|
writable bool
|
|
}
|
|
|
|
// readConfig holds the resolved read mode configuration.
|
|
type readConfig struct {
|
|
pdu []byte
|
|
offset int
|
|
useCache bool
|
|
}
|
|
|
|
// 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 *modbus.Block) (readConfig, error) {
|
|
if id < 0 || id > 255 {
|
|
return readConfig{}, fmt.Errorf("id must be 0-255, got %d", id)
|
|
}
|
|
if count == 0 {
|
|
return readConfig{}, errors.New("count must be ≥ 1")
|
|
}
|
|
|
|
// Block mode: fetch the whole block and extract the target register at its
|
|
// offset. Multiple sources sharing the same block share one UDP exchange.
|
|
if block != nil {
|
|
if block.Count == 0 {
|
|
return readConfig{}, errors.New("block count must be ≥ 1")
|
|
}
|
|
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: block.ByteOffset(register),
|
|
useCache: true,
|
|
}, nil
|
|
}
|
|
|
|
// Register mode: single targeted read, value at offset 0, no caching.
|
|
return readConfig{pdu: buildPDU(byte(id), register, count), useCache: false}, nil
|
|
}
|
|
|
|
// New constructs an AA55UDP from a high-level configuration. The register count
|
|
// and decoder are derived from reg; the read mode is resolved from block.
|
|
func New(log *util.Logger, conn *net.UDPConn, id int, reg modbus.Register, block *modbus.Block, scale float64, delay time.Duration) (*AA55UDP, error) {
|
|
count, err := reg.Length()
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
decode, err := reg.DecodeFunc()
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
cfg, err := buildReadConfig(id, reg.Address, count, block)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
ap := &AA55UDP{
|
|
log: log,
|
|
conn: conn,
|
|
decode: decode,
|
|
length: int(count) * 2,
|
|
scale: scale,
|
|
delay: delay,
|
|
pdu: cfg.pdu,
|
|
offset: cfg.offset,
|
|
}
|
|
if cfg.useCache {
|
|
ap.cacheKey = conn.RemoteAddr().String() + "/" + string(cfg.pdu)
|
|
}
|
|
return ap, nil
|
|
}
|
|
|
|
// NewSetter constructs a write-mode AA55UDP for a single holding register.
|
|
// The register type must be a write type (writesingle/writemultiple).
|
|
func NewSetter(log *util.Logger, conn *net.UDPConn, id int, reg modbus.Register, scale float64, delay time.Duration) (*AA55UDP, error) {
|
|
if id < 0 || id > 255 {
|
|
return nil, fmt.Errorf("id must be 0-255, got %d", id)
|
|
}
|
|
if err := reg.Error(); err != nil {
|
|
return nil, err
|
|
}
|
|
return &AA55UDP{
|
|
log: log,
|
|
conn: conn,
|
|
scale: scale,
|
|
delay: delay,
|
|
id: byte(id),
|
|
reg: reg,
|
|
writable: true,
|
|
}, nil
|
|
}
|
|
|
|
// FloatSetter implements the evcc plugin.FloatSetter interface.
|
|
func (p *AA55UDP) FloatSetter(_ string) (func(float64) error, error) {
|
|
return p.writeFunc()
|
|
}
|
|
|
|
// IntSetter implements the evcc plugin.IntSetter interface.
|
|
func (p *AA55UDP) IntSetter(_ string) (func(int64) error, error) {
|
|
set, err := p.writeFunc()
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
return func(val int64) error {
|
|
return set(float64(val))
|
|
}, nil
|
|
}
|
|
|
|
// writeFunc builds the register writer derived from the register's func code.
|
|
func (p *AA55UDP) writeFunc() (func(float64) error, error) {
|
|
fc, err := p.reg.FuncCode()
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
encode, err := p.reg.EncodeFunc()
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
return func(val float64) error {
|
|
val *= p.scale
|
|
|
|
var pdu []byte
|
|
switch fc {
|
|
case gridx.FuncCodeWriteSingleRegister:
|
|
pdu = buildWriteSinglePDU(p.id, p.reg.Address, uint16(val))
|
|
case gridx.FuncCodeWriteMultipleRegisters:
|
|
b, err := encode(val)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
pdu = buildWriteMultiplePDU(p.id, p.reg.Address, b)
|
|
default:
|
|
return fmt.Errorf("invalid func code: %d", fc)
|
|
}
|
|
|
|
raw, err := p.sendRecv(append(pdu, modbusCRC16(pdu)...))
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if err := validateWriteResponse(raw, fc); err != nil {
|
|
return err
|
|
}
|
|
|
|
// a write invalidates cached reads for all sources of this inverter
|
|
cache.Clear()
|
|
return nil
|
|
}, nil
|
|
}
|
|
|
|
// FloatGetter implements the evcc plugin.FloatGetter interface.
|
|
func (p *AA55UDP) FloatGetter() (func() (float64, error), error) {
|
|
if p.writable {
|
|
return nil, errors.New("register configured for write")
|
|
}
|
|
return p.query, nil
|
|
}
|
|
|
|
// query fetches the payload and returns the decoded, scaled value at p.offset.
|
|
func (p *AA55UDP) query() (float64, error) {
|
|
payload, err := p.fetch()
|
|
if err != nil {
|
|
return 0, err
|
|
}
|
|
|
|
end := p.offset + p.length
|
|
if len(payload) < end {
|
|
return 0, fmt.Errorf("payload too short (len=%d, need=%d)", len(payload), end)
|
|
}
|
|
|
|
return p.decode(payload[p.offset:end]) * p.scale, nil
|
|
}
|
|
|
|
// fetch returns the response payload. In block-read mode the shared cache
|
|
// serves and de-duplicates requests.
|
|
func (p *AA55UDP) fetch() ([]byte, error) {
|
|
// Register mode: single targeted read, no caching.
|
|
if p.cacheKey == "" {
|
|
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
|
|
}
|
|
return stripHeader(raw)
|
|
}
|
|
|
|
// sendRecv sends packet over p.conn and returns the raw response bytes. When a
|
|
// delay is set it serializes and spaces exchanges to the same inverter.
|
|
func (p *AA55UDP) sendRecv(packet []byte) ([]byte, error) {
|
|
if p.delay > 0 {
|
|
g := pace.gate(p.conn.RemoteAddr().String())
|
|
g.mu.Lock()
|
|
defer g.mu.Unlock()
|
|
g.wait(p.delay)
|
|
}
|
|
|
|
p.log.TRACE.Printf("send %s: %x", p.conn.RemoteAddr(), packet)
|
|
|
|
if _, err := p.conn.Write(packet); err != nil {
|
|
return nil, fmt.Errorf("write: %w", err)
|
|
}
|
|
|
|
if err := p.conn.SetReadDeadline(time.Now().Add(readTimeout)); err != nil {
|
|
return nil, fmt.Errorf("deadline: %w", err)
|
|
}
|
|
|
|
buf := make([]byte, 512)
|
|
n, err := p.conn.Read(buf)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("read: %w", err)
|
|
}
|
|
|
|
p.log.TRACE.Printf("recv %s: %x", p.conn.RemoteAddr(), buf[:n])
|
|
|
|
return buf[:n], nil
|
|
}
|