OCPP/EEBus: lazy-start servers so the CLI works against a running evcc (#30839)

Co-authored-by: Michael Geers <michael@geers.tv>
This commit is contained in:
andig 2026-06-25 17:26:08 +02:00 • committed by GitHub
parent 518556e6b8
commit fc84062905
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
23 changed files with 237 additions and 159 deletions

View file

@ -2,7 +2,6 @@ package charger
import ( import (
"context" "context"
"errors"
"fmt" "fmt"
"sync" "sync"
"time" "time"
@ -74,33 +73,34 @@ func NewEEBusFromConfig(ctx context.Context, other map[string]any) (api.Charger,
// newEEBus creates and initializes a raw *EEBus charger. // newEEBus creates and initializes a raw *EEBus charger.
// It registers the device with the EEBus instance and waits for the connection. // It registers the device with the EEBus instance and waits for the connection.
func newEEBus(ctx context.Context, ski, ip string) (*EEBus, error) { func newEEBus(ctx context.Context, ski, ip string) (*EEBus, error) {
if eebus.Instance == nil { inst, err := eebus.Instance()
return nil, errors.New("eebus not configured") if err != nil {
return nil, err
} }
c := &EEBus{ c := &EEBus{
Caps: implement.New(), Caps: implement.New(),
log: util.NewLogger("eebus"), log: util.NewLogger("eebus"),
current: 6, current: 6,
cem: eebus.Instance.CustomerEnergyManagement(), cem: inst.CustomerEnergyManagement(),
} }
c.connector = eebus.NewConnector() c.connector = eebus.NewConnector()
c.minMaxG = util.Cached(c.minMax, time.Second) c.minMaxG = util.Cached(c.minMax, time.Second)
if err := eebus.Instance.RegisterDevice(ski, ip, c); err != nil { if err := inst.RegisterDevice(ski, ip, c); err != nil {
return nil, err return nil, err
} }
if err := c.connector.Wait(ctx); err != nil { if err := c.connector.Wait(ctx); err != nil {
eebus.Instance.UnregisterDevice(ski, c) inst.UnregisterDevice(ski, c)
return nil, err return nil, err
} }
// unregister device when context is cancelled (e.g. UI config validation) // unregister device when context is cancelled (e.g. UI config validation)
go func() { go func() {
<-ctx.Done() <-ctx.Done()
eebus.Instance.UnregisterDevice(ski, c) inst.UnregisterDevice(ski, c)
}() }()
return c, nil return c, nil

View file

@ -139,9 +139,14 @@ func NewOCPP(ctx context.Context,
) (*OCPP, error) { ) (*OCPP, error) {
log := util.NewLogger(fmt.Sprintf("%s-%d", cmp.Or(id, "ocpp"), connector)) log := util.NewLogger(fmt.Sprintf("%s-%d", cmp.Or(id, "ocpp"), connector))
cp, err := ocpp.Instance().RegisterChargepoint(id, cs, err := ocpp.Instance()
if err != nil {
return nil, err
}
cp, err := cs.RegisterChargepoint(id,
func() *ocpp.CP { func() *ocpp.CP {
return ocpp.NewChargePoint(log, id) return ocpp.NewChargePoint(log, cs, id)
}, },
func(cp *ocpp.CP) error { func(cp *ocpp.CP) error {
log.DEBUG.Printf("waiting for chargepoint: %v", connectTimeout) log.DEBUG.Printf("waiting for chargepoint: %v", connectTimeout)

View file

@ -63,7 +63,7 @@ func NewConnector(ctx context.Context, log *util.Logger, id int, cp *CP, idTag s
var ok bool var ok bool
// apply cached status if available // apply cached status if available
instance.WithConnectorStatus(cp.ID(), id, func(status *core.StatusNotificationRequest) { cp.cs.WithConnectorStatus(cp.ID(), id, func(status *core.StatusNotificationRequest) {
if _, err := cp.OnStatusNotification(status); err == nil { if _, err := cp.OnStatusNotification(status); err == nil {
ok = true ok = true
} }

View file

@ -111,7 +111,7 @@ func (conn *Connector) OnStartTransaction(request *core.StartTransactionRequest)
conn.mu.Lock() conn.mu.Lock()
defer conn.mu.Unlock() defer conn.mu.Unlock()
conn.txnId = int(instance.txnId.Add(1)) conn.txnId = int(conn.cp.cs.txnId.Add(1))
conn.idTag = request.IdTag conn.idTag = request.IdTag
res := &core.StartTransactionConfirmation{ res := &core.StartTransactionConfirmation{

View file

@ -26,7 +26,7 @@ type connTestSuite struct {
func (suite *connTestSuite) SetupTest() { func (suite *connTestSuite) SetupTest() {
// setup instance // setup instance
Instance() Instance()
suite.cp = NewChargePoint(util.NewLogger("foo"), "abc") suite.cp = NewChargePoint(util.NewLogger("foo"), instance, "abc")
suite.conn, _ = NewConnector(suite.T().Context(), util.NewLogger("foo"), 1, suite.cp, "", Timeout) suite.conn, _ = NewConnector(suite.T().Context(), util.NewLogger("foo"), 1, suite.cp, "", Timeout)
suite.clock = clock.NewMock() suite.clock = clock.NewMock()

View file

@ -16,6 +16,7 @@ import (
type CP struct { type CP struct {
mu sync.RWMutex mu sync.RWMutex
cs *CS // central system this charge point is registered with
log *util.Logger log *util.Logger
onceConnect sync.Once onceConnect sync.Once
onceMonitor sync.Once onceMonitor sync.Once
@ -43,8 +44,9 @@ type CP struct {
connectors map[int]*Connector connectors map[int]*Connector
} }
func NewChargePoint(log *util.Logger, id string) *CP { func NewChargePoint(log *util.Logger, cs *CS, id string) *CP {
return &CP{ return &CP{
cs: cs,
log: log, log: log,
id: id, id: id,
@ -171,7 +173,7 @@ func (cp *CP) onTransportConnect() {
cp.log.DEBUG.Printf("proactively triggering BootNotification") cp.log.DEBUG.Printf("proactively triggering BootNotification")
if err := Instance().TriggerMessage( if err := cp.cs.TriggerMessage(
cp.id, cp.id,
func(conf *remotetrigger.TriggerMessageConfirmation, err error) { func(conf *remotetrigger.TriggerMessageConfirmation, err error) {
if err != nil { if err != nil {

View file

@ -14,7 +14,7 @@ import (
func TestBootNotificationStoresResultAndConnects(t *testing.T) { func TestBootNotificationStoresResultAndConnects(t *testing.T) {
log := util.NewLogger("test") log := util.NewLogger("test")
cp := NewChargePoint(log, "test-cp") cp := NewChargePoint(log, instance, "test-cp")
assert.False(t, cp.Connected(), "should not be connected initially") assert.False(t, cp.Connected(), "should not be connected initially")
assert.Nil(t, cp.BootNotificationResult, "should have no boot result initially") assert.Nil(t, cp.BootNotificationResult, "should have no boot result initially")
@ -43,7 +43,7 @@ func TestBootNotificationStoresResultAndConnects(t *testing.T) {
func TestBootNotificationStopsTimer(t *testing.T) { func TestBootNotificationStopsTimer(t *testing.T) {
log := util.NewLogger("test") log := util.NewLogger("test")
cp := NewChargePoint(log, "test-cp") cp := NewChargePoint(log, instance, "test-cp")
// simulate WebSocket connect (starts timer) // simulate WebSocket connect (starts timer)
cp.onTransportConnect() cp.onTransportConnect()
@ -73,7 +73,7 @@ func TestBootNotificationStopsTimer(t *testing.T) {
func TestTransportConnectTimeoutFallback(t *testing.T) { func TestTransportConnectTimeoutFallback(t *testing.T) {
log := util.NewLogger("test") log := util.NewLogger("test")
cp := NewChargePoint(log, "test-cp") cp := NewChargePoint(log, instance, "test-cp")
// use a short timeout for testing // use a short timeout for testing
origTimeout := Timeout origTimeout := Timeout
@ -99,7 +99,7 @@ func TestTransportConnectTimeoutFallback(t *testing.T) {
func TestDisconnectCancelsTimer(t *testing.T) { func TestDisconnectCancelsTimer(t *testing.T) {
log := util.NewLogger("test") log := util.NewLogger("test")
cp := NewChargePoint(log, "test-cp") cp := NewChargePoint(log, instance, "test-cp")
// use a short timeout for testing // use a short timeout for testing
origTimeout := Timeout origTimeout := Timeout
@ -126,7 +126,7 @@ func TestDisconnectCancelsTimer(t *testing.T) {
func TestBootNotificationChannelCoalesces(t *testing.T) { func TestBootNotificationChannelCoalesces(t *testing.T) {
log := util.NewLogger("test") log := util.NewLogger("test")
cp := NewChargePoint(log, "test-cp") cp := NewChargePoint(log, instance, "test-cp")
// pre-fill the channel (buffer size 1) with a stale notification // pre-fill the channel (buffer size 1) with a stale notification
cp.bootNotificationRequestC <- &core.BootNotificationRequest{ cp.bootNotificationRequestC <- &core.BootNotificationRequest{
@ -165,7 +165,7 @@ func TestBootNotificationChannelCoalesces(t *testing.T) {
// notification so a later Setup re-runs against the charge point's real state. // notification so a later Setup re-runs against the charge point's real state.
func TestBootNotificationRebootLoopKeepsLatest(t *testing.T) { func TestBootNotificationRebootLoopKeepsLatest(t *testing.T) {
log := util.NewLogger("test") log := util.NewLogger("test")
cp := NewChargePoint(log, "test-cp") cp := NewChargePoint(log, instance, "test-cp")
for i := range 5 { for i := range 5 {
_, err := cp.OnBootNotification(&core.BootNotificationRequest{ _, err := cp.OnBootNotification(&core.BootNotificationRequest{
@ -188,7 +188,7 @@ func TestBootNotificationRebootLoopKeepsLatest(t *testing.T) {
func TestReconnectAfterReboot(t *testing.T) { func TestReconnectAfterReboot(t *testing.T) {
log := util.NewLogger("test") log := util.NewLogger("test")
cp := NewChargePoint(log, "test-cp") cp := NewChargePoint(log, instance, "test-cp")
// simulate initial connection // simulate initial connection
cp.connect(true) cp.connect(true)
@ -220,7 +220,7 @@ func TestReconnectAfterReboot(t *testing.T) {
func TestMonitorRebootOnlyOnce(t *testing.T) { func TestMonitorRebootOnlyOnce(t *testing.T) {
log := util.NewLogger("test") log := util.NewLogger("test")
cp := NewChargePoint(log, "test-cp") cp := NewChargePoint(log, instance, "test-cp")
ctx := t.Context() ctx := t.Context()
var callCount atomic.Int32 var callCount atomic.Int32

View file

@ -13,7 +13,7 @@ import (
func (cp *CP) ChangeAvailabilityRequest(connectorId int, availabilityType core.AvailabilityType) error { func (cp *CP) ChangeAvailabilityRequest(connectorId int, availabilityType core.AvailabilityType) error {
rc := make(chan error, 1) rc := make(chan error, 1)
err := Instance().ChangeAvailability(cp.id, func(request *core.ChangeAvailabilityConfirmation, err error) { err := cp.cs.ChangeAvailability(cp.id, func(request *core.ChangeAvailabilityConfirmation, err error) {
if err == nil && request != nil && request.Status != core.AvailabilityStatusAccepted && request.Status != core.AvailabilityStatusScheduled { if err == nil && request != nil && request.Status != core.AvailabilityStatusAccepted && request.Status != core.AvailabilityStatusScheduled {
err = errors.New(string(request.Status)) err = errors.New(string(request.Status))
} }
@ -28,7 +28,7 @@ func (cp *CP) GetCompositeScheduleRequest(connectorId int, duration int) (*smart
var res *smartcharging.GetCompositeScheduleConfirmation var res *smartcharging.GetCompositeScheduleConfirmation
rc := make(chan error, 1) rc := make(chan error, 1)
err := Instance().GetCompositeSchedule(cp.id, func(request *smartcharging.GetCompositeScheduleConfirmation, err error) { err := cp.cs.GetCompositeSchedule(cp.id, func(request *smartcharging.GetCompositeScheduleConfirmation, err error) {
if err == nil && request != nil && request.Status != smartcharging.GetCompositeScheduleStatusAccepted { if err == nil && request != nil && request.Status != smartcharging.GetCompositeScheduleStatusAccepted {
err = errors.New(string(request.Status)) err = errors.New(string(request.Status))
} }
@ -43,7 +43,7 @@ func (cp *CP) GetCompositeScheduleRequest(connectorId int, duration int) (*smart
func (cp *CP) RemoteStartTransactionRequest(connectorId int, idTag string) error { func (cp *CP) RemoteStartTransactionRequest(connectorId int, idTag string) error {
rc := make(chan error, 1) rc := make(chan error, 1)
err := Instance().RemoteStartTransaction(cp.id, func(request *core.RemoteStartTransactionConfirmation, err error) { err := cp.cs.RemoteStartTransaction(cp.id, func(request *core.RemoteStartTransactionConfirmation, err error) {
if err == nil && request != nil && request.Status != types.RemoteStartStopStatusAccepted { if err == nil && request != nil && request.Status != types.RemoteStartStopStatusAccepted {
err = errors.New(string(request.Status)) err = errors.New(string(request.Status))
} }
@ -61,7 +61,7 @@ func (cp *CP) RemoteStartTransactionRequest(connectorId int, idTag string) error
func (cp *CP) SetChargingProfileRequest(connectorId int, profile *types.ChargingProfile) error { func (cp *CP) SetChargingProfileRequest(connectorId int, profile *types.ChargingProfile) error {
rc := make(chan error, 1) rc := make(chan error, 1)
err := Instance().SetChargingProfile(cp.id, func(request *smartcharging.SetChargingProfileConfirmation, err error) { err := cp.cs.SetChargingProfile(cp.id, func(request *smartcharging.SetChargingProfileConfirmation, err error) {
if err == nil && request != nil && request.Status != smartcharging.ChargingProfileStatusAccepted { if err == nil && request != nil && request.Status != smartcharging.ChargingProfileStatusAccepted {
err = errors.New(string(request.Status)) err = errors.New(string(request.Status))
} }
@ -79,7 +79,7 @@ func (cp *CP) TriggerMessageRequest(connectorId int, requestedMessage remotetrig
rc := make(chan error, 1) rc := make(chan error, 1)
err := Instance().TriggerMessage(cp.id, func(request *remotetrigger.TriggerMessageConfirmation, err error) { err := cp.cs.TriggerMessage(cp.id, func(request *remotetrigger.TriggerMessageConfirmation, err error) {
if err == nil && request != nil && request.Status != remotetrigger.TriggerMessageStatusAccepted { if err == nil && request != nil && request.Status != remotetrigger.TriggerMessageStatusAccepted {
err = errors.New(string(request.Status)) err = errors.New(string(request.Status))
} }
@ -97,7 +97,7 @@ func (cp *CP) TriggerMessageRequest(connectorId int, requestedMessage remotetrig
func (cp *CP) ChangeConfigurationRequest(key, value string) error { func (cp *CP) ChangeConfigurationRequest(key, value string) error {
rc := make(chan error, 1) rc := make(chan error, 1)
err := Instance().ChangeConfiguration(cp.id, func(request *core.ChangeConfigurationConfirmation, err error) { err := cp.cs.ChangeConfiguration(cp.id, func(request *core.ChangeConfigurationConfirmation, err error) {
if err == nil && request != nil && request.Status != core.ConfigurationStatusAccepted { if err == nil && request != nil && request.Status != core.ConfigurationStatusAccepted {
err = errors.New(string(request.Status)) err = errors.New(string(request.Status))
} }
@ -112,7 +112,7 @@ func (cp *CP) GetConfigurationRequest() (*core.GetConfigurationConfirmation, err
rc := make(chan error, 1) rc := make(chan error, 1)
var res *core.GetConfigurationConfirmation var res *core.GetConfigurationConfirmation
err := Instance().GetConfiguration(cp.id, func(request *core.GetConfigurationConfirmation, err error) { err := cp.cs.GetConfiguration(cp.id, func(request *core.GetConfigurationConfirmation, err error) {
res = request res = request
rc <- err rc <- err

View file

@ -9,6 +9,7 @@ import (
"github.com/evcc-io/evcc/util" "github.com/evcc-io/evcc/util"
ocpp16 "github.com/lorenzodonini/ocpp-go/ocpp1.6" ocpp16 "github.com/lorenzodonini/ocpp-go/ocpp1.6"
"github.com/lorenzodonini/ocpp-go/ocpp1.6/core" "github.com/lorenzodonini/ocpp-go/ocpp1.6/core"
"github.com/lorenzodonini/ocpp-go/ocppj"
"github.com/lorenzodonini/ocpp-go/ws" "github.com/lorenzodonini/ocpp-go/ws"
) )
@ -30,7 +31,8 @@ type CS struct {
regs map[string]*registration // guarded by mu mutex regs map[string]*registration // guarded by mu mutex
txnId atomic.Int64 txnId atomic.Int64
publishFunc func() publishFunc func()
server ws.Server // raw server, used by the forwarder to write frames server ws.Server // raw server, used by the forwarder to write frames
dispatcher ocppj.ServerDispatcher // request dispatcher, timeout set at start
} }
// Write sends a raw OCPP frame to the charger with the given station ID. // Write sends a raw OCPP frame to the charger with the given station ID.

View file

@ -420,7 +420,10 @@ func drainPendingWithErrors(id string, sc *sidecar) {
delete(pendingMsgs, id) delete(pendingMsgs, id)
pendingMu.Unlock() pendingMu.Unlock()
cs := Instance() cs, err := Instance()
if err != nil {
return
}
for _, frame := range buffered { for _, frame := range buffered {
msgType, msgID, action, err := parseOCPPFrame(frame) msgType, msgID, action, err := parseOCPPFrame(frame)
@ -587,6 +590,12 @@ func (sc *sidecar) readFromUpstream() {
notifyUpdated() notifyUpdated()
}() }()
cs, err := Instance()
if err != nil {
forwarderLog.ERROR.Printf("forwarder: central system unavailable for %s: %v", sc.chargerID, err)
return
}
for { for {
_, msg, err := sc.conn.Read(context.Background()) _, msg, err := sc.conn.Read(context.Background())
if err != nil { if err != nil {
@ -645,7 +654,7 @@ func (sc *sidecar) readFromUpstream() {
sc.pendingUpstreamCalls[msgID] = struct{}{} sc.pendingUpstreamCalls[msgID] = struct{}{}
sc.pendingUpstreamCallsMu.Unlock() sc.pendingUpstreamCallsMu.Unlock()
if err := Instance().Write(sc.chargerID, msg); err != nil { if err := cs.Write(sc.chargerID, msg); err != nil {
forwarderLog.ERROR.Printf("forwarder: inject upstream call into charger %s: %v", sc.chargerID, err) forwarderLog.ERROR.Printf("forwarder: inject upstream call into charger %s: %v", sc.chargerID, err)
} }
@ -660,7 +669,7 @@ func (sc *sidecar) readFromUpstream() {
if isChargerCall { if isChargerCall {
// relay to charger; its handler was bypassed and it awaits this reply // relay to charger; its handler was bypassed and it awaits this reply
if err := Instance().Write(sc.chargerID, msg); err != nil { if err := cs.Write(sc.chargerID, msg); err != nil {
forwarderLog.ERROR.Printf("forwarder: relay upstream response to charger %s: %v", sc.chargerID, err) forwarderLog.ERROR.Printf("forwarder: relay upstream response to charger %s: %v", sc.chargerID, err)
} }
continue continue

View file

@ -1,6 +1,7 @@
package ocpp package ocpp
import ( import (
"errors"
"fmt" "fmt"
"net/http" "net/http"
"net/url" "net/url"
@ -43,8 +44,8 @@ func (r ForwarderRule) Redacted() ForwarderRule {
} }
var ( var (
once sync.Once
instance *CS instance *CS
started func() error // memoized listen; set in NewServer, runs once
port = 8887 port = 8887
boundPort int boundPort int
externalUrl string externalUrl string
@ -129,64 +130,77 @@ func CurrentConfig() Config {
return Config{Port: port} return Config{Port: port}
} }
// Init initializes the OCPP server // NewServer builds the OCPP central system without starting it.
func Init(cfg Config, networkExternalUrl string) { func NewServer(cfg Config, networkExternalUrl string) {
port = cfg.Port port = cfg.Port
externalUrl = networkExternalUrl externalUrl = networkExternalUrl
}
func Instance() *CS { log := util.NewLogger("ocpp")
once.Do(func() {
log := util.NewLogger("ocpp")
server := &interceptingServer{Server: ws.NewServer()} server := &interceptingServer{Server: ws.NewServer()}
server.SetCheckOriginHandler(func(r *http.Request) bool { return true }) server.SetCheckOriginHandler(func(r *http.Request) bool { return true })
dispatcher := ocppj.NewDefaultServerDispatcher(ocppj.NewFIFOQueueMap(0)) dispatcher := ocppj.NewDefaultServerDispatcher(ocppj.NewFIFOQueueMap(0))
dispatcher.SetTimeout(Timeout)
endpoint := ocppj.NewServer(server, dispatcher, nil, core.Profile, remotetrigger.Profile, smartcharging.Profile, security.Profile, firmware.Profile) endpoint := ocppj.NewServer(server, dispatcher, nil, core.Profile, remotetrigger.Profile, smartcharging.Profile, security.Profile, firmware.Profile)
endpoint.SetInvalidMessageHook(func(client ws.Channel, err *ocpp.Error, rawMessage string, parsedFields []any) *ocpp.Error { endpoint.SetInvalidMessageHook(func(client ws.Channel, err *ocpp.Error, rawMessage string, parsedFields []any) *ocpp.Error {
log.ERROR.Printf("%v (%s)", err, rawMessage) log.ERROR.Printf("%v (%s)", err, rawMessage)
return nil return nil
})
cs := ocpp16.NewCentralSystem(endpoint, server)
instance = &CS{
log: log,
regs: make(map[string]*registration),
CentralSystem: cs,
server: server,
}
instance.txnId.Store(time.Now().UTC().Unix())
ocppj.SetLogger(instance)
cs.SetCoreHandler(instance)
cs.SetSecurityHandler(instance)
cs.SetFirmwareManagementHandler(instance)
cs.SetNewChargePointHandler(instance.NewChargePoint)
cs.SetChargePointDisconnectedHandler(instance.ChargePointDisconnected)
go instance.errorHandler(cs.Errors())
go cs.Start(port, "/{ws}")
// wait for server to start
tick := time.Tick(10 * time.Millisecond)
timeout := time.After(10 * time.Second)
for server.Addr() == nil {
select {
case <-tick:
case <-timeout:
log.ERROR.Println("timeout waiting for server to bind")
return
}
}
boundPort = server.Addr().Port
}) })
return instance cs := ocpp16.NewCentralSystem(endpoint, server)
inst := &CS{
log: log,
regs: make(map[string]*registration),
CentralSystem: cs,
server: server,
dispatcher: dispatcher,
}
inst.txnId.Store(time.Now().UTC().Unix())
ocppj.SetLogger(inst)
cs.SetCoreHandler(inst)
cs.SetSecurityHandler(inst)
cs.SetFirmwareManagementHandler(inst)
cs.SetNewChargePointHandler(inst.NewChargePoint)
cs.SetChargePointDisconnectedHandler(inst.ChargePointDisconnected)
// wire the start memo before publishing instance, so Instance() never sees
// a non-nil instance with a nil started
started = sync.OnceValue(inst.listen)
instance = inst
}
// listen starts the central system and blocks until it has bound its port.
func (cs *CS) listen() error {
cs.dispatcher.SetTimeout(Timeout)
go cs.errorHandler(cs.Errors())
go cs.CentralSystem.Start(port, "/{ws}")
// wait for server to bind
tick := time.Tick(10 * time.Millisecond)
timeout := time.After(10 * time.Second)
for cs.server.Addr() == nil {
select {
case <-tick:
case <-timeout:
return errors.New("timeout waiting for server to bind")
}
}
boundPort = cs.server.Addr().Port
return nil
}
// Instance returns the central system, starting it once on first call. It
// returns an error if the server fails to bind its port.
func Instance() (*CS, error) {
if instance == nil {
return nil, errors.New("ocpp not configured")
}
return instance, started()
} }

View file

@ -9,7 +9,7 @@ func TestMain(m *testing.M) {
// bind the central system to an ephemeral port so this test binary does not // bind the central system to an ephemeral port so this test binary does not
// contend with the charger package test binary for the fixed default port // contend with the charger package test binary for the fixed default port
// when both run in parallel under `go test ./...` // when both run in parallel under `go test ./...`
Init(Config{Port: 0}, "") NewServer(Config{Port: 0}, "")
os.Exit(m.Run()) os.Exit(m.Run())
} }

View file

@ -35,7 +35,7 @@ func TestMain(m *testing.M) {
// bind the OCPP central system to an ephemeral port so this test binary // bind the OCPP central system to an ephemeral port so this test binary
// does not contend with the charger/ocpp package test binary for the fixed // does not contend with the charger/ocpp package test binary for the fixed
// default port when both run in parallel under `go test ./...` // default port when both run in parallel under `go test ./...`
ocpp.Init(ocpp.Config{Port: 0}, "") ocpp.NewServer(ocpp.Config{Port: 0}, "")
os.Exit(m.Run()) os.Exit(m.Run())
} }
@ -57,7 +57,10 @@ func (suite *ocppTestSuite) SetupSuite() {
ocpp.TriggerBootDelay = 100 * time.Millisecond ocpp.TriggerBootDelay = 100 * time.Millisecond
// setup cs so we can overwrite logger afterwards // setup cs so we can overwrite logger afterwards
_ = ocpp.Instance() cs, err := ocpp.Instance()
suite.Require().NoError(err, "instance")
suite.NotNil(cs)
suite.logger = &ocppLogger{t: suite.T()} suite.logger = &ocppLogger{t: suite.T()}
ocppj.SetLogger(suite.logger) ocppj.SetLogger(suite.logger)
@ -65,7 +68,6 @@ func (suite *ocppTestSuite) SetupSuite() {
ocppTestUrl = fmt.Sprintf("ws://localhost:%d", ocpp.Port()) ocppTestUrl = fmt.Sprintf("ws://localhost:%d", ocpp.Port())
suite.clock = clock.NewMock() suite.clock = clock.NewMock()
suite.NotNil(ocpp.Instance())
} }
func (suite *ocppTestSuite) TearDownSuite() { func (suite *ocppTestSuite) TearDownSuite() {

View file

@ -200,28 +200,39 @@ func runRoot(cmd *cobra.Command, args []string) {
valueChan := make(chan util.Param, 64) valueChan := make(chan util.Param, 64)
go tee.Run(valueChan) go tee.Run(valueChan)
// start OCPP server // start OCPP and EEBus servers (skipped in degraded mode where setup failed,
ocppCS := ocpp.Instance() // so a misconfigured instance serves only the offline UI)
ocppCS.SetUpdated(func() { if err == nil {
// republish when OCPP state updates cs, ocppErr := ocpp.Instance()
valueChan <- util.Param{Key: keys.Ocpp, Val: globalconfig.ConfigStatus{ if ocppErr != nil {
Config: ocpp.CurrentConfig(), log.ERROR.Println("ocpp:", ocppErr)
Status: ocpp.GetStatus(), } else {
}} cs.SetUpdated(func() {
}) // republish when OCPP state updates
log.INFO.Printf("OCPP local url: ws://127.0.0.1:%d/<stationId>", conf.Ocpp.Port) valueChan <- util.Param{Key: keys.Ocpp, Val: globalconfig.ConfigStatus{
if ocpp.ExternalUrl() != "" { Config: ocpp.CurrentConfig(),
log.INFO.Printf("OCPP external url: %s/<stationId>", ocpp.ExternalUrl()) Status: ocpp.GetStatus(),
} }}
// register the callback even with no rules so runtime additions are pushed })
ocpp.SetForwarderUpdated(func() { log.INFO.Printf("OCPP local url: ws://127.0.0.1:%d/<stationId>", conf.Ocpp.Port)
valueChan <- util.Param{Key: keys.OcppForwarder, Val: globalconfig.ConfigStatus{ if ocpp.ExternalUrl() != "" {
Config: lo.Map(ocpp.ForwarderRules(), func(r ocpp.ForwarderRule, _ int) ocpp.ForwarderRule { return r.Redacted() }), log.INFO.Printf("OCPP external url: %s/<stationId>", ocpp.ExternalUrl())
Status: ocpp.GetForwarderStatus(), }
}} }
}) // register the callback even with no rules so runtime additions are pushed
if ocpp.ForwarderEnabled() { ocpp.SetForwarderUpdated(func() {
log.INFO.Printf("OCPP forwarder: %d rule(s) active", len(ocpp.ForwarderRules())) valueChan <- util.Param{Key: keys.OcppForwarder, Val: globalconfig.ConfigStatus{
Config: lo.Map(ocpp.ForwarderRules(), func(r ocpp.ForwarderRule, _ int) ocpp.ForwarderRule { return r.Redacted() }),
Status: ocpp.GetForwarderStatus(),
}}
})
if ocpp.ForwarderEnabled() {
log.INFO.Printf("OCPP forwarder: %d rule(s) active", len(ocpp.ForwarderRules()))
}
if _, eebusErr := eebus.Instance(); eebusErr != nil {
log.ERROR.Println("eebus:", eebusErr)
}
} }
// value cache // value cache

View file

@ -853,11 +853,16 @@ func configureOCPP(cfg *ocpp.Config, externalUrl string) {
log.WARN.Printf("ocpp: failed to load settings: %v", err) log.WARN.Printf("ocpp: failed to load settings: %v", err)
} }
} }
ocpp.Init(*cfg, externalUrl) ocpp.NewServer(*cfg, externalUrl)
// Load proxy forwarding rules from DB if present. // Load proxy forwarding rules from DB if present.
var rules []ocpp.ForwarderRule var rules []ocpp.ForwarderRule
if err := settings.Json(keys.OcppForwarder, &rules); err == nil { if err := settings.Json(keys.OcppForwarder, &rules); err == nil && len(rules) > 0 {
// the forwarder relays through the central system; abort if it failed to start
if _, err := ocpp.Instance(); err != nil {
log.ERROR.Printf("ocpp: forwarder disabled: %v", err)
return
}
ocpp.ApplyForwarderRules(rules) ocpp.ApplyForwarderRules(rules)
} }
} }
@ -885,13 +890,12 @@ func configureEEBus(conf *eebus.Config) error {
} }
} }
var err error srv, err := eebus.NewServer(*conf)
if eebus.Instance, err = eebus.NewServer(*conf); err != nil { if err != nil {
return fmt.Errorf("failed configuring eebus: %w", err) return fmt.Errorf("failed configuring eebus: %w", err)
} }
eebus.Instance.Run() shutdown.Register(srv.Shutdown)
shutdown.Register(eebus.Instance.Shutdown)
return nil return nil
} }

View file

@ -2,7 +2,6 @@ package eebus
import ( import (
"context" "context"
"errors"
"sync" "sync"
"time" "time"
@ -89,15 +88,16 @@ func NewFromConfig(ctx context.Context, other map[string]any, site site.API) (*E
// NewEEBus creates EEBus HEMS // NewEEBus creates EEBus HEMS
func NewEEBus(ctx context.Context, ski string, limits Limits, passthrough func(bool) error, site site.API, interval time.Duration) (*EEBus, error) { func NewEEBus(ctx context.Context, ski string, limits Limits, passthrough func(bool) error, site site.API, interval time.Duration) (*EEBus, error) {
if eebus.Instance == nil { inst, err := eebus.Instance()
return nil, errors.New("eebus not configured") if err != nil {
return nil, err
} }
c := &EEBus{ c := &EEBus{
log: util.NewLogger("eebus"), log: util.NewLogger("eebus"),
site: site, site: site,
passthrough: passthrough, passthrough: passthrough,
cs: eebus.Instance.ControllableSystem(), cs: inst.ControllableSystem(),
Connector: eebus.NewConnector(), Connector: eebus.NewConnector(),
heartbeat: util.NewValue[struct{}](2 * time.Minute), // LPC-031 heartbeat: util.NewValue[struct{}](2 * time.Minute), // LPC-031
interval: interval, interval: interval,
@ -111,12 +111,12 @@ func NewEEBus(ctx context.Context, ski string, limits Limits, passthrough func(b
// otherwise a heartbeat timeout is assumed when the state machine is called for the first time // otherwise a heartbeat timeout is assumed when the state machine is called for the first time
c.heartbeat.Set(struct{}{}) c.heartbeat.Set(struct{}{})
if err := eebus.Instance.RegisterDevice(ski, "", c); err != nil { if err := inst.RegisterDevice(ski, "", c); err != nil {
return nil, err return nil, err
} }
if err := c.Wait(ctx); err != nil { if err := c.Wait(ctx); err != nil {
eebus.Instance.UnregisterDevice(ski, c) inst.UnregisterDevice(ski, c)
return nil, err return nil, err
} }

View file

@ -93,11 +93,12 @@ func NewEEBusFromConfig(ctx context.Context, other map[string]any) (api.Meter, e
// NewEEBus creates an EEBus meter // NewEEBus creates an EEBus meter
// Uses MGCP only when usage="grid", otherwise uses MPC (default) // Uses MGCP only when usage="grid", otherwise uses MPC (default)
func NewEEBus(ctx context.Context, ski, ip string, usage *templates.Usage) (api.Meter, error) { func NewEEBus(ctx context.Context, ski, ip string, usage *templates.Usage) (api.Meter, error) {
if eebus.Instance == nil { inst, err := eebus.Instance()
return nil, errors.New("eebus not configured") if err != nil {
return nil, err
} }
ma := eebus.Instance.MonitoringAppliance() ma := inst.MonitoringAppliance()
// Use MGCP only for explicit grid usage, MPC for everything else (default) // Use MGCP only for explicit grid usage, MPC for everything else (default)
useCase := "mpc" useCase := "mpc"
@ -113,25 +114,25 @@ func NewEEBus(ctx context.Context, ski, ip string, usage *templates.Usage) (api.
c := &EEBus{ c := &EEBus{
log: util.NewLogger("eebus-" + useCase), log: util.NewLogger("eebus-" + useCase),
ma: ma, ma: ma,
eg: eebus.Instance.EnergyGuard(), eg: inst.EnergyGuard(),
mm: mm, mm: mm,
scenarios: scenarios, scenarios: scenarios,
connector: eebus.NewConnector(), connector: eebus.NewConnector(),
} }
if err := eebus.Instance.RegisterDevice(ski, ip, c); err != nil { if err := inst.RegisterDevice(ski, ip, c); err != nil {
return nil, err return nil, err
} }
if err := c.connector.Wait(ctx); err != nil { if err := c.connector.Wait(ctx); err != nil {
eebus.Instance.UnregisterDevice(ski, c) inst.UnregisterDevice(ski, c)
return nil, err return nil, err
} }
// unregister device when context is cancelled (e.g. UI config validation) // unregister device when context is cancelled (e.g. UI config validation)
go func() { go func() {
<-ctx.Done() <-ctx.Done()
eebus.Instance.UnregisterDevice(ski, c) inst.UnregisterDevice(ski, c)
}() }()
// monitoring appliance // monitoring appliance

View file

@ -84,14 +84,33 @@ type EEBus struct {
clients map[string][]Device clients map[string][]Device
} }
var Instance *EEBus var (
instance *EEBus
started func() error // memoized service start; set in NewServer, runs once
)
// Instance returns the eebus server, starting the service once on first call
// (OCPP pattern). Returns an error if eebus is unconfigured or start fails.
func Instance() (*EEBus, error) {
if instance == nil {
return nil, errors.New("eebus not configured")
}
if err := started(); err != nil {
return nil, err
}
return instance, nil
}
func GetStatus() any { func GetStatus() any {
var ski string
if instance != nil {
ski = instance.Ski()
}
return struct { return struct {
Ski string `json:"ski"` Ski string `json:"ski"`
QR string `json:"qr,omitempty"` QR string `json:"qr,omitempty"`
}{ }{
Ski: Ski(), Ski: ski,
QR: qrCode(), QR: qrCode(),
} }
} }
@ -99,10 +118,10 @@ func GetStatus() any {
// qrCode returns the SHIP installation QR code text per EEBus SHIP installation // qrCode returns the SHIP installation QR code text per EEBus SHIP installation
// requirements, or empty if unavailable // requirements, or empty if unavailable
func qrCode() string { func qrCode() string {
if Instance == nil { if instance == nil {
return "" return ""
} }
qr, err := Instance.service.QRCodeText() qr, err := instance.service.QRCodeText()
if err != nil { if err != nil {
return "" return ""
} }
@ -228,9 +247,17 @@ func NewServer(other Config) (*EEBus, error) {
c.service.AddUseCase(uc) c.service.AddUseCase(uc)
} }
started = sync.OnceValue(c.service.Start)
instance = c
return c, nil return c, nil
} }
// Ski returns the local service SKI.
func (c *EEBus) Ski() string {
return c.ski
}
func (c *EEBus) RegisterDevice(ski, ip string, device Device) error { func (c *EEBus) RegisterDevice(ski, ip string, device Device) error {
ski = shiputil.NormalizeSKI(ski) ski = shiputil.NormalizeSKI(ski)
c.log.TRACE.Printf("registering ski: %s", ski) c.log.TRACE.Printf("registering ski: %s", ski)
@ -296,10 +323,6 @@ func (c *EEBus) RemoteServices() []shipapi.RemoteMdnsService {
return c.remoteServices return c.remoteServices
} }
func (c *EEBus) Run() {
c.service.Start()
}
func (c *EEBus) Shutdown() { func (c *EEBus) Shutdown() {
c.service.Shutdown() c.service.Shutdown()
} }

View file

@ -16,8 +16,8 @@ func init() {
func getServices(w http.ResponseWriter, req *http.Request) { func getServices(w http.ResponseWriter, req *http.Request) {
var res []string var res []string
if Instance != nil { if instance != nil {
for _, s := range Instance.RemoteServices() { for _, s := range instance.RemoteServices() {
res = append(res, s.Ski) res = append(res, s.Ski)
} }
} }

View file

@ -27,7 +27,7 @@ func TestEEBus(t *testing.T) {
public, private, err := server.GetX509KeyPair(certificate) public, private, err := server.GetX509KeyPair(certificate)
require.NoError(t, err, "decode certificate") require.NoError(t, err, "decode certificate")
srv, err := server.NewServer(server.Config{ _, err = server.NewServer(server.Config{
Certificate: server.Certificate{ Certificate: server.Certificate{
Public: public, Public: public,
Private: private, Private: private,
@ -36,12 +36,12 @@ func TestEEBus(t *testing.T) {
require.NoError(t, err, "server") require.NoError(t, err, "server")
server.Instance = srv inst, err := server.Instance()
go srv.Run() require.NoError(t, err, "instance")
require.NotEmpty(t, server.Ski(), "server ski") require.NotEmpty(t, inst.Ski(), "server ski")
box, err := createControlbox(t.Context(), server.Ski(), remotePort) box, err := createControlbox(t.Context(), inst.Ski(), remotePort)
require.NoError(t, err, "controlbox") require.NoError(t, err, "controlbox")
eventC := make(chan api.EventType, 16) eventC := make(chan api.EventType, 16)

View file

@ -63,8 +63,14 @@ func DefaultConfig(conf *Config) (*Config, error) {
return nil, err return nil, err
} }
// preserve a configured port (e.g. EVCC_EEBUS_PORT), default otherwise
port := conf.Port
if port == 0 {
port = 4712
}
res := Config{ res := Config{
Port: 4712, Port: port,
ShipID: createShipID(), ShipID: createShipID(),
Certificate: Certificate{ Certificate: Certificate{
Public: public, Public: public,
@ -74,11 +80,3 @@ func DefaultConfig(conf *Config) (*Config, error) {
return &res, nil return &res, nil
} }
// Ski returns the EEbus server SKI
func Ski() string {
if Instance == nil {
return ""
}
return Instance.ski
}

View file

@ -1,5 +1,5 @@
import { test, expect } from "@playwright/test"; import { test, expect } from "@playwright/test";
import { start, stop, baseUrl, restart } from "./evcc"; import { start, stop, baseUrl, restart, eebusPort } from "./evcc";
import { expectModalHidden, expectModalVisible } from "./utils"; import { expectModalHidden, expectModalVisible } from "./utils";
test.use({ baseURL: baseUrl() }); test.use({ baseURL: baseUrl() });
@ -21,7 +21,7 @@ test.describe("eebus", async () => {
await expect(modal.getByLabel("SKI")).not.toBeEmpty(); await expect(modal.getByLabel("SKI")).not.toBeEmpty();
await page.getByRole("button", { name: "Show advanced settings" }).click(); await page.getByRole("button", { name: "Show advanced settings" }).click();
await expect(modal.getByLabel("Port")).toHaveValue("4712"); await expect(modal.getByLabel("Port")).toHaveValue(String(eebusPort()));
await expect(modal.getByLabel("Interfaces")).toBeVisible(); await expect(modal.getByLabel("Interfaces")).toBeVisible();
await expect(modal.getByLabel("Public certificate")).not.toBeEmpty(); await expect(modal.getByLabel("Public certificate")).not.toBeEmpty();
await expect(modal.getByLabel("Private key")).toHaveValue("***"); await expect(modal.getByLabel("Private key")).toHaveValue("***");

View file

@ -25,6 +25,11 @@ function ocppPort() {
return 12000 + index; return 12000 + index;
} }
export function eebusPort() {
const index = Number(process.env["TEST_WORKER_INDEX"] ?? 0);
return 13000 + index;
}
function logPrefix() { function logPrefix() {
return `[worker:${process.env["TEST_WORKER_INDEX"]}]`; return `[worker:${process.env["TEST_WORKER_INDEX"]}]`;
} }
@ -105,19 +110,21 @@ async function _start(config?: string, flags: string | string[] = []) {
const configArgs = config ? ["--config", config.includes("/") ? config : `tests/${config}`] : []; const configArgs = config ? ["--config", config.includes("/") ? config : `tests/${config}`] : [];
const port = workerPort(); const port = workerPort();
const ocpp = ocppPort(); const ocpp = ocppPort();
const eebus = eebusPort();
log(`wait until port ${port} is available`); log(`wait until port ${port} is available`);
// wait for port to be available // wait for port to be available
await waitOn({ resources: [`tcp:${port}`], reverse: true, log: LOG_ENABLED }); await waitOn({ resources: [`tcp:${port}`], reverse: true, log: LOG_ENABLED });
log(`starting evcc on ports ${port} (HTTP) and ${ocpp} (OCPP)`); log(`starting evcc on ports ${port} (HTTP), ${ocpp} (OCPP) and ${eebus} (EEBus)`);
const additionalFlags = typeof flags === "string" ? [flags] : flags; const additionalFlags = typeof flags === "string" ? [flags] : flags;
additionalFlags.push("--log", "debug,httpd:trace"); additionalFlags.push("--log", "debug,httpd:trace");
log("starting evcc", { config, port, ocpp, additionalFlags }); log("starting evcc", { config, port, ocpp, eebus, additionalFlags });
const instance = spawn(BINARY, [...configArgs, ...additionalFlags], { const instance = spawn(BINARY, [...configArgs, ...additionalFlags], {
env: { env: {
...process.env, ...process.env,
EVCC_NETWORK_PORT: port.toString(), EVCC_NETWORK_PORT: port.toString(),
EVCC_OCPP_PORT: ocpp.toString(), EVCC_OCPP_PORT: ocpp.toString(),
EVCC_EEBUS_PORT: eebus.toString(),
EVCC_DATABASE_DSN: dbPath(), EVCC_DATABASE_DSN: dbPath(),
}, },
stdio: ["pipe", "pipe", "pipe"], stdio: ["pipe", "pipe", "pipe"],