diff --git a/charger/warp-ws.go b/charger/warp-ws.go index dd79ab2cc..8ac51b530 100644 --- a/charger/warp-ws.go +++ b/charger/warp-ws.go @@ -10,6 +10,7 @@ import ( "net/url" "path" "slices" + "strings" "sync" "time" @@ -22,6 +23,13 @@ import ( "github.com/evcc-io/evcc/util/request" ) +type wsRole int + +const ( + wsRoleMain wsRole = iota + wsRolePM +) + type WarpWS struct { *warp.Connection implement.Caps @@ -48,8 +56,9 @@ type WarpWS struct { chargeTracker warp.ChargeTrackerCurrentCharge // power manager - pmState *warp.PmState - pmLowLevelState *warp.PmLowLevelState + pmState *warp.PmState + pmLowLevelState *warp.PmLowLevelState + lastPhasesWanted int // 0=never set; 1 or 3 } func init() { @@ -98,24 +107,8 @@ func NewWarpWSFromConfig(ctx context.Context, other map[string]any) (api.Charger // Feature: Phase Switching // only setup phase switching methods if power manager endpoint is set if (w.hasFeature(warp.FeaturePhaseSwitch) || cc.EnergyManagerURI != "") && w.pm != nil { - if res, err := w.ensurePmState(); err == nil && res.ExternalControl != warp.ExternalControlDeactivated { - w.pmState = &res - implement.Has(w, implement.PhaseSwitcher(w.phases1p3p)) - implement.Has(w, implement.PhaseGetter(w.getPhases)) - } - } - - // Phase Auto Switching needs to be disabled for WARP3 and WARP2 + EM - // Necessary if charging 1p only vehicles - typ, err := w.getWarpType() - if err != nil { - return nil, err - } - if typ == "warp3" || (typ == "warp2" && w.pm != nil && w.pm != w.Connection) { - if err := w.disablePhaseAutoSwitch(); err != nil { - return nil, err - } - w.log.TRACE.Println("disabled phase auto switching") + implement.Has(w, implement.PhaseSwitcher(w.phases1p3p)) + implement.Has(w, implement.PhaseGetter(w.getPhases)) } return w, nil @@ -142,17 +135,40 @@ func NewWarpWS(ctx context.Context, uri, user, pass, emURI, emUser, emPass strin w.pm = w.Connection } + // Phase Auto Switching needs to be disabled for WARP3 and WARP2 + EM + // Necessary if charging 1p only vehicles + typ, err := w.getWarpType() + if err != nil { + return nil, err + } + if typ == "warp3" || (typ == "warp2" && emURI != "") { + enabled, err := w.disablePhaseAutoSwitch() + if err != nil { + return nil, err + } + if enabled { + w.log.WARN.Println("disabled WARP phase auto switching") + } + } + wsURI, err := parseURI(w.URI) if err != nil { return nil, err } - go w.run(ctx, wsURI) + go w.run(ctx, wsRoleMain, w.Connection.Client, wsURI) + if emURI != "" { + pmWsURI, err := parseURI(w.pm.URI) + if err != nil { + return nil, err + } + go w.run(ctx, wsRolePM, w.pm.Client, pmWsURI) + } return w, nil } -func (w *WarpWS) run(ctx context.Context, wsURI string) { +func (w *WarpWS) run(ctx context.Context, role wsRole, client *http.Client, wsURI string) { bo := backoff.NewExponentialBackOff( backoff.WithMaxElapsedTime(0), backoff.WithMaxInterval(30*time.Second), @@ -161,7 +177,7 @@ func (w *WarpWS) run(ctx context.Context, wsURI string) { for ctx.Err() == nil { w.log.DEBUG.Println("websocket: connecting") - conn, _, err := websocket.Dial(ctx, wsURI, &websocket.DialOptions{HTTPClient: w.Client}) + conn, _, err := websocket.Dial(ctx, wsURI, &websocket.DialOptions{HTTPClient: client}) if err != nil { if !errors.Is(err, context.DeadlineExceeded) { w.log.ERROR.Printf("websocket: %v", err) @@ -178,12 +194,30 @@ func (w *WarpWS) run(ctx context.Context, wsURI string) { bo.Reset() - if err := w.handleConnection(ctx, conn); err != nil { + if role == wsRolePM { + if err := w.resendLastPhasesWantedIfAny(); err != nil { + w.log.WARN.Printf("resend phases_wanted on reconnect: %v", err) + } + } + + if err := w.handleConnection(ctx, role, conn); err != nil { w.log.ERROR.Println(err) } } } +func (w *WarpWS) resendLastPhasesWantedIfAny() error { + w.mu.RLock() + phases := w.lastPhasesWanted + w.mu.RUnlock() + + if phases == 0 { + return nil + } + + return w.postPhasesWanted(phases) +} + // Returns parsed URI and hostname func parseURI(uri string) (string, error) { u, err := url.Parse(uri) @@ -197,7 +231,11 @@ func parseURI(uri string) (string, error) { return u.String(), nil } -func (w *WarpWS) handleConnection(ctx context.Context, conn *websocket.Conn) error { +func isPmTopic(topic string) bool { + return strings.HasPrefix(topic, "power_manager/") +} + +func (w *WarpWS) handleConnection(ctx context.Context, role wsRole, conn *websocket.Conn) error { defer conn.Close(websocket.StatusInternalError, "reconnect") for { msgType, r, err := conn.Reader(ctx) @@ -221,6 +259,10 @@ func (w *WarpWS) handleConnection(ctx context.Context, conn *websocket.Conn) err return err } + if role == wsRoleMain && isPmTopic(event.Topic) { + continue + } + w.log.TRACE.Printf("websocket: event %s: %s", event.Topic, event.Payload) if err := w.handleEvent(event.Topic, event.Payload); err != nil { w.log.ERROR.Printf("bad payload for topic %s: %v", event.Topic, err) @@ -402,24 +444,50 @@ func (w *WarpWS) setCurrent(curr int64) error { return err } -func (w *WarpWS) disablePhaseAutoSwitch() error { +func (w *WarpWS) disablePhaseAutoSwitch() (bool, error) { uri := fmt.Sprintf("%s/evse/phase_auto_switch", w.URI) + var state struct { + Enabled bool `json:"enabled"` + } + if err := w.GetJSON(uri, &state); err != nil { + return false, err + } + if !state.Enabled { + return false, nil + } req, _ := request.New(http.MethodPost, uri, request.MarshalJSON(map[string]bool{"enabled": false}), request.JSONEncoding) _, err := w.Do(req) + return true, err +} + +func (w *WarpWS) postPhasesWanted(phases int) error { + uri := fmt.Sprintf("%s/power_manager/external_control", w.pm.URI) + req, _ := request.New(http.MethodPost, uri, request.MarshalJSON(map[string]int{"phases_wanted": phases}), request.JSONEncoding) + _, err := w.pm.Do(req) return err } // phases1p3p implements the api.PhaseSwitcher interface func (w *WarpWS) phases1p3p(phases int) error { - // ensure that phases can be switched - if ec, err := w.ensurePmState(); err != nil || ec.ExternalControl > warp.ExternalControlAvailable { - return fmt.Errorf("external control not available: %d", ec.ExternalControl) + // ExternalControlDeactivated is the WEM/WARP3 idle state before any + // phases_wanted has been sent — the POST below activates external control. + // Only block on states the POST cannot resolve. + ec, err := w.ensurePmState() + if err != nil { + return err + } + if ec.ExternalControl == warp.ExternalControlRuntimeConditionsNotMet || + ec.ExternalControl == warp.ExternalControlCurrentlySwitching { + return fmt.Errorf("external control %v: %w", ec.ExternalControl, api.ErrNotAvailable) } - uri := fmt.Sprintf("%s/power_manager/external_control", w.pm.URI) - req, _ := request.New(http.MethodPost, uri, request.MarshalJSON(map[string]int{"phases_wanted": phases}), request.JSONEncoding) - _, err := w.pm.Do(req) - return err + if err := w.postPhasesWanted(phases); err != nil { + return err + } + w.mu.Lock() + w.lastPhasesWanted = phases + w.mu.Unlock() + return nil } // getPhases implements the api.PhaseGetter interface @@ -428,7 +496,6 @@ func (w *WarpWS) getPhases() (int, error) { if err != nil { return 0, err } - if s.Is3phase { return 3, nil } @@ -448,6 +515,9 @@ func (w *WarpWS) ensurePmLowLevelState() (warp.PmLowLevelState, error) { return warp.PmLowLevelState{}, err } + w.mu.Lock() + w.pmLowLevelState = &ns + w.mu.Unlock() return ns, nil } @@ -464,6 +534,9 @@ func (w *WarpWS) ensurePmState() (warp.PmState, error) { return warp.PmState{}, err } + w.mu.Lock() + w.pmState = &res + w.mu.Unlock() return res, nil }