Easee: warn on rogue CommandResponse not triggered by evcc (#27916)
Co-authored-by: Claude Sonnet 4.6 <noreply@anthropic.com>
This commit is contained in:
parent
56c1d56626
commit
16d1258795
2 changed files with 344 additions and 45 deletions
143
charger/easee.go
143
charger/easee.go
|
|
@ -61,12 +61,15 @@ type Easee struct {
|
|||
phaseMode int
|
||||
currentPower, sessionEnergy, totalEnergy,
|
||||
currentL1, currentL2, currentL3 float64
|
||||
rfid string
|
||||
lp loadpoint.API
|
||||
cmdC chan easee.SignalRCommandResponse
|
||||
obsC chan easee.Observation
|
||||
obsTime map[easee.ObservationID]time.Time
|
||||
startDone func()
|
||||
rfid string
|
||||
lp loadpoint.API
|
||||
cmdMu sync.Mutex
|
||||
pendingTicks map[int64]chan easee.SignalRCommandResponse
|
||||
pendingByID map[easee.ObservationID]chan easee.SignalRCommandResponse
|
||||
expectedOrphans map[easee.ObservationID]int
|
||||
obsC chan easee.Observation
|
||||
obsTime map[easee.ObservationID]time.Time
|
||||
startDone func()
|
||||
}
|
||||
|
||||
func init() {
|
||||
|
|
@ -107,15 +110,17 @@ func NewEasee(ctx context.Context, user, password, charger string, timeout time.
|
|||
done := make(chan struct{})
|
||||
|
||||
c := &Easee{
|
||||
Helper: request.NewHelper(log),
|
||||
charger: charger,
|
||||
authorize: authorize,
|
||||
log: log,
|
||||
current: 6, // default current
|
||||
startDone: sync.OnceFunc(func() { close(done) }),
|
||||
cmdC: make(chan easee.SignalRCommandResponse),
|
||||
obsC: make(chan easee.Observation),
|
||||
obsTime: make(map[easee.ObservationID]time.Time),
|
||||
Helper: request.NewHelper(log),
|
||||
charger: charger,
|
||||
authorize: authorize,
|
||||
log: log,
|
||||
current: 6, // default current
|
||||
startDone: sync.OnceFunc(func() { close(done) }),
|
||||
pendingTicks: make(map[int64]chan easee.SignalRCommandResponse),
|
||||
pendingByID: make(map[easee.ObservationID]chan easee.SignalRCommandResponse),
|
||||
expectedOrphans: make(map[easee.ObservationID]int),
|
||||
obsC: make(chan easee.Observation),
|
||||
obsTime: make(map[easee.ObservationID]time.Time),
|
||||
}
|
||||
|
||||
c.Client.Timeout = timeout
|
||||
|
|
@ -211,6 +216,48 @@ func (c *Easee) waitForOptionalState() {
|
|||
c.log.WARN.Println("did not receive full state from cloud")
|
||||
}
|
||||
|
||||
func (c *Easee) registerPendingTick(tick int64, ch chan easee.SignalRCommandResponse) {
|
||||
c.cmdMu.Lock()
|
||||
c.pendingTicks[tick] = ch
|
||||
c.cmdMu.Unlock()
|
||||
}
|
||||
|
||||
func (c *Easee) unregisterPendingTick(tick int64) {
|
||||
c.cmdMu.Lock()
|
||||
delete(c.pendingTicks, tick)
|
||||
c.cmdMu.Unlock()
|
||||
}
|
||||
|
||||
func (c *Easee) registerPendingByID(id easee.ObservationID, ch chan easee.SignalRCommandResponse) {
|
||||
c.cmdMu.Lock()
|
||||
defer c.cmdMu.Unlock()
|
||||
c.pendingByID[id] = ch
|
||||
}
|
||||
|
||||
func (c *Easee) unregisterPendingByID(id easee.ObservationID) {
|
||||
c.cmdMu.Lock()
|
||||
defer c.cmdMu.Unlock()
|
||||
delete(c.pendingByID, id)
|
||||
}
|
||||
|
||||
func (c *Easee) registerExpectedOrphan(ids ...easee.ObservationID) {
|
||||
c.cmdMu.Lock()
|
||||
defer c.cmdMu.Unlock()
|
||||
for _, id := range ids {
|
||||
c.expectedOrphans[id]++
|
||||
}
|
||||
}
|
||||
|
||||
func (c *Easee) consumeExpectedOrphan(id easee.ObservationID) bool {
|
||||
c.cmdMu.Lock()
|
||||
defer c.cmdMu.Unlock()
|
||||
if c.expectedOrphans[id] > 0 {
|
||||
c.expectedOrphans[id]--
|
||||
return true
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
// check c.obsTime for presence of ALL of the following keys: easee.SESSION_ENERGY, easee.LIFETIME_ENERGY
|
||||
func (c *Easee) optionalStatePresent() bool {
|
||||
c.mux.Lock()
|
||||
|
|
@ -383,12 +430,34 @@ func (c *Easee) CommandResponse(i json.RawMessage) {
|
|||
c.log.ERROR.Printf("invalid message: %s %v", i, err)
|
||||
return
|
||||
}
|
||||
|
||||
c.log.TRACE.Printf("CommandResponse %s: %+v", res.SerialNumber, res)
|
||||
|
||||
select {
|
||||
case c.cmdC <- res:
|
||||
default:
|
||||
obsID := easee.ObservationID(res.ID)
|
||||
|
||||
c.cmdMu.Lock()
|
||||
chTick, tickOk := c.pendingTicks[res.Ticks]
|
||||
chID, idOk := c.pendingByID[obsID]
|
||||
c.cmdMu.Unlock()
|
||||
|
||||
if tickOk {
|
||||
chTick <- res
|
||||
return
|
||||
}
|
||||
|
||||
if idOk {
|
||||
chID <- res
|
||||
return
|
||||
}
|
||||
|
||||
if c.consumeExpectedOrphan(obsID) {
|
||||
return
|
||||
}
|
||||
|
||||
c.log.WARN.Printf("rogue CommandResponse: charger %s ObservationID=%s Ticks=%d "+
|
||||
"(accepted=%v, resultCode=%d) which was not triggered by evcc — "+
|
||||
"another system may be controlling this charger",
|
||||
res.SerialNumber, obsID, res.Ticks, res.WasAccepted, res.ResultCode)
|
||||
}
|
||||
|
||||
func (c *Easee) chargers() ([]easee.Charger, error) {
|
||||
|
|
@ -545,26 +614,28 @@ func (c *Easee) postJSONAndWait(uri string, data any) (bool, error) {
|
|||
return true, nil
|
||||
}
|
||||
|
||||
return false, c.waitForTickResponse(cmd.Ticks)
|
||||
ch := make(chan easee.SignalRCommandResponse, 1)
|
||||
c.registerPendingTick(cmd.Ticks, ch)
|
||||
defer c.unregisterPendingTick(cmd.Ticks)
|
||||
obsID := easee.ObservationID(cmd.CommandId)
|
||||
c.registerPendingByID(obsID, ch)
|
||||
defer c.unregisterPendingByID(obsID)
|
||||
return false, c.waitForTickResponse(ch)
|
||||
}
|
||||
|
||||
// all other response codes lead to an error
|
||||
return false, fmt.Errorf("invalid status: %d", resp.StatusCode)
|
||||
}
|
||||
|
||||
func (c *Easee) waitForTickResponse(expectedTick int64) error {
|
||||
for {
|
||||
select {
|
||||
case cmdResp := <-c.cmdC:
|
||||
if cmdResp.Ticks == expectedTick {
|
||||
if !cmdResp.WasAccepted {
|
||||
return fmt.Errorf("command rejected: %d", cmdResp.Ticks)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
case <-time.After(c.Client.Timeout):
|
||||
return api.ErrTimeout
|
||||
func (c *Easee) waitForTickResponse(ch <-chan easee.SignalRCommandResponse) error {
|
||||
select {
|
||||
case cmdResp := <-ch:
|
||||
if !cmdResp.WasAccepted {
|
||||
return fmt.Errorf("command rejected: %d", cmdResp.Ticks)
|
||||
}
|
||||
return nil
|
||||
case <-time.After(c.Client.Timeout):
|
||||
return api.ErrTimeout
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -744,7 +815,15 @@ func (c *Easee) Phases1p3p(phases int) error {
|
|||
data.DynamicCircuitCurrentP3 = &max3
|
||||
}
|
||||
|
||||
_, err = c.postJSONAndWait(uri, data)
|
||||
// Register before POST so the SignalR CommandResponse that the Easee
|
||||
// cloud sends on HTTP 200 (sync) responses is silently consumed rather
|
||||
// than logged as rogue. On error we undo the registration; on 202/noop
|
||||
// paths no extra CommandResponse is expected, so the counter is
|
||||
// intentionally left to be consumed by the eventual SignalR message.
|
||||
c.registerExpectedOrphan(easee.CIRCUIT_MAX_CURRENT_P1)
|
||||
if _, err = c.postJSONAndWait(uri, data); err != nil {
|
||||
c.consumeExpectedOrphan(easee.CIRCUIT_MAX_CURRENT_P1)
|
||||
}
|
||||
} else {
|
||||
// charger level
|
||||
if phases == 3 {
|
||||
|
|
|
|||
|
|
@ -13,6 +13,7 @@ import (
|
|||
"github.com/evcc-io/evcc/util/request"
|
||||
"github.com/jarcoal/httpmock"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
// Helper function to create a payload
|
||||
|
|
@ -30,12 +31,14 @@ func createPayload(id easee.ObservationID, timestamp time.Time, dataType easee.D
|
|||
func newEasee() *Easee {
|
||||
log := util.NewLogger("easee")
|
||||
e := Easee{
|
||||
Helper: request.NewHelper(log),
|
||||
obsTime: make(map[easee.ObservationID]time.Time),
|
||||
log: log,
|
||||
startDone: func() {},
|
||||
cmdC: make(chan easee.SignalRCommandResponse),
|
||||
obsC: make(chan easee.Observation),
|
||||
Helper: request.NewHelper(log),
|
||||
obsTime: make(map[easee.ObservationID]time.Time),
|
||||
pendingTicks: make(map[int64]chan easee.SignalRCommandResponse),
|
||||
pendingByID: make(map[easee.ObservationID]chan easee.SignalRCommandResponse),
|
||||
expectedOrphans: make(map[easee.ObservationID]int),
|
||||
log: log,
|
||||
startDone: func() {},
|
||||
obsC: make(chan easee.Observation),
|
||||
}
|
||||
e.Client.Timeout = 500 * time.Millisecond //aggressive timeout to accelerate testing
|
||||
return &e
|
||||
|
|
@ -166,15 +169,12 @@ func TestEasee_waitForTickResponse(t *testing.T) {
|
|||
|
||||
e := newEasee()
|
||||
|
||||
// Set up the command channel for test and Easee to share
|
||||
cmdC := make(chan easee.SignalRCommandResponse, 1) // make it buffered for ease of testing
|
||||
e.cmdC = cmdC
|
||||
|
||||
ch := make(chan easee.SignalRCommandResponse, 1)
|
||||
if tc.cmdCValue != nil {
|
||||
cmdC <- *tc.cmdCValue
|
||||
ch <- *tc.cmdCValue
|
||||
}
|
||||
|
||||
err := e.waitForTickResponse(tc.expectedTick)
|
||||
err := e.waitForTickResponse(ch)
|
||||
|
||||
// Assert the result
|
||||
if tc.expectedErr != nil {
|
||||
|
|
@ -227,7 +227,18 @@ func TestEasee_postJsonAndWait(t *testing.T) {
|
|||
|
||||
if tc.cmdResp != nil {
|
||||
go func() {
|
||||
e.cmdC <- *tc.cmdResp
|
||||
// wait for postJSONAndWait to register the per-tick channel
|
||||
var ch chan easee.SignalRCommandResponse
|
||||
for {
|
||||
e.cmdMu.Lock()
|
||||
ch = e.pendingTicks[tc.cmdResp.Ticks]
|
||||
e.cmdMu.Unlock()
|
||||
if ch != nil {
|
||||
break
|
||||
}
|
||||
time.Sleep(time.Millisecond)
|
||||
}
|
||||
ch <- *tc.cmdResp
|
||||
}()
|
||||
}
|
||||
|
||||
|
|
@ -380,3 +391,212 @@ func TestEasee_MaxCurrent(t *testing.T) {
|
|||
//assert.Equal(t, e.current, e.dynamicChargerCurrent)
|
||||
}
|
||||
}
|
||||
|
||||
func TestEasee_CommandResponse_rogue(t *testing.T) {
|
||||
e := newEasee()
|
||||
|
||||
rogueResp := easee.SignalRCommandResponse{
|
||||
SerialNumber: "EH123456",
|
||||
Ticks: 999999999,
|
||||
WasAccepted: true,
|
||||
ResultCode: 0,
|
||||
}
|
||||
|
||||
raw, err := json.Marshal(rogueResp)
|
||||
require.NoError(t, err)
|
||||
|
||||
// No pending tick registered → should log WARN (not panic, not block)
|
||||
assert.NotPanics(t, func() {
|
||||
e.CommandResponse(raw)
|
||||
})
|
||||
|
||||
// pendingTicks should still be empty
|
||||
e.cmdMu.Lock()
|
||||
assert.Empty(t, e.pendingTicks)
|
||||
e.cmdMu.Unlock()
|
||||
}
|
||||
|
||||
func TestEasee_CommandResponse_legitimate(t *testing.T) {
|
||||
e := newEasee()
|
||||
|
||||
ticks := int64(638798974487432600)
|
||||
ch := make(chan easee.SignalRCommandResponse, 1)
|
||||
e.registerPendingTick(ticks, ch)
|
||||
|
||||
resp := easee.SignalRCommandResponse{
|
||||
SerialNumber: "EH123456",
|
||||
Ticks: ticks,
|
||||
WasAccepted: true,
|
||||
}
|
||||
|
||||
raw, err := json.Marshal(resp)
|
||||
require.NoError(t, err)
|
||||
|
||||
e.CommandResponse(raw)
|
||||
|
||||
// Channel should have received the response
|
||||
select {
|
||||
case got := <-ch:
|
||||
assert.Equal(t, ticks, got.Ticks)
|
||||
assert.True(t, got.WasAccepted)
|
||||
case <-time.After(100 * time.Millisecond):
|
||||
t.Fatal("CommandResponse did not deliver to pending channel")
|
||||
}
|
||||
}
|
||||
|
||||
func TestEasee_CommandResponse_expectedOrphan(t *testing.T) {
|
||||
e := newEasee()
|
||||
|
||||
// Pre-register the expected orphan
|
||||
e.registerExpectedOrphan(easee.CIRCUIT_MAX_CURRENT_P1)
|
||||
|
||||
resp := easee.SignalRCommandResponse{
|
||||
SerialNumber: "EH123456",
|
||||
ID: int(easee.CIRCUIT_MAX_CURRENT_P1),
|
||||
Ticks: 111111111,
|
||||
WasAccepted: true,
|
||||
ResultCode: 0,
|
||||
}
|
||||
|
||||
raw, err := json.Marshal(resp)
|
||||
require.NoError(t, err)
|
||||
|
||||
// Should not panic and should consume the orphan counter
|
||||
assert.NotPanics(t, func() {
|
||||
e.CommandResponse(raw)
|
||||
})
|
||||
|
||||
// Counter should now be zero — a second response would be rogue
|
||||
assert.False(t, e.consumeExpectedOrphan(easee.CIRCUIT_MAX_CURRENT_P1))
|
||||
}
|
||||
|
||||
func TestEasee_CommandResponse_rogueAfterOrphanConsumed(t *testing.T) {
|
||||
e := newEasee()
|
||||
|
||||
// Register and immediately consume via CommandResponse
|
||||
e.registerExpectedOrphan(easee.CIRCUIT_MAX_CURRENT_P1)
|
||||
|
||||
resp := easee.SignalRCommandResponse{
|
||||
SerialNumber: "EH123456",
|
||||
ID: int(easee.CIRCUIT_MAX_CURRENT_P1),
|
||||
Ticks: 111111111,
|
||||
WasAccepted: true,
|
||||
}
|
||||
raw, err := json.Marshal(resp)
|
||||
require.NoError(t, err)
|
||||
e.CommandResponse(raw) // consumes the counter
|
||||
|
||||
// A second identical response with counter=0 should be treated as rogue (not panic)
|
||||
assert.NotPanics(t, func() {
|
||||
e.CommandResponse(raw)
|
||||
})
|
||||
|
||||
// pendingTicks untouched
|
||||
e.cmdMu.Lock()
|
||||
assert.Empty(t, e.pendingTicks)
|
||||
e.cmdMu.Unlock()
|
||||
}
|
||||
|
||||
func TestEasee_CommandResponse_matchedByID(t *testing.T) {
|
||||
e := newEasee()
|
||||
|
||||
ch := make(chan easee.SignalRCommandResponse, 1)
|
||||
e.registerPendingByID(easee.LOCATION, ch)
|
||||
defer e.unregisterPendingByID(easee.LOCATION)
|
||||
|
||||
// Ticks do NOT match any pendingTicks entry — only the ID matches
|
||||
resp := easee.SignalRCommandResponse{
|
||||
SerialNumber: "EH123456",
|
||||
ID: int(easee.LOCATION),
|
||||
Ticks: 999999999, // not in pendingTicks
|
||||
WasAccepted: true,
|
||||
ResultCode: 0,
|
||||
}
|
||||
|
||||
raw, err := json.Marshal(resp)
|
||||
require.NoError(t, err)
|
||||
|
||||
assert.NotPanics(t, func() {
|
||||
e.CommandResponse(raw)
|
||||
})
|
||||
|
||||
select {
|
||||
case got := <-ch:
|
||||
assert.Equal(t, resp.Ticks, got.Ticks)
|
||||
assert.True(t, got.WasAccepted)
|
||||
default:
|
||||
t.Fatal("expected CommandResponse to be delivered to pendingByID channel")
|
||||
}
|
||||
|
||||
// pendingByID consumed the response — expectedOrphans untouched
|
||||
assert.False(t, e.consumeExpectedOrphan(easee.LOCATION))
|
||||
}
|
||||
|
||||
func TestEasee_registerAndConsumeExpectedOrphan(t *testing.T) {
|
||||
e := newEasee()
|
||||
|
||||
// Not registered yet — consume returns false
|
||||
assert.False(t, e.consumeExpectedOrphan(easee.CIRCUIT_MAX_CURRENT_P1))
|
||||
|
||||
// Register once
|
||||
e.registerExpectedOrphan(easee.CIRCUIT_MAX_CURRENT_P1)
|
||||
|
||||
// First consume succeeds
|
||||
assert.True(t, e.consumeExpectedOrphan(easee.CIRCUIT_MAX_CURRENT_P1))
|
||||
|
||||
// Second consume fails (counter back to zero)
|
||||
assert.False(t, e.consumeExpectedOrphan(easee.CIRCUIT_MAX_CURRENT_P1))
|
||||
}
|
||||
|
||||
func TestEasee_registerExpectedOrphan_multipleRegistrations(t *testing.T) {
|
||||
e := newEasee()
|
||||
|
||||
// Register twice (two concurrent calls in flight)
|
||||
e.registerExpectedOrphan(easee.CIRCUIT_MAX_CURRENT_P1)
|
||||
e.registerExpectedOrphan(easee.CIRCUIT_MAX_CURRENT_P1)
|
||||
|
||||
assert.True(t, e.consumeExpectedOrphan(easee.CIRCUIT_MAX_CURRENT_P1))
|
||||
assert.True(t, e.consumeExpectedOrphan(easee.CIRCUIT_MAX_CURRENT_P1))
|
||||
assert.False(t, e.consumeExpectedOrphan(easee.CIRCUIT_MAX_CURRENT_P1))
|
||||
}
|
||||
|
||||
func TestEasee_Phases1p3p_registersExpectedOrphan(t *testing.T) {
|
||||
const siteID = 12345
|
||||
const circuitID = 67890
|
||||
const chargerID = "TESTTEST"
|
||||
|
||||
e := newEasee()
|
||||
e.charger = chargerID
|
||||
e.site = siteID
|
||||
e.circuit = circuitID
|
||||
|
||||
httpmock.ActivateNonDefault(e.Client)
|
||||
defer httpmock.DeactivateAndReset()
|
||||
|
||||
// Mock GET circuit settings
|
||||
getURI := fmt.Sprintf("%s/sites/%d/circuits/%d/settings", easee.API, siteID, circuitID)
|
||||
maxP1, maxP2, maxP3 := 32.0, 32.0, 32.0
|
||||
getResp := easee.CircuitSettings{
|
||||
MaxCircuitCurrentP1: &maxP1,
|
||||
MaxCircuitCurrentP2: &maxP2,
|
||||
MaxCircuitCurrentP3: &maxP3,
|
||||
}
|
||||
body, err := json.Marshal(getResp)
|
||||
require.NoError(t, err)
|
||||
httpmock.RegisterResponder(http.MethodGet, getURI,
|
||||
httpmock.NewBytesResponder(200, body))
|
||||
|
||||
// Mock POST circuit settings — return 200 (sync)
|
||||
httpmock.RegisterResponder(http.MethodPost, getURI,
|
||||
httpmock.NewStringResponder(200, ""))
|
||||
|
||||
err = e.Phases1p3p(1)
|
||||
assert.NoError(t, err)
|
||||
|
||||
// The orphan counter should have been registered before the POST.
|
||||
// Since no CommandResponse arrived in this test, the counter stays at 1.
|
||||
e.cmdMu.Lock()
|
||||
count := e.expectedOrphans[easee.CIRCUIT_MAX_CURRENT_P1]
|
||||
e.cmdMu.Unlock()
|
||||
assert.Equal(t, 1, count, "expected orphan should be registered before the POST")
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue