From fc840629052cd6b4c0f8c542d6681aab02f1b391 Mon Sep 17 00:00:00 2001 From: andig Date: Thu, 25 Jun 2026 17:26:08 +0200 Subject: [PATCH] OCPP/EEBus: lazy-start servers so the CLI works against a running evcc (#30839) Co-authored-by: Michael Geers --- charger/eebus.go | 14 ++-- charger/ocpp.go | 9 ++- charger/ocpp/connector.go | 2 +- charger/ocpp/connector_core.go | 2 +- charger/ocpp/connector_test.go | 2 +- charger/ocpp/cp.go | 6 +- charger/ocpp/cp_core_test.go | 16 ++--- charger/ocpp/cp_requests.go | 14 ++-- charger/ocpp/cs.go | 4 +- charger/ocpp/forwarder.go | 15 ++++- charger/ocpp/instance.go | 120 ++++++++++++++++++--------------- charger/ocpp/instance_test.go | 2 +- charger/ocpp_test.go | 8 ++- cmd/root.go | 55 +++++++++------ cmd/setup.go | 16 +++-- hems/eebus/eebus.go | 12 ++-- meter/eebus.go | 15 +++-- server/eebus/eebus.go | 39 ++++++++--- server/eebus/service.go | 4 +- server/eebus/test/cs_test.go | 10 +-- server/eebus/types.go | 16 ++--- tests/config-eebus.spec.ts | 4 +- tests/evcc.ts | 11 ++- 23 files changed, 237 insertions(+), 159 deletions(-) diff --git a/charger/eebus.go b/charger/eebus.go index 711ac9cae..cc3556ec1 100644 --- a/charger/eebus.go +++ b/charger/eebus.go @@ -2,7 +2,6 @@ package charger import ( "context" - "errors" "fmt" "sync" "time" @@ -74,33 +73,34 @@ func NewEEBusFromConfig(ctx context.Context, other map[string]any) (api.Charger, // newEEBus creates and initializes a raw *EEBus charger. // It registers the device with the EEBus instance and waits for the connection. func newEEBus(ctx context.Context, ski, ip string) (*EEBus, error) { - if eebus.Instance == nil { - return nil, errors.New("eebus not configured") + inst, err := eebus.Instance() + if err != nil { + return nil, err } c := &EEBus{ Caps: implement.New(), log: util.NewLogger("eebus"), current: 6, - cem: eebus.Instance.CustomerEnergyManagement(), + cem: inst.CustomerEnergyManagement(), } c.connector = eebus.NewConnector() 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 } if err := c.connector.Wait(ctx); err != nil { - eebus.Instance.UnregisterDevice(ski, c) + inst.UnregisterDevice(ski, c) return nil, err } // unregister device when context is cancelled (e.g. UI config validation) go func() { <-ctx.Done() - eebus.Instance.UnregisterDevice(ski, c) + inst.UnregisterDevice(ski, c) }() return c, nil diff --git a/charger/ocpp.go b/charger/ocpp.go index 36c466121..91c8f9ae2 100644 --- a/charger/ocpp.go +++ b/charger/ocpp.go @@ -139,9 +139,14 @@ func NewOCPP(ctx context.Context, ) (*OCPP, error) { 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 { - return ocpp.NewChargePoint(log, id) + return ocpp.NewChargePoint(log, cs, id) }, func(cp *ocpp.CP) error { log.DEBUG.Printf("waiting for chargepoint: %v", connectTimeout) diff --git a/charger/ocpp/connector.go b/charger/ocpp/connector.go index fa224327d..87512b0fe 100644 --- a/charger/ocpp/connector.go +++ b/charger/ocpp/connector.go @@ -63,7 +63,7 @@ func NewConnector(ctx context.Context, log *util.Logger, id int, cp *CP, idTag s var ok bool // 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 { ok = true } diff --git a/charger/ocpp/connector_core.go b/charger/ocpp/connector_core.go index 15d45e7f2..6d4a5ab6b 100644 --- a/charger/ocpp/connector_core.go +++ b/charger/ocpp/connector_core.go @@ -111,7 +111,7 @@ func (conn *Connector) OnStartTransaction(request *core.StartTransactionRequest) conn.mu.Lock() defer conn.mu.Unlock() - conn.txnId = int(instance.txnId.Add(1)) + conn.txnId = int(conn.cp.cs.txnId.Add(1)) conn.idTag = request.IdTag res := &core.StartTransactionConfirmation{ diff --git a/charger/ocpp/connector_test.go b/charger/ocpp/connector_test.go index 139ebcd1a..3d79678ef 100644 --- a/charger/ocpp/connector_test.go +++ b/charger/ocpp/connector_test.go @@ -26,7 +26,7 @@ type connTestSuite struct { func (suite *connTestSuite) SetupTest() { // setup 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.clock = clock.NewMock() diff --git a/charger/ocpp/cp.go b/charger/ocpp/cp.go index f200eceae..b4c780055 100644 --- a/charger/ocpp/cp.go +++ b/charger/ocpp/cp.go @@ -16,6 +16,7 @@ import ( type CP struct { mu sync.RWMutex + cs *CS // central system this charge point is registered with log *util.Logger onceConnect sync.Once onceMonitor sync.Once @@ -43,8 +44,9 @@ type CP struct { connectors map[int]*Connector } -func NewChargePoint(log *util.Logger, id string) *CP { +func NewChargePoint(log *util.Logger, cs *CS, id string) *CP { return &CP{ + cs: cs, log: log, id: id, @@ -171,7 +173,7 @@ func (cp *CP) onTransportConnect() { cp.log.DEBUG.Printf("proactively triggering BootNotification") - if err := Instance().TriggerMessage( + if err := cp.cs.TriggerMessage( cp.id, func(conf *remotetrigger.TriggerMessageConfirmation, err error) { if err != nil { diff --git a/charger/ocpp/cp_core_test.go b/charger/ocpp/cp_core_test.go index b68624ee8..f6125690f 100644 --- a/charger/ocpp/cp_core_test.go +++ b/charger/ocpp/cp_core_test.go @@ -14,7 +14,7 @@ import ( func TestBootNotificationStoresResultAndConnects(t *testing.T) { 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.Nil(t, cp.BootNotificationResult, "should have no boot result initially") @@ -43,7 +43,7 @@ func TestBootNotificationStoresResultAndConnects(t *testing.T) { func TestBootNotificationStopsTimer(t *testing.T) { log := util.NewLogger("test") - cp := NewChargePoint(log, "test-cp") + cp := NewChargePoint(log, instance, "test-cp") // simulate WebSocket connect (starts timer) cp.onTransportConnect() @@ -73,7 +73,7 @@ func TestBootNotificationStopsTimer(t *testing.T) { func TestTransportConnectTimeoutFallback(t *testing.T) { log := util.NewLogger("test") - cp := NewChargePoint(log, "test-cp") + cp := NewChargePoint(log, instance, "test-cp") // use a short timeout for testing origTimeout := Timeout @@ -99,7 +99,7 @@ func TestTransportConnectTimeoutFallback(t *testing.T) { func TestDisconnectCancelsTimer(t *testing.T) { log := util.NewLogger("test") - cp := NewChargePoint(log, "test-cp") + cp := NewChargePoint(log, instance, "test-cp") // use a short timeout for testing origTimeout := Timeout @@ -126,7 +126,7 @@ func TestDisconnectCancelsTimer(t *testing.T) { func TestBootNotificationChannelCoalesces(t *testing.T) { 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 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. func TestBootNotificationRebootLoopKeepsLatest(t *testing.T) { log := util.NewLogger("test") - cp := NewChargePoint(log, "test-cp") + cp := NewChargePoint(log, instance, "test-cp") for i := range 5 { _, err := cp.OnBootNotification(&core.BootNotificationRequest{ @@ -188,7 +188,7 @@ func TestBootNotificationRebootLoopKeepsLatest(t *testing.T) { func TestReconnectAfterReboot(t *testing.T) { log := util.NewLogger("test") - cp := NewChargePoint(log, "test-cp") + cp := NewChargePoint(log, instance, "test-cp") // simulate initial connection cp.connect(true) @@ -220,7 +220,7 @@ func TestReconnectAfterReboot(t *testing.T) { func TestMonitorRebootOnlyOnce(t *testing.T) { log := util.NewLogger("test") - cp := NewChargePoint(log, "test-cp") + cp := NewChargePoint(log, instance, "test-cp") ctx := t.Context() var callCount atomic.Int32 diff --git a/charger/ocpp/cp_requests.go b/charger/ocpp/cp_requests.go index 00850c77e..40f717966 100644 --- a/charger/ocpp/cp_requests.go +++ b/charger/ocpp/cp_requests.go @@ -13,7 +13,7 @@ import ( func (cp *CP) ChangeAvailabilityRequest(connectorId int, availabilityType core.AvailabilityType) error { 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 { err = errors.New(string(request.Status)) } @@ -28,7 +28,7 @@ func (cp *CP) GetCompositeScheduleRequest(connectorId int, duration int) (*smart var res *smartcharging.GetCompositeScheduleConfirmation 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 { 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 { 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 { 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 { 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 { err = errors.New(string(request.Status)) } @@ -79,7 +79,7 @@ func (cp *CP) TriggerMessageRequest(connectorId int, requestedMessage remotetrig 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 { 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 { 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 { err = errors.New(string(request.Status)) } @@ -112,7 +112,7 @@ func (cp *CP) GetConfigurationRequest() (*core.GetConfigurationConfirmation, err rc := make(chan error, 1) 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 rc <- err diff --git a/charger/ocpp/cs.go b/charger/ocpp/cs.go index beca61356..fc1626cd3 100644 --- a/charger/ocpp/cs.go +++ b/charger/ocpp/cs.go @@ -9,6 +9,7 @@ import ( "github.com/evcc-io/evcc/util" ocpp16 "github.com/lorenzodonini/ocpp-go/ocpp1.6" "github.com/lorenzodonini/ocpp-go/ocpp1.6/core" + "github.com/lorenzodonini/ocpp-go/ocppj" "github.com/lorenzodonini/ocpp-go/ws" ) @@ -30,7 +31,8 @@ type CS struct { regs map[string]*registration // guarded by mu mutex txnId atomic.Int64 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. diff --git a/charger/ocpp/forwarder.go b/charger/ocpp/forwarder.go index 89365391a..598a3501c 100644 --- a/charger/ocpp/forwarder.go +++ b/charger/ocpp/forwarder.go @@ -420,7 +420,10 @@ func drainPendingWithErrors(id string, sc *sidecar) { delete(pendingMsgs, id) pendingMu.Unlock() - cs := Instance() + cs, err := Instance() + if err != nil { + return + } for _, frame := range buffered { msgType, msgID, action, err := parseOCPPFrame(frame) @@ -587,6 +590,12 @@ func (sc *sidecar) readFromUpstream() { notifyUpdated() }() + cs, err := Instance() + if err != nil { + forwarderLog.ERROR.Printf("forwarder: central system unavailable for %s: %v", sc.chargerID, err) + return + } + for { _, msg, err := sc.conn.Read(context.Background()) if err != nil { @@ -645,7 +654,7 @@ func (sc *sidecar) readFromUpstream() { sc.pendingUpstreamCalls[msgID] = struct{}{} 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) } @@ -660,7 +669,7 @@ func (sc *sidecar) readFromUpstream() { if isChargerCall { // 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) } continue diff --git a/charger/ocpp/instance.go b/charger/ocpp/instance.go index dfb3aadeb..5d5351dca 100644 --- a/charger/ocpp/instance.go +++ b/charger/ocpp/instance.go @@ -1,6 +1,7 @@ package ocpp import ( + "errors" "fmt" "net/http" "net/url" @@ -43,8 +44,8 @@ func (r ForwarderRule) Redacted() ForwarderRule { } var ( - once sync.Once instance *CS + started func() error // memoized listen; set in NewServer, runs once port = 8887 boundPort int externalUrl string @@ -129,64 +130,77 @@ func CurrentConfig() Config { return Config{Port: port} } -// Init initializes the OCPP server -func Init(cfg Config, networkExternalUrl string) { +// NewServer builds the OCPP central system without starting it. +func NewServer(cfg Config, networkExternalUrl string) { port = cfg.Port externalUrl = networkExternalUrl -} -func Instance() *CS { - once.Do(func() { - log := util.NewLogger("ocpp") + log := util.NewLogger("ocpp") - server := &interceptingServer{Server: ws.NewServer()} - server.SetCheckOriginHandler(func(r *http.Request) bool { return true }) + server := &interceptingServer{Server: ws.NewServer()} + server.SetCheckOriginHandler(func(r *http.Request) bool { return true }) - dispatcher := ocppj.NewDefaultServerDispatcher(ocppj.NewFIFOQueueMap(0)) - dispatcher.SetTimeout(Timeout) + dispatcher := ocppj.NewDefaultServerDispatcher(ocppj.NewFIFOQueueMap(0)) - 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 { - log.ERROR.Printf("%v (%s)", err, rawMessage) - 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 + 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 { + log.ERROR.Printf("%v (%s)", err, rawMessage) + return nil }) - 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() } diff --git a/charger/ocpp/instance_test.go b/charger/ocpp/instance_test.go index 1e8ffefab..c6b5d316e 100644 --- a/charger/ocpp/instance_test.go +++ b/charger/ocpp/instance_test.go @@ -9,7 +9,7 @@ func TestMain(m *testing.M) { // 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 // when both run in parallel under `go test ./...` - Init(Config{Port: 0}, "") + NewServer(Config{Port: 0}, "") os.Exit(m.Run()) } diff --git a/charger/ocpp_test.go b/charger/ocpp_test.go index 64c486689..5e800167f 100644 --- a/charger/ocpp_test.go +++ b/charger/ocpp_test.go @@ -35,7 +35,7 @@ func TestMain(m *testing.M) { // 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 // 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()) } @@ -57,7 +57,10 @@ func (suite *ocppTestSuite) SetupSuite() { ocpp.TriggerBootDelay = 100 * time.Millisecond // 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()} ocppj.SetLogger(suite.logger) @@ -65,7 +68,6 @@ func (suite *ocppTestSuite) SetupSuite() { ocppTestUrl = fmt.Sprintf("ws://localhost:%d", ocpp.Port()) suite.clock = clock.NewMock() - suite.NotNil(ocpp.Instance()) } func (suite *ocppTestSuite) TearDownSuite() { diff --git a/cmd/root.go b/cmd/root.go index c3270ea00..1afa88364 100644 --- a/cmd/root.go +++ b/cmd/root.go @@ -200,28 +200,39 @@ func runRoot(cmd *cobra.Command, args []string) { valueChan := make(chan util.Param, 64) go tee.Run(valueChan) - // start OCPP server - ocppCS := ocpp.Instance() - ocppCS.SetUpdated(func() { - // republish when OCPP state updates - valueChan <- util.Param{Key: keys.Ocpp, Val: globalconfig.ConfigStatus{ - Config: ocpp.CurrentConfig(), - Status: ocpp.GetStatus(), - }} - }) - log.INFO.Printf("OCPP local url: ws://127.0.0.1:%d/", conf.Ocpp.Port) - if ocpp.ExternalUrl() != "" { - log.INFO.Printf("OCPP external url: %s/", ocpp.ExternalUrl()) - } - // register the callback even with no rules so runtime additions are pushed - ocpp.SetForwarderUpdated(func() { - 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())) + // start OCPP and EEBus servers (skipped in degraded mode where setup failed, + // so a misconfigured instance serves only the offline UI) + if err == nil { + cs, ocppErr := ocpp.Instance() + if ocppErr != nil { + log.ERROR.Println("ocpp:", ocppErr) + } else { + cs.SetUpdated(func() { + // republish when OCPP state updates + valueChan <- util.Param{Key: keys.Ocpp, Val: globalconfig.ConfigStatus{ + Config: ocpp.CurrentConfig(), + Status: ocpp.GetStatus(), + }} + }) + log.INFO.Printf("OCPP local url: ws://127.0.0.1:%d/", conf.Ocpp.Port) + if ocpp.ExternalUrl() != "" { + log.INFO.Printf("OCPP external url: %s/", ocpp.ExternalUrl()) + } + } + // register the callback even with no rules so runtime additions are pushed + ocpp.SetForwarderUpdated(func() { + 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 diff --git a/cmd/setup.go b/cmd/setup.go index 17205ac62..852610751 100644 --- a/cmd/setup.go +++ b/cmd/setup.go @@ -853,11 +853,16 @@ func configureOCPP(cfg *ocpp.Config, externalUrl string) { 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. 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) } } @@ -885,13 +890,12 @@ func configureEEBus(conf *eebus.Config) error { } } - var err error - if eebus.Instance, err = eebus.NewServer(*conf); err != nil { + srv, err := eebus.NewServer(*conf) + if err != nil { return fmt.Errorf("failed configuring eebus: %w", err) } - eebus.Instance.Run() - shutdown.Register(eebus.Instance.Shutdown) + shutdown.Register(srv.Shutdown) return nil } diff --git a/hems/eebus/eebus.go b/hems/eebus/eebus.go index 259c95e95..5a1c43551 100644 --- a/hems/eebus/eebus.go +++ b/hems/eebus/eebus.go @@ -2,7 +2,6 @@ package eebus import ( "context" - "errors" "sync" "time" @@ -89,15 +88,16 @@ func NewFromConfig(ctx context.Context, other map[string]any, site site.API) (*E // 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) { - if eebus.Instance == nil { - return nil, errors.New("eebus not configured") + inst, err := eebus.Instance() + if err != nil { + return nil, err } c := &EEBus{ log: util.NewLogger("eebus"), site: site, passthrough: passthrough, - cs: eebus.Instance.ControllableSystem(), + cs: inst.ControllableSystem(), Connector: eebus.NewConnector(), heartbeat: util.NewValue[struct{}](2 * time.Minute), // LPC-031 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 c.heartbeat.Set(struct{}{}) - if err := eebus.Instance.RegisterDevice(ski, "", c); err != nil { + if err := inst.RegisterDevice(ski, "", c); err != nil { return nil, err } if err := c.Wait(ctx); err != nil { - eebus.Instance.UnregisterDevice(ski, c) + inst.UnregisterDevice(ski, c) return nil, err } diff --git a/meter/eebus.go b/meter/eebus.go index f06a6ec0e..c9069b12a 100644 --- a/meter/eebus.go +++ b/meter/eebus.go @@ -93,11 +93,12 @@ func NewEEBusFromConfig(ctx context.Context, other map[string]any) (api.Meter, e // NewEEBus creates an EEBus meter // Uses MGCP only when usage="grid", otherwise uses MPC (default) func NewEEBus(ctx context.Context, ski, ip string, usage *templates.Usage) (api.Meter, error) { - if eebus.Instance == nil { - return nil, errors.New("eebus not configured") + inst, err := eebus.Instance() + 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) useCase := "mpc" @@ -113,25 +114,25 @@ func NewEEBus(ctx context.Context, ski, ip string, usage *templates.Usage) (api. c := &EEBus{ log: util.NewLogger("eebus-" + useCase), ma: ma, - eg: eebus.Instance.EnergyGuard(), + eg: inst.EnergyGuard(), mm: mm, scenarios: scenarios, 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 } if err := c.connector.Wait(ctx); err != nil { - eebus.Instance.UnregisterDevice(ski, c) + inst.UnregisterDevice(ski, c) return nil, err } // unregister device when context is cancelled (e.g. UI config validation) go func() { <-ctx.Done() - eebus.Instance.UnregisterDevice(ski, c) + inst.UnregisterDevice(ski, c) }() // monitoring appliance diff --git a/server/eebus/eebus.go b/server/eebus/eebus.go index f672de997..57afae406 100644 --- a/server/eebus/eebus.go +++ b/server/eebus/eebus.go @@ -84,14 +84,33 @@ type EEBus struct { 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 { + var ski string + if instance != nil { + ski = instance.Ski() + } return struct { Ski string `json:"ski"` QR string `json:"qr,omitempty"` }{ - Ski: Ski(), + Ski: ski, QR: qrCode(), } } @@ -99,10 +118,10 @@ func GetStatus() any { // qrCode returns the SHIP installation QR code text per EEBus SHIP installation // requirements, or empty if unavailable func qrCode() string { - if Instance == nil { + if instance == nil { return "" } - qr, err := Instance.service.QRCodeText() + qr, err := instance.service.QRCodeText() if err != nil { return "" } @@ -228,9 +247,17 @@ func NewServer(other Config) (*EEBus, error) { c.service.AddUseCase(uc) } + started = sync.OnceValue(c.service.Start) + instance = c + 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 { ski = shiputil.NormalizeSKI(ski) c.log.TRACE.Printf("registering ski: %s", ski) @@ -296,10 +323,6 @@ func (c *EEBus) RemoteServices() []shipapi.RemoteMdnsService { return c.remoteServices } -func (c *EEBus) Run() { - c.service.Start() -} - func (c *EEBus) Shutdown() { c.service.Shutdown() } diff --git a/server/eebus/service.go b/server/eebus/service.go index 9e7592ca1..687241e59 100644 --- a/server/eebus/service.go +++ b/server/eebus/service.go @@ -16,8 +16,8 @@ func init() { func getServices(w http.ResponseWriter, req *http.Request) { var res []string - if Instance != nil { - for _, s := range Instance.RemoteServices() { + if instance != nil { + for _, s := range instance.RemoteServices() { res = append(res, s.Ski) } } diff --git a/server/eebus/test/cs_test.go b/server/eebus/test/cs_test.go index a25adbdf2..46c5ea55d 100644 --- a/server/eebus/test/cs_test.go +++ b/server/eebus/test/cs_test.go @@ -27,7 +27,7 @@ func TestEEBus(t *testing.T) { public, private, err := server.GetX509KeyPair(certificate) require.NoError(t, err, "decode certificate") - srv, err := server.NewServer(server.Config{ + _, err = server.NewServer(server.Config{ Certificate: server.Certificate{ Public: public, Private: private, @@ -36,12 +36,12 @@ func TestEEBus(t *testing.T) { require.NoError(t, err, "server") - server.Instance = srv - go srv.Run() + inst, err := server.Instance() + 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") eventC := make(chan api.EventType, 16) diff --git a/server/eebus/types.go b/server/eebus/types.go index c1ea9964a..d9ad47f0d 100644 --- a/server/eebus/types.go +++ b/server/eebus/types.go @@ -63,8 +63,14 @@ func DefaultConfig(conf *Config) (*Config, error) { return nil, err } + // preserve a configured port (e.g. EVCC_EEBUS_PORT), default otherwise + port := conf.Port + if port == 0 { + port = 4712 + } + res := Config{ - Port: 4712, + Port: port, ShipID: createShipID(), Certificate: Certificate{ Public: public, @@ -74,11 +80,3 @@ func DefaultConfig(conf *Config) (*Config, error) { return &res, nil } - -// Ski returns the EEbus server SKI -func Ski() string { - if Instance == nil { - return "" - } - return Instance.ski -} diff --git a/tests/config-eebus.spec.ts b/tests/config-eebus.spec.ts index fcc16fb93..0c9050ee2 100644 --- a/tests/config-eebus.spec.ts +++ b/tests/config-eebus.spec.ts @@ -1,5 +1,5 @@ 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"; test.use({ baseURL: baseUrl() }); @@ -21,7 +21,7 @@ test.describe("eebus", async () => { await expect(modal.getByLabel("SKI")).not.toBeEmpty(); 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("Public certificate")).not.toBeEmpty(); await expect(modal.getByLabel("Private key")).toHaveValue("***"); diff --git a/tests/evcc.ts b/tests/evcc.ts index 5ce51b8ca..94cd59f7c 100644 --- a/tests/evcc.ts +++ b/tests/evcc.ts @@ -25,6 +25,11 @@ function ocppPort() { return 12000 + index; } +export function eebusPort() { + const index = Number(process.env["TEST_WORKER_INDEX"] ?? 0); + return 13000 + index; +} + function logPrefix() { 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 port = workerPort(); const ocpp = ocppPort(); + const eebus = eebusPort(); log(`wait until port ${port} is available`); // wait for port to be available 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; 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], { env: { ...process.env, EVCC_NETWORK_PORT: port.toString(), EVCC_OCPP_PORT: ocpp.toString(), + EVCC_EEBUS_PORT: eebus.toString(), EVCC_DATABASE_DSN: dbPath(), }, stdio: ["pipe", "pipe", "pipe"],