diff --git a/charger/easee.go b/charger/easee.go index 4929e564b..9b32e4e50 100644 --- a/charger/easee.go +++ b/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 { diff --git a/charger/easee_test.go b/charger/easee_test.go index 30c543c71..3c6b19943 100644 --- a/charger/easee_test.go +++ b/charger/easee_test.go @@ -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") +}