diff --git a/charger/pulsatrix.go b/charger/pulsatrix.go index 76512a5e4..665aea049 100644 --- a/charger/pulsatrix.go +++ b/charger/pulsatrix.go @@ -3,8 +3,8 @@ package charger /* MIT License -Copyright (c) 2023-2024 pulsatrix gmbh -Copyright (c) 2019-2024 andig +Copyright (c) 2023-2025 pulsatrix gmbh +Copyright (c) 2019-2025 andig Permission is hereby granted, free of charge, to any person obtaining a copy of this software and associated documentation files (the "Software"), to deal @@ -25,13 +25,27 @@ OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE. */ +/* +This module integrates the pulsatrix Supply Equipment Charge Controller (SECC) +with evcc.io enabling dynamic PV surplus charging. Communication is handled via +a bidirectional WebSocket connection, exchanging state data (e.g. vehicle +status, voltages, currents, energy counters) and sending control commands (e.g. +start/stop charging, current limits, phase switching) to the controller. + +Robust operation is ensured by automatic detection and handling of connection +losses, with reconnection based on exponential backoff strategies. In addition +to real-time data exchange, periodic heartbeats maintain connectivity. + +For further details, see: https://docs.pulsatrix.com or https://pulsatrix.de +*/ + import ( "bytes" "context" "encoding/json" "fmt" - "strconv" "sync" + "sync/atomic" "time" "github.com/cenkalti/backoff/v4" @@ -42,14 +56,30 @@ import ( "github.com/evcc-io/evcc/util/sponsor" ) +const ( + dataTimeout = 15 * time.Second + heartbeatInterval = 3 * time.Minute + maxRetries = 3 + syncRetries = 3 + backoffInitial = 2 * time.Second + backoffMax = 30 * time.Second + backoffMultiplier = 1.5 + errorLogInterval = 15 * time.Minute +) + // pulsatrix charger implementation type Pulsatrix struct { - log *util.Logger - mu sync.Mutex - conn *websocket.Conn - uri string - enabled bool - data *util.Monitor[pulsatrixData] + log *util.Logger + mu sync.RWMutex + conn *websocket.Conn + hostname string + uri string + enabled int32 // atomic for thread-safe access + data *util.Monitor[pulsatrixData] + cancel context.CancelFunc // for graceful shutdown + wg sync.WaitGroup // for goroutine synchronization + consecutiveReadErrors int32 // atomic counter + consecutiveHeartbeatErrors int32 // atomic counter } type pulsatrixData struct { @@ -65,7 +95,7 @@ func init() { registry.Add("pulsatrix", NewPulsatrixFromConfig) } -// NewPulsatrixtFromConfig creates a pulsatrix charger from generic config +// NewPulsatrixFromConfig creates a pulsatrix charger from generic config func NewPulsatrixFromConfig(other map[string]interface{}) (api.Charger, error) { var cc struct { Host string @@ -80,155 +110,301 @@ func NewPulsatrixFromConfig(other map[string]interface{}) (api.Charger, error) { // NewPulsatrix creates pulsatrix charger func NewPulsatrix(hostname string) (*Pulsatrix, error) { - wb := Pulsatrix{ - log: util.NewLogger("pulsatrix"), - uri: fmt.Sprintf("ws://%s/api/ws", hostname), - data: util.NewMonitor[pulsatrixData](15 * time.Second), - } - - if err := wb.connectWs(); err != nil { - return nil, err - } - + // check sponsor authorization early (fail fast) if !sponsor.IsAuthorized() { return nil, api.ErrSponsorRequired } - return &wb, nil + wb := &Pulsatrix{ + log: util.NewLogger("pulsatrix"), + hostname: hostname, + uri: fmt.Sprintf("ws://%s/api/ws", hostname), + data: util.NewMonitor[pulsatrixData](dataTimeout), + } + + if err := wb.connect(); err != nil { + return nil, fmt.Errorf("initial connection failed: %w", err) + } + + return wb, nil } -// ConnectWs connects to a pulsatrix SECC websocket -func (c *Pulsatrix) connectWs() error { +// connect connects to a pulsatrix SECC via websocket +func (c *Pulsatrix) connect() error { ctx, cancel := context.WithTimeout(context.Background(), request.Timeout) defer cancel() - c.log.TRACE.Printf("connecting to %s", c.uri) - conn, _, err := websocket.Dial(ctx, c.uri, nil) + conn, _, err := websocket.Dial(ctx, c.uri, &websocket.DialOptions{ + CompressionMode: websocket.CompressionDisabled, + }) if err != nil { + return fmt.Errorf("websocket dial to pulsatrix SECC at %s failed: %w", c.hostname, err) + } + + c.mu.Lock() + // close existing connection if present + if c.conn != nil { + c.conn.Close(websocket.StatusNormalClosure, "replacing connection") + } + c.conn = conn + + // create context for shutdown handling + ctx, c.cancel = context.WithCancel(context.Background()) + c.mu.Unlock() + + // sync with retry mechanism + if err := c.sync(); err != nil { + conn.Close(websocket.StatusInternalError, "sync failed") + return fmt.Errorf("sync failed: %w", err) + } + + // start background routines + c.wg.Add(2) + go c.reader(ctx) + go c.heartbeat(ctx) + + // reset error counters on successful connection + atomic.StoreInt32(&c.consecutiveReadErrors, 0) + atomic.StoreInt32(&c.consecutiveHeartbeatErrors, 0) + + c.log.INFO.Printf("connected to pulsatrix SECC at %s", c.hostname) + return nil +} + +// sync attempts synchronization with retry mechanism +func (c *Pulsatrix) sync() error { + for i := 0; i < syncRetries; i++ { + if err := c.Enable(false); err == nil { + return nil + } + if i < syncRetries-1 { + time.Sleep(time.Second) + } + } + return fmt.Errorf("sync with pulsatrix SECC at %s failed after %d attempts", c.hostname, syncRetries) +} + +// reconnect reconnects to a pulsatrix SECC websocket +func (c *Pulsatrix) reconnect() { + bo := backoff.NewExponentialBackOff( + backoff.WithInitialInterval(backoffInitial), + backoff.WithMaxInterval(backoffMax), + backoff.WithMultiplier(backoffMultiplier), + ) + bo.MaxElapsedTime = 0 // 0 means no time limit - retry indefinitely + + var lastErrorLog time.Time + operation := func() error { + err := c.connect() + if err != nil { + // log error every errorLogInterval + if time.Since(lastErrorLog) >= errorLogInterval { + c.log.ERROR.Printf("reconnect to pulsatrix SECC at %s still failing: %v", c.hostname, err) + lastErrorLog = time.Now() + } + } return err } - c.conn = conn - - // ensure evcc and SECC are in sync - if err := c.Enable(false); err != nil { - c.log.ERROR.Println(err) - } - - go c.wsReader() - go c.heartbeat(ctx) - - return nil -} - -// ReconnectWs reconnects to a pulsatrix SECC websocket -func (c *Pulsatrix) reconnectWs() { - bo := backoff.NewExponentialBackOff( - backoff.WithInitialInterval(time.Second), - backoff.WithMaxInterval(time.Minute), - backoff.WithMaxElapsedTime(0)) // retry forever; default is 15 min - if err := backoff.Retry(c.connectWs, bo); err != nil { - c.log.ERROR.Println(err) + if err := backoff.Retry(operation, bo); err != nil { + // should never be reached with MaxElapsedTime = 0 + c.log.ERROR.Printf("unexpected backoff failure: %v", err) } } -// WsReader runs a loop that reads messages from the websocket -func (c *Pulsatrix) wsReader() { - for { - ctx, cancel := context.WithTimeout(context.Background(), request.Timeout) - defer cancel() - - messageType, message, err := c.conn.Read(ctx) - if err != nil { - c.log.ERROR.Println("read message:", err) - break - } else { - c.parseWsMessage(messageType, message) +// reader runs a loop that reads messages from the websocket +func (c *Pulsatrix) reader(ctx context.Context) { + defer c.wg.Done() + defer func() { + c.mu.Lock() + if c.conn != nil { + c.conn.Close(websocket.StatusNormalClosure, "websocket reader shutting down") + c.conn = nil } - } + c.mu.Unlock() - c.mu.Lock() - c.conn.Close(websocket.StatusNormalClosure, "Reconnecting") - c.conn = nil - c.mu.Unlock() - - c.reconnectWs() -} - -// wsWrite writes a message to the websocket -func (c *Pulsatrix) write(message string) error { - c.mu.Lock() - defer c.mu.Unlock() - if c.conn != nil { - ctx, cancel := context.WithTimeout(context.Background(), request.Timeout) - defer cancel() - - if err := c.conn.Write(ctx, websocket.MessageText, []byte(message)); err != nil { - return err - } - } - return nil -} - -// ParseWsMessage parses a message from the websocket -func (c *Pulsatrix) parseWsMessage(messageType websocket.MessageType, message []byte) { - if messageType == websocket.MessageText { - b := bytes.ReplaceAll(message, []byte(":NaN"), []byte(":null")) - var parsedMessage struct { - Message json.RawMessage `json:"message"` - } - - if err := json.Unmarshal(b, &parsedMessage); err != nil { - c.log.ERROR.Println(err) - return - } - - val, _ := c.data.Get() - if err := json.Unmarshal(parsedMessage.Message, &val); err != nil { - c.log.ERROR.Println(err) - } else { - c.data.Set(val) - } - } -} - -// Heartbeat sends a heartbeat to the pulsatrix SECC -func (c *Pulsatrix) heartbeat(ctx context.Context) { - for tick := time.Tick(3 * time.Minute); ; { + // only reconnect if not explicitly stopped + select { + case <-ctx.Done(): + return // shutdown requested + default: + time.Sleep(time.Second) + go c.reconnect() + } + }() + + for { + // check for context cancellation select { - case <-tick: case <-ctx.Done(): return + default: } - if err := c.Enable(c.enabled); err != nil { - c.log.ERROR.Println(err) + readCtx, cancel := context.WithTimeout(ctx, request.Timeout) + messageType, message, err := c.getConn().Read(readCtx) + cancel() + + if err != nil { + // check if context was cancelled (graceful shutdown) + if ctx.Err() != nil { + return + } + + atomic.AddInt32(&c.consecutiveReadErrors, 1) + consecutiveErrors := atomic.LoadInt32(&c.consecutiveReadErrors) + + // warn only after consecutive errors + if consecutiveErrors >= maxRetries { + c.log.WARN.Printf("websocket read on pulsatrix SECC at %s failed %d times consecutively: %v", + c.hostname, consecutiveErrors, err) + } else { + c.log.TRACE.Printf("websocket read error on pulsatrix SECC at %s (attempt %d of %d): %v", + c.hostname, consecutiveErrors, maxRetries, err) + } + return // trigger defer reconnect + } + + // reset error counter after successful read + atomic.StoreInt32(&c.consecutiveReadErrors, 0) + + c.parseMessage(messageType, message) + } +} + +// getConn returns the current connection in a thread-safe manner +func (c *Pulsatrix) getConn() *websocket.Conn { + c.mu.RLock() + defer c.mu.RUnlock() + return c.conn +} + +// write writes a message to the websocket +func (c *Pulsatrix) write(message string) error { + conn := c.getConn() + if conn == nil { + return fmt.Errorf("websocket not connected to pulsatrix SECC at %s", c.hostname) + } + + ctx, cancel := context.WithTimeout(context.Background(), request.Timeout) + defer cancel() + + if err := conn.Write(ctx, websocket.MessageText, []byte(message)); err != nil { + c.log.WARN.Printf("write to pulsatrix SECC at %s failed: %v - trying reconnect", c.hostname, err) + go c.reconnect() // async reconnect + return err + } + return nil +} + +// parseMessage parses a message from the websocket +func (c *Pulsatrix) parseMessage(messageType websocket.MessageType, message []byte) { + if messageType != websocket.MessageText { + return + } + + if bytes.Contains(message, []byte(":NaN")) { + message = bytes.ReplaceAll(message, []byte(":NaN"), []byte(":null")) + } + + var parsedMessage struct { + Message json.RawMessage `json:"message"` + } + + if err := json.Unmarshal(message, &parsedMessage); err != nil { + c.log.DEBUG.Printf("failed to unmarshal websocket message: %v", err) + return + } + + val, _ := c.data.Get() + if err := json.Unmarshal(parsedMessage.Message, &val); err != nil { + c.log.DEBUG.Printf("failed to unmarshal message content: %v", err) + } else { + c.data.Set(val) + } +} + +// heartbeat sends a heartbeat to the pulsatrix SECC to keep remote control active +func (c *Pulsatrix) heartbeat(ctx context.Context) { + defer c.wg.Done() + + ticker := time.NewTicker(heartbeatInterval) + defer ticker.Stop() + + for { + select { + case <-ctx.Done(): + return + case <-ticker.C: + enabled := atomic.LoadInt32(&c.enabled) != 0 + if err := c.Enable(enabled); err != nil { + atomic.AddInt32(&c.consecutiveHeartbeatErrors, 1) + consecutiveErrors := atomic.LoadInt32(&c.consecutiveHeartbeatErrors) + + // warn only after consecutive failures + if consecutiveErrors >= maxRetries { + c.log.WARN.Printf("heartbeat with pulsatrix SECC at %s failed %d times consecutively: %v", + c.hostname, consecutiveErrors, err) + } + } else { + // reset error counter after successful heartbeat + atomic.StoreInt32(&c.consecutiveHeartbeatErrors, 0) + } } } } +// Shutdown gracefully closes the connection and stops all goroutines +func (c *Pulsatrix) Shutdown() error { + if c.cancel != nil { + c.cancel() + } + + // wait for all goroutines to finish + done := make(chan struct{}) + go func() { + c.wg.Wait() + close(done) + }() + + select { + case <-done: + return nil + case <-time.After(30 * time.Second): + return fmt.Errorf("shutdown timeout") + } +} + +// evcc.io API functions + // Status implements the api.Charger interface func (c *Pulsatrix) Status() (api.ChargeStatus, error) { res, err := c.data.Get() if err != nil { return api.StatusNone, err } - return api.ChargeStatusString(res.VehicleStatus) } // Enabled implements the api.Charger interface func (c *Pulsatrix) Enabled() (bool, error) { - return verifyEnabled(c, c.enabled) + enabled := atomic.LoadInt32(&c.enabled) != 0 + return verifyEnabled(c, enabled) } // Enable implements the api.Charger interface func (c *Pulsatrix) Enable(enable bool) error { - err := c.write("setEnabled\n" + strconv.FormatBool(enable)) - if err == nil { - c.enabled = enable + message := fmt.Sprintf("setEnabled\n%t", enable) + if err := c.write(message); err != nil { + return err } - return err + + var enabledVal int32 + if enable { + enabledVal = 1 + } + atomic.StoreInt32(&c.enabled, enabledVal) + return nil } // MaxCurrent implements the api.CurrentLimiter interface @@ -238,7 +414,8 @@ func (c *Pulsatrix) MaxCurrent(current int64) error { // MaxCurrentMillis implements the api.ChargerEx interface func (c *Pulsatrix) MaxCurrentMillis(current float64) error { - return c.write("setCurrentLimit\n" + strconv.FormatFloat(current, 'f', 10, 64)) + message := fmt.Sprintf("setCurrentLimit\n%g", current) + return c.write(message) } var _ api.CurrentGetter = (*Pulsatrix)(nil) @@ -265,7 +442,8 @@ func (c *Pulsatrix) TotalEnergy() (float64, error) { // Phases1p3p implements the api.PhaseSwitcher interface func (c *Pulsatrix) Phases1p3p(phases int) error { - return c.write("set1p3p\n" + strconv.FormatBool(phases == 1)) + message := fmt.Sprintf("set1p3p\n%t", phases == 1) + return c.write(message) } var _ api.PhaseCurrents = (*Pulsatrix)(nil) @@ -276,7 +454,6 @@ func (c *Pulsatrix) Currents() (float64, float64, float64, error) { if err != nil { return 0, 0, 0, err } - return res.PhaseAmperage[0], res.PhaseAmperage[1], res.PhaseAmperage[2], nil } @@ -288,6 +465,5 @@ func (c *Pulsatrix) Voltages() (float64, float64, float64, error) { if err != nil { return 0, 0, 0, err } - return res.PhaseVoltage[0], res.PhaseVoltage[1], res.PhaseVoltage[2], nil }