Easee: extract CommandDispatcher for async command correlation (#28036)
This commit is contained in:
parent
452a1301b1
commit
5877d0acd5
4 changed files with 525 additions and 386 deletions
192
charger/easee.go
192
charger/easee.go
|
|
@ -24,7 +24,6 @@ import (
|
|||
"fmt"
|
||||
"net/http"
|
||||
"os"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
|
|
@ -68,12 +67,11 @@ type Easee struct {
|
|||
phaseMode int
|
||||
currentPower, sessionEnergy, totalEnergy,
|
||||
currentL1, currentL2, currentL3 float64
|
||||
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
|
||||
rfid string
|
||||
lp loadpoint.API
|
||||
|
||||
dispatcher *easee.CommandDispatcher
|
||||
|
||||
obsC chan easee.Observation
|
||||
obsTime map[easee.ObservationID]time.Time
|
||||
lastObsReceived time.Time
|
||||
|
|
@ -118,21 +116,20 @@ 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) }),
|
||||
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),
|
||||
Helper: request.NewHelper(log),
|
||||
charger: charger,
|
||||
authorize: authorize,
|
||||
log: log,
|
||||
current: 6, // default current
|
||||
startDone: sync.OnceFunc(func() { close(done) }),
|
||||
obsC: make(chan easee.Observation),
|
||||
obsTime: make(map[easee.ObservationID]time.Time),
|
||||
}
|
||||
|
||||
c.Client.Timeout = timeout
|
||||
|
||||
c.dispatcher = easee.NewCommandDispatcher(c.Helper, log, timeout)
|
||||
|
||||
ts, err := easee.TokenSource(log, user, password)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
|
|
@ -224,48 +221,6 @@ 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()
|
||||
|
|
@ -439,39 +394,12 @@ func (c *Easee) SubscribeToMyProduct(i json.RawMessage) {
|
|||
// CommandResponse implements the signalr receiver
|
||||
func (c *Easee) CommandResponse(i json.RawMessage) {
|
||||
var res easee.SignalRCommandResponse
|
||||
|
||||
if err := json.Unmarshal(i, &res); err != nil {
|
||||
c.log.ERROR.Printf("invalid message: %s %v", i, err)
|
||||
return
|
||||
}
|
||||
|
||||
c.log.TRACE.Printf("CommandResponse %s: %+v", res.SerialNumber, res)
|
||||
|
||||
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)
|
||||
c.dispatcher.Dispatch(res)
|
||||
}
|
||||
|
||||
func (c *Easee) chargers() ([]easee.Charger, error) {
|
||||
|
|
@ -532,7 +460,7 @@ func (c *Easee) Enable(enable bool) (err error) {
|
|||
}
|
||||
|
||||
uri := fmt.Sprintf("%s/chargers/%s/settings", easee.API, c.charger)
|
||||
if _, err := c.postJSONAndWait(uri, data); err != nil {
|
||||
if err := c.dispatcher.Send(uri, data); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
|
@ -554,7 +482,7 @@ func (c *Easee) Enable(enable bool) (err error) {
|
|||
}
|
||||
|
||||
uri := fmt.Sprintf("%s/chargers/%s/commands/%s", easee.API, c.charger, action)
|
||||
if _, err := c.postJSONAndWait(uri, nil); err != nil {
|
||||
if err := c.dispatcher.Send(uri, nil); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
|
|
@ -594,65 +522,6 @@ func (c *Easee) inExpectedOpMode(enable bool) bool {
|
|||
return c.opMode == easee.ModeAwaitingStart || c.opMode == easee.ModeAwaitingAuthentication
|
||||
}
|
||||
|
||||
// posts JSON to the Easee API endpoint and waits for the async response
|
||||
func (c *Easee) postJSONAndWait(uri string, data any) (bool, error) {
|
||||
resp, err := c.Post(uri, request.JSONContent, request.MarshalJSON(data))
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
|
||||
if resp.StatusCode == 200 { // sync call
|
||||
return false, nil
|
||||
}
|
||||
|
||||
if resp.StatusCode == 202 { // async call, wait for response
|
||||
var cmd easee.RestCommandResponse
|
||||
|
||||
if strings.Contains(uri, "/commands/") { // command endpoint
|
||||
if err := json.NewDecoder(resp.Body).Decode(&cmd); err != nil {
|
||||
return false, err
|
||||
}
|
||||
} else { // settings endpoint
|
||||
var cmdArr []easee.RestCommandResponse
|
||||
if err := json.NewDecoder(resp.Body).Decode(&cmdArr); err != nil {
|
||||
return false, err
|
||||
}
|
||||
|
||||
if len(cmdArr) != 0 {
|
||||
cmd = cmdArr[0]
|
||||
}
|
||||
}
|
||||
|
||||
if cmd.Ticks == 0 { // api thinks this was a noop
|
||||
return true, nil
|
||||
}
|
||||
|
||||
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(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
|
||||
}
|
||||
}
|
||||
|
||||
// wait for opMode become expected op mode
|
||||
func (c *Easee) waitForChargerEnabledState(expEnabled bool) error {
|
||||
// check any updates received meanwhile
|
||||
|
|
@ -725,15 +594,11 @@ func (c *Easee) MaxCurrent(current int64) error {
|
|||
}
|
||||
|
||||
uri := fmt.Sprintf("%s/chargers/%s/settings", easee.API, c.charger)
|
||||
noop, err := c.postJSONAndWait(uri, data)
|
||||
if err != nil {
|
||||
if err := c.dispatcher.Send(uri, data); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
if !noop {
|
||||
if err := c.waitForDynamicChargerCurrent(float64(current)); err != nil {
|
||||
return err
|
||||
}
|
||||
if err := c.waitForDynamicChargerCurrent(cur); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
c.mux.Lock()
|
||||
|
|
@ -840,12 +705,11 @@ func (c *Easee) Phases1p3p(phases int) error {
|
|||
|
||||
// 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)
|
||||
// than logged as rogue. On error we undo the registration.
|
||||
c.dispatcher.ExpectOrphan(easee.CIRCUIT_MAX_CURRENT_P1)
|
||||
err = c.dispatcher.Send(uri, data)
|
||||
if err != nil {
|
||||
c.dispatcher.CancelOrphan(easee.CIRCUIT_MAX_CURRENT_P1)
|
||||
}
|
||||
} else {
|
||||
// charger level
|
||||
|
|
@ -861,7 +725,7 @@ func (c *Easee) Phases1p3p(phases int) error {
|
|||
|
||||
uri := fmt.Sprintf("%s/chargers/%s/settings", easee.API, c.charger)
|
||||
|
||||
if _, err = c.postJSONAndWait(uri, data); err != nil {
|
||||
if err = c.dispatcher.Send(uri, data); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
|
@ -941,7 +805,7 @@ func (c *Easee) updateSmartCharging() {
|
|||
|
||||
uri := fmt.Sprintf("%s/chargers/%s/settings", easee.API, c.charger)
|
||||
|
||||
if _, err := c.postJSONAndWait(uri, data); err != nil {
|
||||
if err := c.dispatcher.Send(uri, data); err != nil {
|
||||
c.log.WARN.Printf("smart charging: %v", err)
|
||||
return
|
||||
}
|
||||
|
|
|
|||
186
charger/easee/dispatcher.go
Normal file
186
charger/easee/dispatcher.go
Normal file
|
|
@ -0,0 +1,186 @@
|
|||
package easee
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/evcc-io/evcc/api"
|
||||
"github.com/evcc-io/evcc/util"
|
||||
"github.com/evcc-io/evcc/util/request"
|
||||
)
|
||||
|
||||
// CommandDispatcher owns the full lifecycle of an Easee command:
|
||||
// HTTP POST → response parsing → SignalR CommandResponse correlation.
|
||||
type CommandDispatcher struct {
|
||||
helper *request.Helper
|
||||
mu sync.Mutex
|
||||
pendingTicks map[int64]chan SignalRCommandResponse
|
||||
pendingByID map[ObservationID]chan SignalRCommandResponse
|
||||
expectedOrphans map[ObservationID]int
|
||||
log *util.Logger
|
||||
timeout time.Duration
|
||||
}
|
||||
|
||||
// NewCommandDispatcher creates a dispatcher. helper must be the authenticated
|
||||
// HTTP client used for all Easee API calls.
|
||||
func NewCommandDispatcher(helper *request.Helper, log *util.Logger, timeout time.Duration) *CommandDispatcher {
|
||||
return &CommandDispatcher{
|
||||
helper: helper,
|
||||
log: log,
|
||||
timeout: timeout,
|
||||
pendingTicks: make(map[int64]chan SignalRCommandResponse),
|
||||
pendingByID: make(map[ObservationID]chan SignalRCommandResponse),
|
||||
expectedOrphans: make(map[ObservationID]int),
|
||||
}
|
||||
}
|
||||
|
||||
// Dispatch routes an incoming CommandResponse to the appropriate waiter.
|
||||
// Must be called from the Easee.CommandResponse SignalR handler.
|
||||
// Logs a WARN if no pending registration or expected orphan matches.
|
||||
func (d *CommandDispatcher) Dispatch(res SignalRCommandResponse) {
|
||||
obsID := ObservationID(res.ID)
|
||||
|
||||
d.mu.Lock()
|
||||
chTick, tickOk := d.pendingTicks[res.Ticks]
|
||||
chID, idOk := d.pendingByID[obsID]
|
||||
d.mu.Unlock()
|
||||
|
||||
// Tick lookup takes priority over ID lookup (primary correlation).
|
||||
// ID lookup is a fallback for backend clock drift / load balancer
|
||||
// scenarios where the delivered Ticks differs from the HTTP 202 body.
|
||||
if tickOk {
|
||||
// Channels are buffered (capacity 1) — this send never blocks even if
|
||||
// the waiter has timed out and unregistered the channel already.
|
||||
chTick <- res
|
||||
return
|
||||
}
|
||||
|
||||
if idOk {
|
||||
chID <- res // buffered (capacity 1), see comment above
|
||||
return
|
||||
}
|
||||
|
||||
if d.CancelOrphan(obsID) {
|
||||
return
|
||||
}
|
||||
|
||||
d.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)
|
||||
}
|
||||
|
||||
// ExpectOrphan pre-registers one expected CommandResponse per id for a
|
||||
// sync (HTTP 200) endpoint that still produces a CommandResponse on the wire.
|
||||
// Must be called before Send to avoid a race with the arriving CommandResponse.
|
||||
func (d *CommandDispatcher) ExpectOrphan(ids ...ObservationID) {
|
||||
d.mu.Lock()
|
||||
defer d.mu.Unlock()
|
||||
for _, id := range ids {
|
||||
d.expectedOrphans[id]++
|
||||
}
|
||||
}
|
||||
|
||||
// CancelOrphan decrements the expected-orphan counter for id.
|
||||
// Returns true if a counter entry was consumed, false if none existed.
|
||||
// Used by call sites to undo an ExpectOrphan registration when the POST fails.
|
||||
func (d *CommandDispatcher) CancelOrphan(id ObservationID) bool {
|
||||
d.mu.Lock()
|
||||
defer d.mu.Unlock()
|
||||
if d.expectedOrphans[id] > 0 {
|
||||
d.expectedOrphans[id]--
|
||||
if d.expectedOrphans[id] == 0 {
|
||||
delete(d.expectedOrphans, id)
|
||||
}
|
||||
return true
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
// Send posts to uri with data, parses the Easee-specific response body, and
|
||||
// if the response is asynchronous (HTTP 202), waits for the matching SignalR
|
||||
// CommandResponse.
|
||||
//
|
||||
// Returns nil on success (both synchronous HTTP 200 and confirmed async HTTP 202,
|
||||
// including noops where Ticks == 0). Returns an error on HTTP failure, decode
|
||||
// failure, command rejection, or timeout.
|
||||
func (d *CommandDispatcher) Send(uri string, data any) error {
|
||||
resp, err := d.helper.Post(uri, request.JSONContent, request.MarshalJSON(data))
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
|
||||
if resp.StatusCode == 200 {
|
||||
return nil
|
||||
}
|
||||
|
||||
// Any status other than 200 or 202 is unexpected — return an error.
|
||||
// Note: http.Client.Post only errors on transport failures (DNS, TLS, etc.),
|
||||
// not on HTTP error responses, so this guard is the actual defense against
|
||||
// 4xx/5xx responses from the Easee API.
|
||||
if resp.StatusCode != 202 {
|
||||
return fmt.Errorf("unexpected status: %d", resp.StatusCode)
|
||||
}
|
||||
|
||||
// HTTP 202: parse the response body to get the command correlation info.
|
||||
var cmd RestCommandResponse
|
||||
if strings.Contains(uri, "/commands/") {
|
||||
// Command endpoints return a single object.
|
||||
if err := json.NewDecoder(resp.Body).Decode(&cmd); err != nil {
|
||||
return err
|
||||
}
|
||||
} else {
|
||||
// Settings endpoints return an array; take index 0 if present.
|
||||
var cmdArr []RestCommandResponse
|
||||
if err := json.NewDecoder(resp.Body).Decode(&cmdArr); err != nil {
|
||||
return err
|
||||
}
|
||||
if len(cmdArr) != 0 {
|
||||
cmd = cmdArr[0]
|
||||
for _, extra := range cmdArr[1:] {
|
||||
d.log.TRACE.Printf("ignoring additional CommandResponse in settings reply: %+v", extra)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if cmd.Ticks == 0 {
|
||||
// Noop: the API indicates no state change was needed.
|
||||
return nil
|
||||
}
|
||||
|
||||
// Create a buffered channel (capacity 1) so Dispatch never blocks even if
|
||||
// Send has already returned due to timeout.
|
||||
ch := make(chan SignalRCommandResponse, 1)
|
||||
|
||||
d.mu.Lock()
|
||||
d.pendingTicks[cmd.Ticks] = ch
|
||||
// Note: if two concurrent Send calls share the same ObservationID, the
|
||||
// second would overwrite the first's pendingByID entry. In practice this
|
||||
// cannot occur because the loadpoint serializes Enable/MaxCurrent calls.
|
||||
d.pendingByID[ObservationID(cmd.CommandId)] = ch
|
||||
d.mu.Unlock()
|
||||
|
||||
defer func() {
|
||||
d.mu.Lock()
|
||||
delete(d.pendingTicks, cmd.Ticks)
|
||||
delete(d.pendingByID, ObservationID(cmd.CommandId))
|
||||
d.mu.Unlock()
|
||||
}()
|
||||
|
||||
timer := time.NewTimer(d.timeout)
|
||||
defer timer.Stop()
|
||||
|
||||
select {
|
||||
case res := <-ch:
|
||||
if !res.WasAccepted {
|
||||
return fmt.Errorf("command rejected: %d", res.Ticks)
|
||||
}
|
||||
return nil
|
||||
case <-timer.C:
|
||||
return api.ErrTimeout
|
||||
}
|
||||
}
|
||||
247
charger/easee/dispatcher_test.go
Normal file
247
charger/easee/dispatcher_test.go
Normal file
|
|
@ -0,0 +1,247 @@
|
|||
package easee
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"net/http"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/evcc-io/evcc/api"
|
||||
"github.com/evcc-io/evcc/util"
|
||||
"github.com/evcc-io/evcc/util/request"
|
||||
"github.com/jarcoal/httpmock"
|
||||
"github.com/stretchr/testify/assert"
|
||||
)
|
||||
|
||||
const (
|
||||
testURI = API + "/chargers/TESTTEST/settings"
|
||||
testCmdURI = API + "/chargers/TESTTEST/commands/resume_charging"
|
||||
)
|
||||
|
||||
func newTestDispatcher(t *testing.T) *CommandDispatcher {
|
||||
t.Helper()
|
||||
log := util.NewLogger("test")
|
||||
h := request.NewHelper(log)
|
||||
h.Client.Timeout = 500 * time.Millisecond
|
||||
return NewCommandDispatcher(h, log, 500*time.Millisecond)
|
||||
}
|
||||
|
||||
// waitForPendingTick blocks until d.pendingTicks contains ticks, or the test fails.
|
||||
func waitForPendingTick(t *testing.T, d *CommandDispatcher, ticks int64) {
|
||||
t.Helper()
|
||||
deadline := time.After(2 * time.Second)
|
||||
for {
|
||||
select {
|
||||
case <-deadline:
|
||||
t.Fatalf("timed out waiting for pendingTicks registration for ticks %d", ticks)
|
||||
default:
|
||||
}
|
||||
d.mu.Lock()
|
||||
_, ok := d.pendingTicks[ticks]
|
||||
d.mu.Unlock()
|
||||
if ok {
|
||||
return
|
||||
}
|
||||
time.Sleep(time.Millisecond)
|
||||
}
|
||||
}
|
||||
|
||||
func TestDispatcher_Dispatch_Rogue(t *testing.T) {
|
||||
// Intentionally triggers a WARN — suppress it to keep test output clean.
|
||||
util.LogLevel("error", nil)
|
||||
t.Cleanup(func() { util.LogLevel("info", nil) })
|
||||
|
||||
d := newTestDispatcher(t)
|
||||
assert.NotPanics(t, func() {
|
||||
d.Dispatch(SignalRCommandResponse{
|
||||
SerialNumber: "EH123456",
|
||||
Ticks: 999999999,
|
||||
WasAccepted: true,
|
||||
})
|
||||
})
|
||||
}
|
||||
|
||||
func TestDispatcher_Dispatch_ExpectedOrphan(t *testing.T) {
|
||||
d := newTestDispatcher(t)
|
||||
d.ExpectOrphan(CIRCUIT_MAX_CURRENT_P1)
|
||||
|
||||
assert.NotPanics(t, func() {
|
||||
d.Dispatch(SignalRCommandResponse{
|
||||
ID: int(CIRCUIT_MAX_CURRENT_P1),
|
||||
Ticks: 111111111,
|
||||
WasAccepted: true,
|
||||
})
|
||||
})
|
||||
|
||||
// Counter consumed — a second call to CancelOrphan returns false
|
||||
assert.False(t, d.CancelOrphan(CIRCUIT_MAX_CURRENT_P1))
|
||||
}
|
||||
|
||||
func TestDispatcher_CancelOrphan_Rollback(t *testing.T) {
|
||||
d := newTestDispatcher(t)
|
||||
d.ExpectOrphan(CIRCUIT_MAX_CURRENT_P1)
|
||||
assert.True(t, d.CancelOrphan(CIRCUIT_MAX_CURRENT_P1))
|
||||
assert.False(t, d.CancelOrphan(CIRCUIT_MAX_CURRENT_P1))
|
||||
}
|
||||
|
||||
func TestDispatcher_CancelOrphan_DoubleConsume(t *testing.T) {
|
||||
d := newTestDispatcher(t)
|
||||
d.ExpectOrphan(CIRCUIT_MAX_CURRENT_P1)
|
||||
// Dispatch consumes the orphan counter
|
||||
d.Dispatch(SignalRCommandResponse{ID: int(CIRCUIT_MAX_CURRENT_P1), Ticks: 111})
|
||||
// CancelOrphan now finds nothing
|
||||
assert.False(t, d.CancelOrphan(CIRCUIT_MAX_CURRENT_P1))
|
||||
}
|
||||
|
||||
// --- Send tests ---
|
||||
|
||||
func TestDispatcher_Send_HTTP200Sync(t *testing.T) {
|
||||
d := newTestDispatcher(t)
|
||||
httpmock.ActivateNonDefault(d.helper.Client)
|
||||
// per-client mock; no global teardown needed
|
||||
|
||||
httpmock.RegisterResponder(http.MethodPost, testURI,
|
||||
httpmock.NewStringResponder(http.StatusOK, ""))
|
||||
|
||||
assert.NoError(t, d.Send(testURI, nil))
|
||||
}
|
||||
|
||||
func TestDispatcher_Send_Noop(t *testing.T) {
|
||||
d := newTestDispatcher(t)
|
||||
httpmock.ActivateNonDefault(d.helper.Client)
|
||||
// per-client mock; no global teardown needed
|
||||
|
||||
// Empty array body → Ticks == 0 → noop
|
||||
httpmock.RegisterResponder(http.MethodPost, testURI,
|
||||
httpmock.NewStringResponder(http.StatusAccepted, "[]"))
|
||||
|
||||
assert.NoError(t, d.Send(testURI, nil))
|
||||
}
|
||||
|
||||
func TestDispatcher_Send_HTTP202_InvalidJSON(t *testing.T) {
|
||||
d := newTestDispatcher(t)
|
||||
httpmock.ActivateNonDefault(d.helper.Client)
|
||||
|
||||
httpmock.RegisterResponder(http.MethodPost, testURI,
|
||||
httpmock.NewStringResponder(http.StatusAccepted, "{not-json"))
|
||||
|
||||
assert.Error(t, d.Send(testURI, nil))
|
||||
}
|
||||
|
||||
func TestDispatcher_Send_HTTP202_NonNumericTicks(t *testing.T) {
|
||||
d := newTestDispatcher(t)
|
||||
httpmock.ActivateNonDefault(d.helper.Client)
|
||||
|
||||
httpmock.RegisterResponder(http.MethodPost, testURI,
|
||||
httpmock.NewStringResponder(http.StatusAccepted, `[{"ticks":"NaN"}]`))
|
||||
|
||||
assert.Error(t, d.Send(testURI, nil))
|
||||
}
|
||||
|
||||
func TestDispatcher_Send_HTTPError(t *testing.T) {
|
||||
d := newTestDispatcher(t)
|
||||
httpmock.ActivateNonDefault(d.helper.Client)
|
||||
// per-client mock; no global teardown needed
|
||||
|
||||
httpmock.RegisterResponder(http.MethodPost, testURI,
|
||||
httpmock.NewStringResponder(http.StatusBadRequest, ""))
|
||||
|
||||
err := d.Send(testURI, nil)
|
||||
assert.Error(t, err)
|
||||
}
|
||||
|
||||
func TestDispatcher_Send_TicksMatch(t *testing.T) {
|
||||
const ticks int64 = 638798974487432600
|
||||
|
||||
d := newTestDispatcher(t)
|
||||
httpmock.ActivateNonDefault(d.helper.Client)
|
||||
// per-client mock; no global teardown needed
|
||||
|
||||
body := fmt.Sprintf(`[{"device":"TESTTEST","commandId":48,"ticks":%d}]`, ticks)
|
||||
httpmock.RegisterResponder(http.MethodPost, testURI,
|
||||
httpmock.NewStringResponder(http.StatusAccepted, body))
|
||||
|
||||
go func() {
|
||||
waitForPendingTick(t, d, ticks)
|
||||
d.Dispatch(SignalRCommandResponse{Ticks: ticks, WasAccepted: true})
|
||||
}()
|
||||
|
||||
assert.NoError(t, d.Send(testURI, nil))
|
||||
}
|
||||
|
||||
func TestDispatcher_Send_IDFallback(t *testing.T) {
|
||||
const ticks int64 = 638798974487432600
|
||||
const obsID = DYNAMIC_CHARGER_CURRENT // ObservationID = 48
|
||||
|
||||
d := newTestDispatcher(t)
|
||||
httpmock.ActivateNonDefault(d.helper.Client)
|
||||
// per-client mock; no global teardown needed
|
||||
|
||||
body := fmt.Sprintf(`[{"device":"TESTTEST","commandId":%d,"ticks":%d}]`, int(obsID), ticks)
|
||||
httpmock.RegisterResponder(http.MethodPost, testURI,
|
||||
httpmock.NewStringResponder(http.StatusAccepted, body))
|
||||
|
||||
go func() {
|
||||
waitForPendingTick(t, d, ticks)
|
||||
// Wrong Ticks (T+1), correct ID — triggers the ID fallback path
|
||||
d.Dispatch(SignalRCommandResponse{ID: int(obsID), Ticks: ticks + 1, WasAccepted: true})
|
||||
}()
|
||||
|
||||
assert.NoError(t, d.Send(testURI, nil))
|
||||
}
|
||||
|
||||
func TestDispatcher_Send_Timeout(t *testing.T) {
|
||||
const ticks int64 = 789
|
||||
|
||||
d := newTestDispatcher(t)
|
||||
httpmock.ActivateNonDefault(d.helper.Client)
|
||||
// per-client mock; no global teardown needed
|
||||
|
||||
body := fmt.Sprintf(`[{"device":"TESTTEST","commandId":48,"ticks":%d}]`, ticks)
|
||||
httpmock.RegisterResponder(http.MethodPost, testURI,
|
||||
httpmock.NewStringResponder(http.StatusAccepted, body))
|
||||
|
||||
// No Dispatch call → Send times out
|
||||
assert.ErrorIs(t, d.Send(testURI, nil), api.ErrTimeout)
|
||||
}
|
||||
|
||||
func TestDispatcher_Send_Rejected(t *testing.T) {
|
||||
const ticks int64 = 456
|
||||
|
||||
d := newTestDispatcher(t)
|
||||
httpmock.ActivateNonDefault(d.helper.Client)
|
||||
// per-client mock; no global teardown needed
|
||||
|
||||
body := fmt.Sprintf(`[{"device":"TESTTEST","commandId":48,"ticks":%d}]`, ticks)
|
||||
httpmock.RegisterResponder(http.MethodPost, testURI,
|
||||
httpmock.NewStringResponder(http.StatusAccepted, body))
|
||||
|
||||
go func() {
|
||||
waitForPendingTick(t, d, ticks)
|
||||
d.Dispatch(SignalRCommandResponse{Ticks: ticks, WasAccepted: false})
|
||||
}()
|
||||
|
||||
err := d.Send(testURI, nil)
|
||||
assert.Error(t, err)
|
||||
assert.Contains(t, err.Error(), "rejected")
|
||||
}
|
||||
|
||||
func TestDispatcher_Send_CommandURI(t *testing.T) {
|
||||
const ticks int64 = 638798974487432600
|
||||
|
||||
d := newTestDispatcher(t)
|
||||
httpmock.ActivateNonDefault(d.helper.Client)
|
||||
// per-client mock; no global teardown needed
|
||||
|
||||
// /commands/ endpoint → body is a JSON object, not an array
|
||||
body := fmt.Sprintf(`{"device":"TESTTEST","commandId":48,"ticks":%d}`, ticks)
|
||||
httpmock.RegisterResponder(http.MethodPost, testCmdURI,
|
||||
httpmock.NewStringResponder(http.StatusAccepted, body))
|
||||
|
||||
go func() {
|
||||
waitForPendingTick(t, d, ticks)
|
||||
d.Dispatch(SignalRCommandResponse{Ticks: ticks, WasAccepted: true})
|
||||
}()
|
||||
|
||||
assert.NoError(t, d.Send(testCmdURI, nil))
|
||||
}
|
||||
|
|
@ -30,15 +30,14 @@ func createPayload(id easee.ObservationID, timestamp time.Time, dataType easee.D
|
|||
|
||||
func newEasee() *Easee {
|
||||
log := util.NewLogger("easee")
|
||||
helper := request.NewHelper(log)
|
||||
e := Easee{
|
||||
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),
|
||||
Helper: helper,
|
||||
obsTime: make(map[easee.ObservationID]time.Time),
|
||||
log: log,
|
||||
startDone: func() {},
|
||||
obsC: make(chan easee.Observation),
|
||||
dispatcher: easee.NewCommandDispatcher(helper, log, 500*time.Millisecond),
|
||||
}
|
||||
e.Client.Timeout = 500 * time.Millisecond //aggressive timeout to accelerate testing
|
||||
return &e
|
||||
|
|
@ -137,120 +136,6 @@ func TestInExpectedOpMode(t *testing.T) {
|
|||
}
|
||||
}
|
||||
|
||||
func TestEasee_waitForTickResponse(t *testing.T) {
|
||||
testCases := []struct {
|
||||
name string
|
||||
expectedTick int64
|
||||
cmdCValue *easee.SignalRCommandResponse
|
||||
expectedErr error
|
||||
}{
|
||||
{
|
||||
name: "Success - Tick Found",
|
||||
expectedTick: 123,
|
||||
cmdCValue: &easee.SignalRCommandResponse{Ticks: 123, WasAccepted: true},
|
||||
expectedErr: nil,
|
||||
},
|
||||
{
|
||||
name: "Success - Tick Found, but Rejected",
|
||||
expectedTick: 456,
|
||||
cmdCValue: &easee.SignalRCommandResponse{Ticks: 456, WasAccepted: false},
|
||||
expectedErr: fmt.Errorf("command rejected: %d", 456),
|
||||
},
|
||||
{
|
||||
name: "Timeout",
|
||||
expectedTick: 789,
|
||||
expectedErr: api.ErrTimeout,
|
||||
},
|
||||
}
|
||||
|
||||
for _, tc := range testCases {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
t.Logf("%+v", tc)
|
||||
|
||||
e := newEasee()
|
||||
|
||||
ch := make(chan easee.SignalRCommandResponse, 1)
|
||||
if tc.cmdCValue != nil {
|
||||
ch <- *tc.cmdCValue
|
||||
}
|
||||
|
||||
err := e.waitForTickResponse(ch)
|
||||
|
||||
// Assert the result
|
||||
if tc.expectedErr != nil {
|
||||
assert.EqualError(t, err, tc.expectedErr.Error())
|
||||
} else {
|
||||
assert.NoError(t, err)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestEasee_postJsonAndWait(t *testing.T) {
|
||||
const chargerID string = "TESTTEST"
|
||||
const ticks int64 = 638798974487432600
|
||||
|
||||
settingsUri := fmt.Sprintf("%s/chargers/%s/settings", easee.API, chargerID)
|
||||
commandUri := fmt.Sprintf("%s/chargers/%s/commands/resume_charging", easee.API, chargerID)
|
||||
|
||||
settingsReply := fmt.Sprintf("{\"device\":\"%s\",\"commandId\":48,\"ticks\":%d}", chargerID, ticks)
|
||||
|
||||
cmdResponse := easee.SignalRCommandResponse{
|
||||
WasAccepted: true,
|
||||
Ticks: ticks,
|
||||
}
|
||||
|
||||
testCases := []struct {
|
||||
uri string
|
||||
httpRc int
|
||||
respBody string
|
||||
cmdResp *easee.SignalRCommandResponse
|
||||
noop bool
|
||||
err error
|
||||
}{
|
||||
{settingsUri, 200, "", nil, false, nil}, //sync reply
|
||||
{settingsUri, 202, "[]", nil, true, nil}, //noop reply
|
||||
{settingsUri, 202, "[" + settingsReply + "]", nil, false, api.ErrTimeout}, //timeout
|
||||
{commandUri, 202, "{}", nil, true, nil}, //noop command reply
|
||||
{commandUri, 202, settingsReply, &cmdResponse, false, nil}, //command reply
|
||||
{commandUri, 400, "", nil, false, fmt.Errorf("invalid status: %d", 400)}, //unexpected result
|
||||
}
|
||||
|
||||
for _, tc := range testCases {
|
||||
t.Logf("%+v", tc)
|
||||
|
||||
e := newEasee()
|
||||
|
||||
httpmock.ActivateNonDefault(e.Client)
|
||||
httpmock.RegisterResponder(http.MethodPost, tc.uri,
|
||||
httpmock.NewStringResponder(tc.httpRc, tc.respBody))
|
||||
|
||||
if tc.cmdResp != nil {
|
||||
go func() {
|
||||
// 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
|
||||
}()
|
||||
}
|
||||
|
||||
noop, err := e.postJSONAndWait(tc.uri, nil)
|
||||
|
||||
assert.Equal(t, tc.noop, noop)
|
||||
assert.Equal(t, tc.err, err)
|
||||
|
||||
httpmock.Reset()
|
||||
}
|
||||
}
|
||||
|
||||
func TestEasee_waitForChargerEnabledState(t *testing.T) {
|
||||
testCases := []struct {
|
||||
expEnabled bool
|
||||
|
|
@ -374,11 +259,11 @@ func TestEasee_MaxCurrent(t *testing.T) {
|
|||
e := newEasee()
|
||||
e.charger = "CHARGERID"
|
||||
e.maxChargerCurrent = 16
|
||||
e.dynamicChargerCurrent = 6
|
||||
e.dynamicChargerCurrent = tc.expectCurrent // noop: already at target (capped) value
|
||||
|
||||
uriPattern := fmt.Sprintf("=~%s.*", easee.API)
|
||||
|
||||
//register mock NoOp reply, suffices for this test case
|
||||
// register mock NoOp reply (HTTP 202, empty array → Ticks==0)
|
||||
httpmock.ActivateNonDefault(e.Client)
|
||||
httpmock.RegisterResponder(http.MethodPost, uriPattern,
|
||||
httpmock.NewStringResponder(202, "[]"))
|
||||
|
|
@ -393,6 +278,10 @@ func TestEasee_MaxCurrent(t *testing.T) {
|
|||
}
|
||||
|
||||
func TestEasee_CommandResponse_rogue(t *testing.T) {
|
||||
// Intentionally triggers a WARN — suppress it to keep test output clean.
|
||||
util.LogLevel("error", nil)
|
||||
t.Cleanup(func() { util.LogLevel("info", nil) })
|
||||
|
||||
e := newEasee()
|
||||
|
||||
rogueResp := easee.SignalRCommandResponse{
|
||||
|
|
@ -410,45 +299,47 @@ func TestEasee_CommandResponse_rogue(t *testing.T) {
|
|||
e.CommandResponse(raw)
|
||||
})
|
||||
|
||||
// pendingTicks should still be empty
|
||||
e.cmdMu.Lock()
|
||||
assert.Empty(t, e.pendingTicks)
|
||||
e.cmdMu.Unlock()
|
||||
// The meaningful guarantee is no panic and no block.
|
||||
// Dispatcher's internal state is not directly inspectable from outside the package.
|
||||
}
|
||||
|
||||
func TestEasee_CommandResponse_legitimate(t *testing.T) {
|
||||
e := newEasee()
|
||||
httpmock.ActivateNonDefault(e.Client)
|
||||
|
||||
ticks := int64(638798974487432600)
|
||||
ch := make(chan easee.SignalRCommandResponse, 1)
|
||||
e.registerPendingTick(ticks, ch)
|
||||
const ticks int64 = 638798974487432600
|
||||
const uri = easee.API + "/chargers/EH123456/settings"
|
||||
|
||||
body := fmt.Sprintf(`[{"device":"EH123456","commandId":48,"ticks":%d}]`, ticks)
|
||||
httpmock.RegisterResponder(http.MethodPost, uri,
|
||||
httpmock.NewStringResponder(http.StatusAccepted, body))
|
||||
|
||||
resp := easee.SignalRCommandResponse{
|
||||
SerialNumber: "EH123456",
|
||||
Ticks: ticks,
|
||||
WasAccepted: true,
|
||||
}
|
||||
|
||||
raw, err := json.Marshal(resp)
|
||||
require.NoError(t, err)
|
||||
|
||||
e.CommandResponse(raw)
|
||||
errCh := make(chan error, 1)
|
||||
go func() {
|
||||
// Small delay to let Send register the pending tick before we dispatch
|
||||
time.Sleep(10 * time.Millisecond)
|
||||
e.CommandResponse(raw)
|
||||
errCh <- nil
|
||||
}()
|
||||
|
||||
// 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")
|
||||
}
|
||||
err = e.dispatcher.Send(uri, nil)
|
||||
assert.NoError(t, err)
|
||||
<-errCh
|
||||
}
|
||||
|
||||
func TestEasee_CommandResponse_expectedOrphan(t *testing.T) {
|
||||
e := newEasee()
|
||||
|
||||
// Pre-register the expected orphan
|
||||
e.registerExpectedOrphan(easee.CIRCUIT_MAX_CURRENT_P1)
|
||||
// Pre-register the expected orphan via the dispatcher's public API
|
||||
e.dispatcher.ExpectOrphan(easee.CIRCUIT_MAX_CURRENT_P1)
|
||||
|
||||
resp := easee.SignalRCommandResponse{
|
||||
SerialNumber: "EH123456",
|
||||
|
|
@ -467,14 +358,18 @@ func TestEasee_CommandResponse_expectedOrphan(t *testing.T) {
|
|||
})
|
||||
|
||||
// Counter should now be zero — a second response would be rogue
|
||||
assert.False(t, e.consumeExpectedOrphan(easee.CIRCUIT_MAX_CURRENT_P1))
|
||||
assert.False(t, e.dispatcher.CancelOrphan(easee.CIRCUIT_MAX_CURRENT_P1))
|
||||
}
|
||||
|
||||
func TestEasee_CommandResponse_rogueAfterOrphanConsumed(t *testing.T) {
|
||||
// Intentionally triggers a WARN — suppress it to keep test output clean.
|
||||
util.LogLevel("error", nil)
|
||||
t.Cleanup(func() { util.LogLevel("info", nil) })
|
||||
|
||||
e := newEasee()
|
||||
|
||||
// Register and immediately consume via CommandResponse
|
||||
e.registerExpectedOrphan(easee.CIRCUIT_MAX_CURRENT_P1)
|
||||
e.dispatcher.ExpectOrphan(easee.CIRCUIT_MAX_CURRENT_P1)
|
||||
|
||||
resp := easee.SignalRCommandResponse{
|
||||
SerialNumber: "EH123456",
|
||||
|
|
@ -491,95 +386,43 @@ func TestEasee_CommandResponse_rogueAfterOrphanConsumed(t *testing.T) {
|
|||
e.CommandResponse(raw)
|
||||
})
|
||||
|
||||
// pendingTicks untouched
|
||||
e.cmdMu.Lock()
|
||||
assert.Empty(t, e.pendingTicks)
|
||||
e.cmdMu.Unlock()
|
||||
// The meaningful guarantee is no panic and no block.
|
||||
// Dispatcher's internal state is not directly inspectable from outside the package.
|
||||
}
|
||||
|
||||
func TestEasee_CommandResponse_matchedByID(t *testing.T) {
|
||||
e := newEasee()
|
||||
httpmock.ActivateNonDefault(e.Client)
|
||||
|
||||
ch := make(chan easee.SignalRCommandResponse, 1)
|
||||
e.registerPendingByID(easee.LOCATION, ch)
|
||||
defer e.unregisterPendingByID(easee.LOCATION)
|
||||
const ticks int64 = 638798974487432601
|
||||
const obsID = easee.LOCATION // commandId in body
|
||||
const uri = easee.API + "/chargers/EH123456/settings"
|
||||
|
||||
// Ticks do NOT match any pendingTicks entry — only the ID matches
|
||||
body := fmt.Sprintf(`[{"device":"EH123456","commandId":%d,"ticks":%d}]`, int(obsID), ticks)
|
||||
httpmock.RegisterResponder(http.MethodPost, uri,
|
||||
httpmock.NewStringResponder(http.StatusAccepted, body))
|
||||
|
||||
// Ticks do NOT match — only the ID matches (wrong ticks value in SignalR response)
|
||||
resp := easee.SignalRCommandResponse{
|
||||
SerialNumber: "EH123456",
|
||||
ID: int(easee.LOCATION),
|
||||
Ticks: 999999999, // not in pendingTicks
|
||||
ID: int(obsID),
|
||||
Ticks: ticks + 1, // wrong ticks → forces ID fallback path
|
||||
WasAccepted: true,
|
||||
ResultCode: 0,
|
||||
}
|
||||
|
||||
raw, err := json.Marshal(resp)
|
||||
require.NoError(t, err)
|
||||
|
||||
assert.NotPanics(t, func() {
|
||||
errCh := make(chan error, 1)
|
||||
go func() {
|
||||
time.Sleep(10 * time.Millisecond)
|
||||
e.CommandResponse(raw)
|
||||
})
|
||||
errCh <- nil
|
||||
}()
|
||||
|
||||
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 TestProductUpdate_updatesLastObsReceived_freshTimestamp(t *testing.T) {
|
||||
e := newEasee()
|
||||
assert.True(t, e.lastObsReceived.IsZero())
|
||||
|
||||
// Observation with a fresh charger-side timestamp (seconds ago)
|
||||
now := time.Now().UTC().Truncate(0)
|
||||
e.ProductUpdate(createPayload(easee.TOTAL_POWER, now, easee.Double, "3.5"))
|
||||
|
||||
assert.False(t, e.lastObsReceived.IsZero())
|
||||
assert.WithinDuration(t, time.Now(), e.lastObsReceived, 5*time.Second)
|
||||
}
|
||||
|
||||
func TestProductUpdate_doesNotUpdateLastObsReceived_staleTimestamp(t *testing.T) {
|
||||
e := newEasee()
|
||||
|
||||
// Observation with a charger-side timestamp older than observationTimeout
|
||||
stale := time.Now().UTC().Add(-(observationTimeout + time.Minute))
|
||||
e.ProductUpdate(createPayload(easee.TOTAL_POWER, stale, easee.Double, "3.5"))
|
||||
|
||||
assert.True(t, e.lastObsReceived.IsZero(), "stale replay must not update lastObsReceived")
|
||||
err = e.dispatcher.Send(uri, nil)
|
||||
assert.NoError(t, err)
|
||||
<-errCh
|
||||
}
|
||||
|
||||
func TestEasee_Phases1p3p_registersExpectedOrphan(t *testing.T) {
|
||||
|
|
@ -617,10 +460,9 @@ func TestEasee_Phases1p3p_registersExpectedOrphan(t *testing.T) {
|
|||
|
||||
// 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")
|
||||
// CancelOrphan returns true iff a counter entry was consumed.
|
||||
assert.True(t, e.dispatcher.CancelOrphan(easee.CIRCUIT_MAX_CURRENT_P1),
|
||||
"expected orphan should be registered before the POST")
|
||||
}
|
||||
|
||||
func TestLivenessCheck_staleObservations(t *testing.T) {
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue