diff --git a/AGENTS.md b/AGENTS.md
index 38fffe45b..e78c50b56 100644
--- a/AGENTS.md
+++ b/AGENTS.md
@@ -33,6 +33,7 @@ Deep documentation on specific subsystems is available in `docs/agents/`. Load w
| [Core Domain](docs/agents/core-domain.md) | Control loop, loadpoint logic, PV surplus, charge modes, tariffs, interfaces |
| [Hardware Integrations](docs/agents/hardware-integrations.md) | Charger/meter/vehicle implementations, adding new devices |
| [Easee Architecture](docs/agents/easee-architecture.md) | Easee charger (REST+SignalR, async correlation, concurrency) |
+| [OCPP Forwarder](docs/agents/ocpp-forwarder.md) | OCPP proxy/forwarder (sidecar relay to upstream OCPP server, read-only mode) |
| [Plugin System](docs/agents/plugin-system.md) | Plugin layer (HTTP, MQTT, Modbus, SunSpec, JS) |
| [Web UI & API](docs/agents/web-ui-api.md) | REST API, WebSocket, Vue frontend, authentication |
| [API Security](docs/agents/api-security.md) | Auth modes, JWT/API key/session, two-tier checks, credential storage |
diff --git a/assets/css/app.css b/assets/css/app.css
index 944048c48..22caccc07 100644
--- a/assets/css/app.css
+++ b/assets/css/app.css
@@ -73,9 +73,6 @@
--bs-danger: var(--evcc-red);
--bs-danger-rgb: var(--evcc-red-rgb);
- --bs-danger: var(--evcc-red);
- --bs-danger-rgb: var(--evcc-red-rgb);
-
--bs-form-invalid-border-color: var(--bs-danger);
--bs-body-font-size: 14px;
diff --git a/assets/js/components/Config/JsonModal.vue b/assets/js/components/Config/JsonModal.vue
index 4db6182f9..d62696227 100644
--- a/assets/js/components/Config/JsonModal.vue
+++ b/assets/js/components/Config/JsonModal.vue
@@ -166,7 +166,8 @@ export default {
if (shouldClose) {
await closeModal();
} else {
- await this.load();
+ // keep open: saved values become the new baseline (no longer dirty)
+ this.serverValues = deepClone(this.values);
}
}
if (res.status === 400) {
diff --git a/assets/js/components/Config/OcppForwarderButton.vue b/assets/js/components/Config/OcppForwarderButton.vue
new file mode 100644
index 000000000..8e3f12d38
--- /dev/null
+++ b/assets/js/components/Config/OcppForwarderButton.vue
@@ -0,0 +1,107 @@
+
+
+ {{ host }}
+
+
+
+
+
+
+
diff --git a/assets/js/components/Config/OcppForwarderModal.vue b/assets/js/components/Config/OcppForwarderModal.vue
new file mode 100644
index 000000000..f82f4b71d
--- /dev/null
+++ b/assets/js/components/Config/OcppForwarderModal.vue
@@ -0,0 +1,280 @@
+
+
+ void;
+ }"
+ >
+
+ {{ $t("config.ocppforwarder.status") }}:
+ {{
+ connectionLabel
+ }}
+
+
+ {{ sessionError }}
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
diff --git a/assets/js/components/Config/OcppModal.vue b/assets/js/components/Config/OcppModal.vue
index 1a41b2533..9ac301e1e 100644
--- a/assets/js/components/Config/OcppModal.vue
+++ b/assets/js/components/Config/OcppModal.vue
@@ -23,52 +23,40 @@
+
{{ $t("config.ocpp.noChargers") }}
-
-
-
{{ $t("config.ocpp.connectionStatus") }}
-
{{ $t("config.ocpp.connectionStatusHelp") }}
+
+
+
{{ $t("config.ocpp.stations") }}
+
{{ $t("config.ocpp.stationsHelp") }}
-
+
-
-
{{ station.id }}
-
- {{ $t(`config.ocpp.status.${station.status}`) }}
-
-
-
-
-
-
-
-
{{ $t("config.ocpp.detectedChargers") }}
-
{{ $t("config.ocpp.detectedHelp") }}
-
-
- -
-
{{ station.id }}
-
- {{ $t(`config.ocpp.status.${station.status}`) }}
-
+
+
+
+ {{ entry.title }}
+
+
{{
+ entry.id
+ }}
+
+
@@ -81,8 +69,25 @@ import { defineComponent, type PropType } from "vue";
import GenericModal from "../Helper/GenericModal.vue";
import FormRow from "./FormRow.vue";
import Markdown from "./Markdown.vue";
-import type { Ocpp, OcppStationStatus } from "@/types/evcc";
+import OcppForwarderButton from "./OcppForwarderButton.vue";
+import StatusIndicator from "./StatusIndicator.vue";
+import type {
+ Ocpp,
+ OcppForwarderRule,
+ OcppForwarderSession,
+ OcppStationStatus,
+} from "@/types/evcc";
+import { OCPP_STATION_STATUS } from "@/types/evcc";
import { getOcppUrl, getOcppUrlWithStationId } from "@/utils/ocpp";
+import store from "@/store";
+
+type StationEntry = {
+ id: string;
+ status: OcppStationStatus["status"];
+ title?: string;
+ rule?: OcppForwarderRule;
+ error?: string;
+};
export default defineComponent({
name: "OcppModal",
@@ -90,46 +95,74 @@ export default defineComponent({
GenericModal,
FormRow,
Markdown,
+ OcppForwarderButton,
+ StatusIndicator,
},
props: {
ocpp: {
type: Object as PropType
,
default: () => ({ config: { port: 0 }, status: { stations: [] } }),
},
+ stationTitles: {
+ type: Object as PropType>,
+ default: () => ({}),
+ },
},
computed: {
- status() {
- return this.ocpp.status;
- },
- stations() {
- return this.status.stations;
- },
ocppUrl(): string {
return getOcppUrl(this.ocpp);
},
ocppUrlWithStationId(): string {
return getOcppUrlWithStationId(this.ocpp);
},
- connectedStations(): OcppStationStatus[] {
- return this.stations.filter(
- (s) => s.status === "connected" || s.status === "configured"
- );
+ rules(): OcppForwarderRule[] {
+ return store.state?.ocppforwarder?.config || [];
},
- detectedStations(): OcppStationStatus[] {
- return this.stations.filter((s) => s.status === "unknown");
+ sessions(): OcppForwarderSession[] {
+ return store.state?.ocppforwarder?.status || [];
+ },
+ // merge of published stations and configured forwarder rules, keyed by id.
+ // a rule without a matching station is shown as "unknown".
+ entries(): StationEntry[] {
+ const byId = new Map();
+ const published = new Set();
+ for (const station of this.ocpp.status.stations) {
+ byId.set(station.id, {
+ id: station.id,
+ status: station.status,
+ title: this.stationTitles[station.id],
+ });
+ published.add(station.id);
+ }
+ for (const rule of this.rules) {
+ if (rule.stationId === "*") continue;
+ const entry = byId.get(rule.stationId) || {
+ id: rule.stationId,
+ status: OCPP_STATION_STATUS.UNKNOWN,
+ title: this.stationTitles[rule.stationId],
+ };
+ entry.rule = rule;
+ entry.error = this.sessions.find((s) => s.chargerId === rule.stationId)?.error;
+ byId.set(rule.stationId, entry);
+ }
+ // published stations first, then rule-only entries; each sorted by id
+ const byIdAlpha = (a: StationEntry, b: StationEntry) => a.id.localeCompare(b.id);
+ const all = [...byId.values()];
+ return [
+ ...all.filter((e) => published.has(e.id)).sort(byIdAlpha),
+ ...all.filter((e) => !published.has(e.id)).sort(byIdAlpha),
+ ];
},
},
methods: {
- statusBadgeClass(status: string): string {
+ statusVariant(status: string): "success" | "warning" | "muted" {
switch (status) {
case "connected":
- return "bg-success";
+ return "success";
case "configured":
- return "bg-warning";
- case "unknown":
- return "bg-secondary";
+ return "warning";
default:
- return "bg-secondary";
+ return "muted";
}
},
},
@@ -142,7 +175,12 @@ export default defineComponent({
padding-left: 0;
}
-.list-group-item code {
- font-size: 0.9rem;
+.station-bar {
+ --bs-border-color: var(--bs-gray-light);
+}
+
+/* min-width:0 lets the identity actually truncate */
+.station-identity {
+ min-width: 0;
}
diff --git a/assets/js/components/Config/StatusIndicator.vue b/assets/js/components/Config/StatusIndicator.vue
new file mode 100644
index 000000000..17c92d037
--- /dev/null
+++ b/assets/js/components/Config/StatusIndicator.vue
@@ -0,0 +1,75 @@
+
+
+
+
+
+
+
+
+
+
diff --git a/assets/js/components/MaterialIcon/OcppForwardStatus.vue b/assets/js/components/MaterialIcon/OcppForwardStatus.vue
new file mode 100644
index 000000000..97b844beb
--- /dev/null
+++ b/assets/js/components/MaterialIcon/OcppForwardStatus.vue
@@ -0,0 +1,37 @@
+
+
+
+
+
diff --git a/assets/js/configModal.ts b/assets/js/configModal.ts
index 715020b3c..68f1256a6 100644
--- a/assets/js/configModal.ts
+++ b/assets/js/configModal.ts
@@ -7,6 +7,7 @@ export interface ModalEntry {
id?: number;
type?: string;
choices?: string[];
+ station?: string;
}
export interface ModalResult {
@@ -106,7 +107,12 @@ function syncAllModals(): void {
// Parse brackets: "meter[type:grid]" => { name: "meter", type: "grid" }
// "meter[choices:pv,battery]" => { name: "meter", choices: ["pv", "battery"] }
-export function parseKey(key: string): { name: string; type?: string; choices?: string[] } {
+export function parseKey(key: string): {
+ name: string;
+ type?: string;
+ choices?: string[];
+ station?: string;
+} {
const bracketMatch = key.match(/^([^[]+)\[([^\]]+)\]$/);
if (!bracketMatch) {
return { name: key };
@@ -126,6 +132,9 @@ export function parseKey(key: string): { name: string; type?: string; choices?:
if (paramKey === "choices") {
return { name, choices: paramValue.split(",") };
}
+ if (paramKey === "station") {
+ return { name, station: paramValue };
+ }
return { name };
}
@@ -158,6 +167,7 @@ export function parseQueryString(queryString: string): ModalEntry[] {
}
if (parsed.type) entry.type = parsed.type;
if (parsed.choices) entry.choices = parsed.choices;
+ if (parsed.station) entry.station = parsed.station;
entries.push(entry);
}
return entries;
@@ -172,6 +182,8 @@ export function buildQuery(stack: ModalEntry[]): Record {
key += `[type:${entry.type}]`;
} else if (entry.choices?.length) {
key += `[choices:${entry.choices.join(",")}]`;
+ } else if (entry.station) {
+ key += `[station:${entry.station}]`;
}
query[key] = entry.id !== undefined ? String(entry.id) : "";
}
@@ -225,7 +237,7 @@ export function initConfigModal(router: Router): void {
export function openModal(
name: string,
- params?: { id?: number; type?: string; choices?: string[] }
+ params?: { id?: number; type?: string; choices?: string[]; station?: string }
): Promise {
if (!_router) {
return Promise.resolve({ action: "cancelled" });
@@ -235,6 +247,7 @@ export function openModal(
if (params?.id !== undefined) entry.id = params.id;
if (params?.type) entry.type = params.type;
if (params?.choices) entry.choices = params.choices;
+ if (params?.station) entry.station = params.station;
const newStack = [...configModal.stack, entry];
const query = buildQuery(newStack);
@@ -273,7 +286,7 @@ export async function closeModal(result?: ModalResult): Promise {
export function replaceModal(
name: string,
- params?: { id?: number; type?: string; choices?: string[] }
+ params?: { id?: number; type?: string; choices?: string[]; station?: string }
): void {
if (!_router) return;
@@ -281,6 +294,7 @@ export function replaceModal(
if (params?.id !== undefined) entry.id = params.id;
if (params?.type) entry.type = params.type;
if (params?.choices) entry.choices = params.choices;
+ if (params?.station) entry.station = params.station;
const newStack = [...configModal.stack.slice(0, -1), entry];
const query = buildQuery(newStack);
diff --git a/assets/js/types/evcc.ts b/assets/js/types/evcc.ts
index c5cbdc156..9d4a3439c 100644
--- a/assets/js/types/evcc.ts
+++ b/assets/js/types/evcc.ts
@@ -121,6 +121,7 @@ export interface State {
config?: string;
database?: string;
ocpp?: Ocpp;
+ ocppforwarder?: ConfigStatus;
optimizer?: boolean;
mcp?: boolean;
}
@@ -137,6 +138,24 @@ export interface OcppConfig {
port: number;
}
+export interface OcppForwarderRule {
+ stationId: string;
+ upstreamUrl: string;
+ password?: string;
+ upstreamStationId?: string;
+ username?: string;
+ insecure?: boolean;
+ caCert?: string;
+ readOnly?: boolean;
+}
+
+export interface OcppForwarderSession {
+ chargerId: string;
+ upstreamUrl: string;
+ upstreamConnected: boolean;
+ error?: string;
+}
+
export interface OcppStatus {
externalUrl?: string;
stations: OcppStationStatus[];
diff --git a/assets/js/views/Config.vue b/assets/js/views/Config.vue
index d858b4071..20d9f2d6c 100644
--- a/assets/js/views/Config.vue
+++ b/assets/js/views/Config.vue
@@ -449,7 +449,8 @@
:yamlSource="eebus?.yamlSource"
@changed="loadDirty"
/>
-
+
+
@@ -479,6 +480,7 @@ import EebusIcon from "../components/MaterialIcon/Eebus.vue";
import EebusModal from "../components/Config/EebusModal.vue";
import OcppIcon from "../components/MaterialIcon/Ocpp.vue";
import OcppModal from "../components/Config/OcppModal.vue";
+import OcppForwarderModal from "../components/Config/OcppForwarderModal.vue";
import formatter from "../mixins/formatter";
import GeneralConfig from "../components/Config/GeneralConfig.vue";
import HemsIcon from "../components/MaterialIcon/Hems.vue";
@@ -571,6 +573,7 @@ export default defineComponent({
EebusModal,
OcppIcon,
OcppModal,
+ OcppForwarderModal,
GeneralConfig,
HemsIcon,
HemsModal,
@@ -856,6 +859,18 @@ export default defineComponent({
}
return { configured: { value: false } };
},
+ // maps an OCPP station id to its loadpoint title (fallback: charger title)
+ stationTitles(): Record {
+ const map: Record = {};
+ this.chargers.forEach((charger) => {
+ const stationId = charger.config?.["stationid"];
+ if (typeof stationId !== "string" || !stationId) return;
+ const loadpoint = this.loadpoints.find((lp) => lp.charger === charger.name);
+ const title = loadpoint?.title || charger.config?.title;
+ if (title) map[stationId] = title;
+ });
+ return map;
+ },
messagingTags(): DeviceTags {
if (this.messagingUiConfigured) {
const events = store.state?.messagingEvents || [];
diff --git a/charger/ocpp/cs.go b/charger/ocpp/cs.go
index bd582892c..beca61356 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/ws"
)
type registration struct {
@@ -29,6 +30,12 @@ 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
+}
+
+// Write sends a raw OCPP frame to the charger with the given station ID.
+func (cs *CS) Write(id string, data []byte) error {
+ return cs.server.Write(id, data)
}
type stationStatus struct {
diff --git a/charger/ocpp/forwarder.go b/charger/ocpp/forwarder.go
new file mode 100644
index 000000000..7f5e4d0d6
--- /dev/null
+++ b/charger/ocpp/forwarder.go
@@ -0,0 +1,746 @@
+package ocpp
+
+// Hybrid OCPP proxy. Chargers connect to evcc's central system; for each charger
+// with a matching rule a "sidecar" WebSocket to the upstream OCPP server runs in
+// parallel. Billing-critical Calls (actionsRelayedToUpstream) are relayed to
+// upstream as the authoritative responder and evcc's handler is bypassed; all
+// other messages are processed by evcc and mirrored to upstream for observation.
+// Upstream Calls are injected into the charger unless the rule is read-only.
+//
+// Design: docs/agents/ocpp-forwarder.md
+
+import (
+ "context"
+ "crypto/tls"
+ "crypto/x509"
+ "encoding/base64"
+ "encoding/json"
+ "fmt"
+ "net/http"
+ "slices"
+ "strconv"
+ "strings"
+ "sync"
+ "time"
+
+ "github.com/coder/websocket"
+ "github.com/evcc-io/evcc/util"
+ "github.com/lorenzodonini/ocpp-go/ocpp1.6/core"
+ "github.com/lorenzodonini/ocpp-go/ocppj"
+ "github.com/lorenzodonini/ocpp-go/ws"
+)
+
+// forwarder rules, guarded by forwarderMu
+var (
+ forwarderMu sync.RWMutex
+ forwarderRules []ForwarderRule
+)
+
+// ForwarderEnabled returns true when at least one rule is configured.
+func ForwarderEnabled() bool {
+ forwarderMu.RLock()
+ defer forwarderMu.RUnlock()
+ return len(forwarderRules) > 0
+}
+
+// ApplyForwarderRules replaces the forwarding rules at runtime and republishes status.
+func ApplyForwarderRules(rules []ForwarderRule) {
+ forwarderMu.Lock()
+ forwarderRules = rules
+ forwarderMu.Unlock()
+
+ valid := make(map[string]bool, len(rules))
+ hasWildcard := false
+ for _, r := range rules {
+ valid[r.StationID] = true
+ if r.StationID == "*" {
+ hasWildcard = true
+ }
+ }
+
+ // drop sidecars/errors for removed rules
+ var stale []*sidecar
+ sidecarsMu.Lock()
+ for id, sc := range sidecars {
+ if !valid[id] && !hasWildcard {
+ stale = append(stale, sc)
+ delete(sidecars, id)
+ }
+ }
+ for id := range forwarderErrors {
+ if !valid[id] && !hasWildcard {
+ delete(forwarderErrors, id)
+ }
+ }
+ sidecarsMu.Unlock()
+ for _, sc := range stale {
+ sc.conn.CloseNow()
+ }
+
+ // connected charger: (re)dial sidecar; otherwise test-dial to surface config errors
+ for _, r := range rules {
+ if r.StationID == "*" || r.UpstreamURL == "" {
+ continue
+ }
+ sidecarsMu.Lock()
+ sc, active := sidecars[r.StationID]
+ connected := connectedChargers[r.StationID]
+ var changed *sidecar
+ if active && !sc.rule.sameConnection(r) {
+ changed = sc
+ delete(sidecars, r.StationID)
+ active = false
+ }
+ sidecarsMu.Unlock()
+ if changed != nil {
+ changed.conn.CloseNow()
+ }
+ if active {
+ continue
+ }
+ if connected {
+ pendingMu.Lock()
+ pendingMsgs[r.StationID] = nil
+ pendingMu.Unlock()
+ go dialUpstreamSidecar(r.StationID, r)
+ } else {
+ go validateUpstream(r.StationID, r)
+ }
+ }
+
+ notifyUpdated()
+}
+
+// validateUpstream test-dials a rule's upstream and records/clears the error so
+// the UI reflects unreachable hosts.
+func validateUpstream(id string, rule ForwarderRule) {
+ upstreamBase := strings.TrimRight(rule.UpstreamURL, "/")
+ upstreamPath := rule.upstreamPath(id)
+
+ tlsConfig := &tls.Config{InsecureSkipVerify: rule.Insecure}
+ if rule.CaCert != "" {
+ caCertPool := x509.NewCertPool()
+ if ok := caCertPool.AppendCertsFromPEM([]byte(rule.CaCert)); !ok {
+ recordForwarderError(id, "invalid CA certificate")
+ notifyUpdated()
+ return
+ }
+ tlsConfig.RootCAs = caCertPool
+ }
+
+ var header http.Header
+ if rule.Username != "" || rule.Password != "" {
+ header = authHeader(rule.Username, rule.Password)
+ }
+
+ ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
+ conn, _, err := websocket.Dial(ctx, upstreamBase+upstreamPath, &websocket.DialOptions{
+ Subprotocols: []string{"ocpp1.6"},
+ HTTPHeader: header,
+ HTTPClient: &http.Client{Transport: &http.Transport{TLSClientConfig: tlsConfig}},
+ })
+ cancel()
+
+ // a charger may have connected meanwhile; its sidecar is authoritative
+ sidecarsMu.Lock()
+ _, active := sidecars[id]
+ sidecarsMu.Unlock()
+ if active {
+ if conn != nil {
+ conn.CloseNow()
+ }
+ return
+ }
+
+ if err != nil {
+ recordForwarderError(id, err.Error())
+ notifyUpdated()
+ return
+ }
+ conn.Close(websocket.StatusNormalClosure, "")
+ if clearForwarderError(id) {
+ notifyUpdated()
+ }
+}
+
+// ForwarderRules returns the current forwarding rules.
+func ForwarderRules() []ForwarderRule {
+ forwarderMu.RLock()
+ defer forwarderMu.RUnlock()
+ return forwarderRules
+}
+
+// StartForwarder is a no-op; hooks fire on every charger connection.
+func StartForwarder() {}
+
+// actionsRelayedToUpstream lists actions for which upstream is the authoritative
+// Central System: evcc's handler is bypassed and upstream's reply relayed back.
+var actionsRelayedToUpstream = map[string]bool{
+ "Authorize": true,
+ "StartTransaction": true,
+ "StopTransaction": true,
+ "DataTransfer": true,
+}
+
+var forwarderLog = util.NewLogger("ocpp-forwarder")
+
+// init wires the forwarder hooks declared in instance.go.
+func init() {
+ chargerConnectHook = onChargerConnect
+ chargerDisconnectHook = onChargerDisconnect
+ chargerMessageHook = onChargerMessage
+}
+
+// sidecar holds an upstream connection for a single charger.
+type sidecar struct {
+ chargerID string
+ upstreamURL string
+ rule ForwarderRule // rule used to dial; detects param changes
+ conn *websocket.Conn
+
+ // message IDs of upstream-initiated Calls; the charger's reply is routed back to upstream
+ pendingUpstreamCallsMu sync.Mutex
+ pendingUpstreamCalls map[string]struct{}
+
+ // message IDs of charger Calls whose evcc handler was bypassed; upstream's reply is relayed to the charger
+ pendingChargerCallsMu sync.Mutex
+ pendingChargerCalls map[string]struct{}
+
+ // when > 0, forward at most one MeterValues per interval to upstream; evcc still sees every frame
+ meterInterval time.Duration
+ lastMeterFwdMu sync.Mutex
+ lastMeterFwd time.Time
+}
+
+var (
+ sidecarsMu sync.Mutex
+ sidecars = make(map[string]*sidecar)
+
+ // connected charger ids so a runtime-added rule can start a sidecar; guarded by sidecarsMu
+ connectedChargers = make(map[string]bool)
+
+ // last upstream failure per charger, surfaced to the UI; guarded by sidecarsMu
+ forwarderErrors = make(map[string]string)
+
+ // raw frames buffered per charger while the sidecar dials; flushed in order
+ // on connect so BootNotification reaches upstream
+ pendingMu sync.Mutex
+ pendingMsgs = make(map[string][][]byte)
+)
+
+// resolveRule returns the forwarding rule for chargerID, or the "*" fallback.
+func resolveRule(chargerID string) (ForwarderRule, bool) {
+ forwarderMu.RLock()
+ rules := forwarderRules
+ forwarderMu.RUnlock()
+ var fallback ForwarderRule
+ var hasFallback bool
+ for _, r := range rules {
+ if r.StationID == chargerID {
+ return r, true
+ }
+ if r.StationID == "*" {
+ fallback = r
+ hasFallback = true
+ }
+ }
+ return fallback, hasFallback
+}
+
+// onChargerConnect dials a sidecar for the connecting charger when a rule matches.
+func onChargerConnect(ch ws.Channel) {
+ id := ch.ID()
+
+ sidecarsMu.Lock()
+ connectedChargers[id] = true
+ sidecarsMu.Unlock()
+
+ rule, ok := resolveRule(id)
+ if !ok {
+ return
+ }
+ // open the buffer before dialling so early messages are captured, not dropped
+ pendingMu.Lock()
+ pendingMsgs[id] = nil
+ pendingMu.Unlock()
+
+ go dialUpstreamSidecar(id, rule)
+}
+
+func dialUpstreamSidecar(id string, rule ForwarderRule) {
+ upstreamBase := strings.TrimRight(rule.UpstreamURL, "/")
+ upstreamPath := rule.upstreamPath(id)
+
+ var header http.Header
+ if rule.Username != "" || rule.Password != "" {
+ header = authHeader(rule.Username, rule.Password)
+ }
+
+ tlsConfig := &tls.Config{InsecureSkipVerify: rule.Insecure}
+ if rule.CaCert != "" {
+ caCertPool := x509.NewCertPool()
+ if ok := caCertPool.AppendCertsFromPEM([]byte(rule.CaCert)); !ok {
+ forwarderLog.WARN.Printf("forwarder: failed to parse CA cert for %s; forwarding disabled", id)
+ recordForwarderError(id, "invalid CA certificate")
+ notifyUpdated()
+ drainPendingWithErrors(id, nil)
+ return
+ }
+ tlsConfig.RootCAs = caCertPool
+ }
+
+ ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
+ conn, _, err := websocket.Dial(ctx, upstreamBase+upstreamPath, &websocket.DialOptions{
+ Subprotocols: []string{"ocpp1.6"},
+ HTTPHeader: header,
+ HTTPClient: &http.Client{Transport: &http.Transport{TLSClientConfig: tlsConfig}},
+ })
+ cancel()
+ if err != nil {
+ forwarderLog.WARN.Printf("forwarder: dial upstream for %s: %v; forwarding disabled", id, err)
+ recordForwarderError(id, err.Error())
+ notifyUpdated()
+ drainPendingWithErrors(id, nil)
+ return
+ }
+ conn.SetReadLimit(-1) // no limit; OCPP frames can be large
+
+ sc := &sidecar{
+ chargerID: id,
+ upstreamURL: upstreamBase,
+ rule: rule,
+ conn: conn,
+ pendingUpstreamCalls: make(map[string]struct{}),
+ pendingChargerCalls: make(map[string]struct{}),
+ }
+
+ // install sidecar and drain the pending buffer
+ sidecarsMu.Lock()
+ pendingMu.Lock()
+ sidecars[id] = sc
+ buffered := pendingMsgs[id]
+ delete(pendingMsgs, id)
+ pendingMu.Unlock()
+ sidecarsMu.Unlock()
+
+ // flush buffered Calls; register relay actions so upstream's reply routes back.
+ // CallResults/Errors are skipped: evcc already answered them.
+ flushed := 0
+ for _, frame := range buffered {
+ msgType, msgID, action, err := parseOCPPFrame(frame)
+ if err != nil || msgType != ocppj.CALL {
+ continue
+ }
+ if actionsRelayedToUpstream[action] {
+ sc.pendingChargerCallsMu.Lock()
+ sc.pendingChargerCalls[msgID] = struct{}{}
+ sc.pendingChargerCallsMu.Unlock()
+ }
+ if err := conn.Write(context.Background(), websocket.MessageText, frame); err != nil {
+ forwarderLog.ERROR.Printf("forwarder: write buffered frame to upstream for %s: %v", id, err)
+ conn.CloseNow()
+ recordForwarderError(id, err.Error())
+ notifyUpdated()
+ return
+ }
+ flushed++
+ }
+ if flushed > 0 {
+ forwarderLog.DEBUG.Printf("forwarder: flushed %d call(s) to upstream for %s (skipped %d non-call frame(s))",
+ flushed, id, len(buffered)-flushed)
+ }
+
+ clearForwarderError(id)
+ notifyUpdated()
+ forwarderLog.INFO.Printf("forwarder: %s → %s", id, upstreamBase+upstreamPath)
+
+ sc.readFromUpstream(rule.ReadOnly)
+}
+
+// drainPendingWithErrors discards charger id's pending buffer, sending a CallError
+// to the charger for each buffered relay Call so it is not left hanging. sc may be nil.
+func drainPendingWithErrors(id string, sc *sidecar) {
+ pendingMu.Lock()
+ buffered := pendingMsgs[id]
+ delete(pendingMsgs, id)
+ pendingMu.Unlock()
+
+ cs := Instance()
+
+ for _, frame := range buffered {
+ msgType, msgID, action, err := parseOCPPFrame(frame)
+ if err != nil || msgType != ocppj.CALL || !actionsRelayedToUpstream[action] {
+ continue
+ }
+ errFrame, _ := (&ocppj.CallError{
+ MessageTypeId: ocppj.CALL_ERROR,
+ UniqueId: msgID,
+ ErrorCode: ocppj.GenericError,
+ ErrorDescription: "Upstream OCPP server unavailable",
+ }).MarshalJSON()
+ if writeErr := cs.Write(id, errFrame); writeErr != nil {
+ forwarderLog.WARN.Printf("forwarder: send error for pending call %s to %s: %v", msgID, id, writeErr)
+ }
+ }
+
+ // error out tracked pendingChargerCalls (sidecar dropped mid-session)
+ if sc != nil {
+ sc.pendingChargerCallsMu.Lock()
+ for msgID := range sc.pendingChargerCalls {
+ errFrame, _ := (&ocppj.CallError{
+ MessageTypeId: ocppj.CALL_ERROR,
+ UniqueId: msgID,
+ ErrorCode: ocppj.GenericError,
+ ErrorDescription: "Upstream OCPP server disconnected",
+ }).MarshalJSON()
+ if writeErr := cs.Write(sc.chargerID, errFrame); writeErr != nil {
+ forwarderLog.WARN.Printf("forwarder: send disconnect error for %s to %s: %v", msgID, sc.chargerID, writeErr)
+ }
+ }
+ sc.pendingChargerCalls = make(map[string]struct{})
+ sc.pendingChargerCallsMu.Unlock()
+ }
+}
+
+// onChargerDisconnect closes the charger's sidecar connection.
+func onChargerDisconnect(ch ws.Channel) {
+ id := ch.ID()
+
+ // discard pending buffer
+ pendingMu.Lock()
+ delete(pendingMsgs, id)
+ pendingMu.Unlock()
+
+ sidecarsMu.Lock()
+ delete(connectedChargers, id)
+ sc, ok := sidecars[id]
+ if ok {
+ delete(sidecars, id)
+ }
+ _, hadErr := forwarderErrors[id]
+ delete(forwarderErrors, id)
+ sidecarsMu.Unlock()
+ if ok {
+ sc.conn.CloseNow()
+ forwarderLog.DEBUG.Printf("forwarder: %s upstream connection closed", id)
+ }
+ if ok || hadErr {
+ notifyUpdated()
+ }
+}
+
+// onChargerMessage handles a raw OCPP frame from a charger and reports whether
+// evcc's handler should be bypassed (upstream is the authoritative responder).
+// Relay-action Calls are forwarded and bypassed; other Calls are forwarded and
+// also handled by evcc; CallResults/Errors are forwarded only when they answer
+// an upstream-initiated Call.
+func onChargerMessage(ch ws.Channel, data []byte) bool {
+ id := ch.ID()
+
+ msgType, msgID, action, err := parseOCPPFrame(data)
+ if err != nil {
+ return false
+ }
+
+ sidecarsMu.Lock()
+ sc := sidecars[id]
+ sidecarsMu.Unlock()
+
+ switch msgType {
+ case ocppj.CALL:
+ relay := actionsRelayedToUpstream[action]
+
+ if sc != nil {
+ // throttle MeterValues to upstream; evcc still processes every frame
+ if action == "MeterValues" && sc.meterInterval > 0 {
+ sc.lastMeterFwdMu.Lock()
+ elapsed := time.Since(sc.lastMeterFwd)
+ if elapsed < sc.meterInterval {
+ sc.lastMeterFwdMu.Unlock()
+ return false // evcc handles normally; skip upstream
+ }
+ sc.lastMeterFwd = time.Now()
+ sc.lastMeterFwdMu.Unlock()
+ }
+ if relay {
+ sc.pendingChargerCallsMu.Lock()
+ sc.pendingChargerCalls[msgID] = struct{}{}
+ sc.pendingChargerCallsMu.Unlock()
+ }
+ if err := sc.conn.Write(context.Background(), websocket.MessageText, data); err != nil {
+ forwarderLog.ERROR.Printf("forwarder: write to upstream for %s: %v", id, err)
+ }
+ return relay
+ }
+
+ // sidecar not ready: buffer if a pending slot exists
+ pendingMu.Lock()
+ _, hasPending := pendingMsgs[id]
+ if hasPending {
+ pendingMsgs[id] = append(pendingMsgs[id], slices.Clone(data))
+ }
+ pendingMu.Unlock()
+
+ // bypass evcc for relay actions while buffering; upstream answers after flush
+ return relay && hasPending
+
+ case ocppj.CALL_RESULT, ocppj.CALL_ERROR:
+ // forward only if this answers an upstream-initiated Call
+ if sc == nil {
+ return false
+ }
+ sc.pendingUpstreamCallsMu.Lock()
+ _, isUpstream := sc.pendingUpstreamCalls[msgID]
+ if isUpstream {
+ delete(sc.pendingUpstreamCalls, msgID)
+ }
+ sc.pendingUpstreamCallsMu.Unlock()
+ if !isUpstream {
+ return false
+ }
+ if err := sc.conn.Write(context.Background(), websocket.MessageText, data); err != nil {
+ forwarderLog.ERROR.Printf("forwarder: write upstream reply for %s: %v", id, err)
+ }
+ return false // evcc may also see it; harmless for unknown IDs
+ }
+ return false
+}
+
+// readFromUpstream relays frames from upstream: Calls are injected into the
+// charger (reply routed back), responses to bypassed charger Calls are relayed
+// to the charger, others discarded. Calls are rejected in read-only mode.
+func (sc *sidecar) readFromUpstream(readOnly bool) {
+ defer func() {
+ sidecarsMu.Lock()
+ if sidecars[sc.chargerID] == sc {
+ delete(sidecars, sc.chargerID)
+ }
+ sidecarsMu.Unlock()
+ // error out pending relayed calls; upstream is gone
+ drainPendingWithErrors(sc.chargerID, sc)
+ sc.conn.CloseNow()
+ notifyUpdated()
+ }()
+
+ for {
+ _, msg, err := sc.conn.Read(context.Background())
+ if err != nil {
+ forwarderLog.DEBUG.Printf("forwarder: upstream disconnected for %s: %v", sc.chargerID, err)
+ // charger disconnect deletes the sidecar first, so a still-current sidecar means upstream dropped
+ sidecarsMu.Lock()
+ upstreamDrop := sidecars[sc.chargerID] == sc
+ sidecarsMu.Unlock()
+ if upstreamDrop {
+ recordForwarderError(sc.chargerID, err.Error())
+ }
+ return
+ }
+
+ msgType, msgID, action, err := parseOCPPFrame(msg)
+ if err != nil {
+ forwarderLog.WARN.Printf("forwarder: upstream parse error for %s: %v", sc.chargerID, err)
+ continue
+ }
+
+ switch msgType {
+ case ocppj.CALL:
+ if readOnly {
+ forwarderLog.DEBUG.Printf("forwarder: blocking upstream call %s in read-only session %s", msgID, sc.chargerID)
+ errFrame, _ := (&ocppj.CallError{
+ MessageTypeId: ocppj.CALL_ERROR,
+ UniqueId: msgID,
+ ErrorCode: ocppj.SecurityError,
+ ErrorDescription: "Charger control not allowed: forwarder is in read-only mode",
+ }).MarshalJSON()
+ _ = sc.conn.Write(context.Background(), websocket.MessageText, errFrame)
+ continue
+ }
+
+ // absorb upstream's MeterValueSampleInterval as our throttle; reply
+ // Accepted without touching the charger's config
+ if action == "ChangeConfiguration" {
+ if interval, ok := extractMeterValueSampleInterval(msg); ok {
+ sc.lastMeterFwdMu.Lock()
+ sc.meterInterval = interval
+ sc.lastMeterFwdMu.Unlock()
+ forwarderLog.DEBUG.Printf("forwarder: upstream set MeterValueSampleInterval=%v for %s", interval, sc.chargerID)
+ accepted, _ := (&ocppj.CallResult{
+ MessageTypeId: ocppj.CALL_RESULT,
+ UniqueId: msgID,
+ Payload: core.NewChangeConfigurationConfirmation(core.ConfigurationStatusAccepted),
+ }).MarshalJSON()
+ _ = sc.conn.Write(context.Background(), websocket.MessageText, accepted)
+ continue
+ }
+ }
+
+ // track id so the charger's reply is forwarded back to upstream
+ sc.pendingUpstreamCallsMu.Lock()
+ sc.pendingUpstreamCalls[msgID] = struct{}{}
+ sc.pendingUpstreamCallsMu.Unlock()
+
+ if err := Instance().Write(sc.chargerID, msg); err != nil {
+ forwarderLog.ERROR.Printf("forwarder: inject upstream call into charger %s: %v", sc.chargerID, err)
+ }
+
+ case ocppj.CALL_RESULT, ocppj.CALL_ERROR:
+ // upstream's authoritative response to a bypassed charger Call?
+ sc.pendingChargerCallsMu.Lock()
+ _, isChargerCall := sc.pendingChargerCalls[msgID]
+ if isChargerCall {
+ delete(sc.pendingChargerCalls, msgID)
+ }
+ sc.pendingChargerCallsMu.Unlock()
+
+ if isChargerCall {
+ // relay to charger; its handler was bypassed and it awaits this reply
+ if err := Instance().Write(sc.chargerID, msg); err != nil {
+ forwarderLog.ERROR.Printf("forwarder: relay upstream response to charger %s: %v", sc.chargerID, err)
+ }
+ continue
+ }
+ // discard: evcc already replied for non-relay Calls
+ }
+ }
+}
+
+// ForwarderSessionStatus is the observable state of one forwarder session.
+type ForwarderSessionStatus struct {
+ ChargerID string `json:"chargerId"`
+ UpstreamURL string `json:"upstreamUrl"`
+ UpstreamConnected bool `json:"upstreamConnected"`
+ Error string `json:"error,omitempty"`
+}
+
+// recordForwarderError stores a charger's last upstream failure; caller must notifyUpdated.
+func recordForwarderError(id, msg string) {
+ sidecarsMu.Lock()
+ forwarderErrors[id] = msg
+ sidecarsMu.Unlock()
+}
+
+// clearForwarderError drops a charger's stored error, reporting whether one existed; caller must notifyUpdated.
+func clearForwarderError(id string) bool {
+ sidecarsMu.Lock()
+ _, ok := forwarderErrors[id]
+ delete(forwarderErrors, id)
+ sidecarsMu.Unlock()
+ return ok
+}
+
+var (
+ forwarderCbMu sync.Mutex
+ forwarderUpdatedCb func()
+)
+
+func notifyUpdated() {
+ status := GetForwarderStatus()
+ forwarderLog.DEBUG.Printf("forwarder: notifyUpdated sessions=%d", len(status))
+ forwarderCbMu.Lock()
+ cb := forwarderUpdatedCb
+ forwarderCbMu.Unlock()
+ if cb != nil {
+ cb()
+ }
+}
+
+// SetForwarderUpdated registers a callback fired when a session connects or disconnects.
+func SetForwarderUpdated(cb func()) {
+ forwarderCbMu.Lock()
+ forwarderUpdatedCb = cb
+ forwarderCbMu.Unlock()
+}
+
+// GetForwarderStatus returns a snapshot of all active forwarder sessions.
+func GetForwarderStatus() []ForwarderSessionStatus {
+ forwarderMu.RLock()
+ rules := append([]ForwarderRule(nil), forwarderRules...)
+ forwarderMu.RUnlock()
+
+ sidecarsMu.Lock()
+ defer sidecarsMu.Unlock()
+ out := make([]ForwarderSessionStatus, 0, len(rules))
+ for _, r := range rules {
+ if r.StationID == "*" {
+ continue
+ }
+ st := ForwarderSessionStatus{
+ ChargerID: r.StationID,
+ UpstreamURL: strings.TrimRight(r.UpstreamURL, "/"),
+ }
+ if _, ok := sidecars[r.StationID]; ok {
+ st.UpstreamConnected = true
+ } else if msg, ok := forwarderErrors[r.StationID]; ok {
+ st.Error = msg
+ }
+ out = append(out, st)
+ }
+ return out
+}
+
+// sameConnection reports whether two rules dial upstream identically (no reconnect
+// needed). ReadOnly is excluded; it applies live per message.
+func (r ForwarderRule) sameConnection(o ForwarderRule) bool {
+ return r.UpstreamURL == o.UpstreamURL &&
+ r.UpstreamStationID == o.UpstreamStationID &&
+ r.Username == o.Username &&
+ r.Password == o.Password &&
+ r.Insecure == o.Insecure &&
+ r.CaCert == o.CaCert
+}
+
+// upstreamPath returns the upstream WebSocket path, defaulting to the charger's own ID.
+func (r ForwarderRule) upstreamPath(chargerID string) string {
+ sid := r.UpstreamStationID
+ if sid == "" {
+ sid = chargerID
+ }
+ return "/" + strings.TrimLeft(sid, "/")
+}
+
+// authHeader returns a Basic Auth header for the given credentials.
+func authHeader(username, password string) http.Header {
+ creds := base64.StdEncoding.EncodeToString([]byte(username + ":" + password))
+ h := make(http.Header)
+ h.Set("Authorization", "Basic "+creds)
+ return h
+}
+
+// extractMeterValueSampleInterval returns the interval from a ChangeConfiguration
+// Call for key MeterValueSampleInterval.
+func extractMeterValueSampleInterval(msg []byte) (time.Duration, bool) {
+ var frame []json.RawMessage
+ if err := json.Unmarshal(msg, &frame); err != nil || len(frame) < 4 {
+ return 0, false
+ }
+ var req core.ChangeConfigurationRequest
+ if err := json.Unmarshal(frame[3], &req); err != nil {
+ return 0, false
+ }
+ if req.Key != "MeterValueSampleInterval" {
+ return 0, false
+ }
+ secs, err := strconv.Atoi(req.Value)
+ if err != nil || secs <= 0 {
+ return 0, false
+ }
+ return time.Duration(secs) * time.Second, true
+}
+
+// parseOCPPFrame extracts the message type, id and (for Calls) action from a raw frame.
+func parseOCPPFrame(msg []byte) (msgType ocppj.MessageType, msgID string, action string, err error) {
+ var frame []json.RawMessage
+ if err = json.Unmarshal(msg, &frame); err != nil || len(frame) < 2 {
+ return 0, "", "", fmt.Errorf("invalid OCPP frame")
+ }
+ if err = json.Unmarshal(frame[0], &msgType); err != nil {
+ return 0, "", "", fmt.Errorf("invalid message type: %w", err)
+ }
+ if err = json.Unmarshal(frame[1], &msgID); err != nil {
+ return 0, "", "", fmt.Errorf("invalid message id: %w", err)
+ }
+ if msgType == ocppj.CALL && len(frame) >= 3 {
+ _ = json.Unmarshal(frame[2], &action) // best-effort; empty string if missing
+ }
+ return msgType, msgID, action, nil
+}
diff --git a/charger/ocpp/instance.go b/charger/ocpp/instance.go
index a3a61f5ff..dfb3aadeb 100644
--- a/charger/ocpp/instance.go
+++ b/charger/ocpp/instance.go
@@ -24,6 +24,24 @@ type Config struct {
Port int `json:"port"`
}
+// ForwarderRule maps a station ID (or "*" for all chargers) to an upstream OCPP server URL.
+type ForwarderRule struct {
+ StationID string `json:"stationId" yaml:"stationId"`
+ UpstreamURL string `json:"upstreamUrl" yaml:"upstreamUrl"`
+ Password string `json:"password,omitempty" yaml:"password,omitempty"`
+ UpstreamStationID string `json:"upstreamStationId,omitempty" yaml:"upstreamStationId,omitempty"`
+ Username string `json:"username,omitempty" yaml:"username,omitempty"`
+ Insecure bool `json:"insecure,omitempty" yaml:"insecure,omitempty"`
+ CaCert string `json:"caCert,omitempty" yaml:"caCert,omitempty"`
+ ReadOnly bool `json:"readOnly,omitempty" yaml:"readOnly,omitempty"`
+}
+
+func (r ForwarderRule) Redacted() ForwarderRule {
+ r.Password = util.Masked(r.Password)
+ r.CaCert = util.Masked(r.CaCert)
+ return r
+}
+
var (
once sync.Once
instance *CS
@@ -32,6 +50,47 @@ var (
externalUrl string
)
+// Forwarder hooks, nil unless the forwarder is built in (set once in init()
+// before any charger connects, so reads need no lock).
+var (
+ chargerConnectHook func(ws.Channel)
+ chargerDisconnectHook func(ws.Channel)
+ chargerMessageHook func(ws.Channel, []byte) bool
+)
+
+// interceptingServer routes connect/disconnect/message events through the
+// forwarder hooks. The message hook returns true to bypass evcc's OCPP handler.
+type interceptingServer struct {
+ ws.Server
+}
+
+func (s *interceptingServer) SetMessageHandler(handler ws.MessageHandler) {
+ s.Server.SetMessageHandler(func(ch ws.Channel, data []byte) error {
+ if chargerMessageHook != nil && chargerMessageHook(ch, data) {
+ return nil
+ }
+ return handler(ch, data)
+ })
+}
+
+func (s *interceptingServer) SetNewClientHandler(handler ws.ConnectedHandler) {
+ s.Server.SetNewClientHandler(func(ch ws.Channel) {
+ if chargerConnectHook != nil {
+ chargerConnectHook(ch)
+ }
+ handler(ch)
+ })
+}
+
+func (s *interceptingServer) SetDisconnectedClientHandler(handler func(ws.Channel)) {
+ s.Server.SetDisconnectedClientHandler(func(ch ws.Channel) {
+ if chargerDisconnectHook != nil {
+ chargerDisconnectHook(ch)
+ }
+ handler(ch)
+ })
+}
+
// Port returns the TCP port the central system is bound to. With the default
// configuration this equals the configured port; when port 0 is configured
// (as in tests) it is the OS-assigned ephemeral port. It returns 0 while the
@@ -65,6 +124,11 @@ func ExternalUrl() string {
return u.String()
}
+// CurrentConfig returns the current runtime OCPP configuration.
+func CurrentConfig() Config {
+ return Config{Port: port}
+}
+
// Init initializes the OCPP server
func Init(cfg Config, networkExternalUrl string) {
port = cfg.Port
@@ -75,7 +139,7 @@ func Instance() *CS {
once.Do(func() {
log := util.NewLogger("ocpp")
- server := ws.NewServer()
+ server := &interceptingServer{Server: ws.NewServer()}
server.SetCheckOriginHandler(func(r *http.Request) bool { return true })
dispatcher := ocppj.NewDefaultServerDispatcher(ocppj.NewFIFOQueueMap(0))
@@ -93,6 +157,7 @@ func Instance() *CS {
log: log,
regs: make(map[string]*registration),
CentralSystem: cs,
+ server: server,
}
instance.txnId.Store(time.Now().UTC().Unix())
diff --git a/cmd/root.go b/cmd/root.go
index 8960eabe4..780e08216 100644
--- a/cmd/root.go
+++ b/cmd/root.go
@@ -32,6 +32,7 @@ import (
"github.com/evcc-io/evcc/util/telemetry"
_ "github.com/joho/godotenv/autoload"
"github.com/prometheus/client_golang/prometheus/promhttp"
+ "github.com/samber/lo"
"github.com/spf13/cast"
"github.com/spf13/cobra"
vpr "github.com/spf13/viper"
@@ -204,7 +205,7 @@ func runRoot(cmd *cobra.Command, args []string) {
ocppCS.SetUpdated(func() {
// republish when OCPP state updates
valueChan <- util.Param{Key: keys.Ocpp, Val: globalconfig.ConfigStatus{
- Config: conf.Ocpp,
+ Config: ocpp.CurrentConfig(),
Status: ocpp.GetStatus(),
}}
})
@@ -212,6 +213,16 @@ func runRoot(cmd *cobra.Command, args []string) {
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()))
+ }
// value cache
cache := util.NewParamCache()
@@ -376,9 +387,13 @@ func runRoot(cmd *cobra.Command, args []string) {
valueChan <- util.Param{Key: keys.Mqtt, Val: conf.Mqtt}
valueChan <- util.Param{Key: keys.Network, Val: conf.Network}
valueChan <- util.Param{Key: keys.Ocpp, Val: globalconfig.ConfigStatus{
- Config: conf.Ocpp,
+ Config: ocpp.CurrentConfig(),
Status: ocpp.GetStatus(),
}}
+ 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(),
+ }}
valueChan <- util.Param{Key: keys.Sponsor, Val: globalconfig.ConfigStatus{
Status: sponsor.RedactedStatus(),
YamlSource: yamlSource.sponsor,
diff --git a/cmd/setup.go b/cmd/setup.go
index c9601b1f9..1cf9d1fe0 100644
--- a/cmd/setup.go
+++ b/cmd/setup.go
@@ -842,7 +842,18 @@ func configureMDNS(conf globalconfig.Network) error {
// setup OCPP
func configureOCPP(cfg *ocpp.Config, externalUrl string) {
+ if settings.Exists(keys.Ocpp) {
+ if err := settings.Json(keys.Ocpp, cfg); err != nil {
+ log.WARN.Printf("ocpp: failed to load settings: %v", err)
+ }
+ }
ocpp.Init(*cfg, externalUrl)
+
+ // Load proxy forwarding rules from DB if present.
+ var rules []ocpp.ForwarderRule
+ if err := settings.Json(keys.OcppForwarder, &rules); err == nil {
+ ocpp.ApplyForwarderRules(rules)
+ }
}
// setup EEBus
diff --git a/core/keys/global.go b/core/keys/global.go
index 695a92bb3..232ec18b4 100644
--- a/core/keys/global.go
+++ b/core/keys/global.go
@@ -16,6 +16,7 @@ const (
MessagingEvents = "messagingEvents"
ModbusProxy = "modbusproxy"
Ocpp = "ocpp"
+ OcppForwarder = "ocppforwarder"
Tariffs = "tariffs"
TariffRefs = "tariffRefs"
Version = "version"
diff --git a/docs/agents/ocpp-forwarder.md b/docs/agents/ocpp-forwarder.md
new file mode 100644
index 000000000..173bf356a
--- /dev/null
+++ b/docs/agents/ocpp-forwarder.md
@@ -0,0 +1,44 @@
+# OCPP Forwarder Architecture
+
+The OCPP forwarder (`charger/ocpp/forwarder.go`) is a hybrid proxy that lets a charger talk to evcc and an upstream OCPP server at the same time. Chargers connect directly to evcc's central system on the normal port. For each charger with a matching `ForwarderRule`, a "sidecar" WebSocket connection to the upstream server is opened and kept in parallel for the lifetime of the charger connection.
+
+The forwarder is opt-in: hooks (`chargerConnectHook`, `chargerDisconnectHook`, `chargerMessageHook` in `instance.go`) are nil unless a rule matches, so a charger without a rule behaves exactly as before.
+
+## Forwarding modes
+
+Two modes apply at the same time, selected per message by its action.
+
+### Transparent relay (billing-critical)
+
+For the actions in `actionsRelayedToUpstream` (`Authorize`, `StartTransaction`, `StopTransaction`, `DataTransfer`), upstream is the authoritative Central System:
+
+1. Charger sends the Call to evcc.
+2. The message hook forwards it to the upstream sidecar and bypasses evcc's OCPP handler.
+3. Upstream's `CallResult`/`CallError` is relayed back to the charger.
+
+evcc's handler is never invoked for these. This lets the pay backend control authorization, issue its own transaction IDs, and see consistent Start/Stop pairs.
+
+### Sidecar observation (informational)
+
+For all other messages (`BootNotification`, `StatusNotification`, `MeterValues`, `Heartbeat`, etc.):
+
+1. Charger sends the Call to evcc, which processes it normally.
+2. The same frame is also mirrored to the upstream sidecar.
+
+Upstream observes the session while evcc manages the charger as usual.
+
+## Upstream to charger (commands)
+
+Calls (type 2) initiated by upstream are injected into the charger via `CS.Write`. The charger's `CallResult`/`CallError` is routed back to upstream. Examples: `RemoteStartTransaction`, `RemoteStopTransaction`, `GetConfiguration`, `ChangeConfiguration`, `TriggerMessage`, `SetChargingProfile`.
+
+`ChangeConfiguration` for `MeterValueSampleInterval` is intercepted: the forwarder absorbs it as a local throttle on `MeterValues` forwarded to upstream and replies `Accepted` without touching the charger's own config. evcc still processes every `MeterValues` frame for energy management.
+
+## Read-only mode
+
+When a rule sets `ReadOnly`, upstream may observe but cannot control the charger. Any incoming Call from upstream is answered with a `SecurityError` and not forwarded. `ReadOnly` is applied live per message, so toggling it does not require reconnecting the sidecar.
+
+## Connection lifecycle
+
+Frames that arrive from a charger before its sidecar finishes dialling are buffered (`pendingMsgs`) and flushed in order once the sidecar connects, so early messages such as `BootNotification` still reach upstream. If the dial fails or upstream drops mid-session, any buffered or in-flight relay Calls are answered to the charger with a `CallError` so it is not left hanging, and the failure is surfaced to the UI via `forwarderErrors`.
+
+Rules can be changed at runtime through `ApplyForwarderRules`. Sidecars for removed rules are closed, rules with changed connection parameters are re-dialled, and rules for chargers that are not connected are test-dialled to surface unreachable hosts immediately.
diff --git a/i18n/de.json b/i18n/de.json
index d07d8fb52..427c973ac 100644
--- a/i18n/de.json
+++ b/i18n/de.json
@@ -618,13 +618,12 @@
"warningUrlPath": "Die URL benötigt normalerweise keinen Pfad. Bist du dir sicher?"
},
"ocpp": {
- "connectedChargers": "Verbundene Wallboxen",
- "connectionStatus": "Konfigurierte Station-IDs",
- "connectionStatusHelp": "Verbindungsstatus der konfigurierten Wallboxen.",
- "detectedChargers": "Erkannte Station-IDs",
- "detectedHelp": "Diese Wallboxen haben versucht, sich mit evcc zu verbinden. Um eine Wallbox zu verwenden, erstelle einen Ladepunkt mit ihrer Station-ID.",
+ "forwardingConfigured": "Weiterleitung konfiguriert",
+ "forwardingError": "Weiterleitungsfehler",
+ "forwardingOff": "Keine Weiterleitung konfiguriert",
"noChargers": "Keine OCPP-Wallboxen erkannt.",
- "noStations": "Keine Stationen verbunden",
+ "stations": "Station-IDs",
+ "stationsHelp": "Konfigurierte und erkannte OCPP-Wallboxen. Zusätzlich zur lokalen Steuerung kannst du eine Nachrichtenweiterleitung an einen externen Server (z. B. für die Abrechnung) konfigurieren.",
"status": {
"configured": "Nicht verbunden",
"connected": "Verbunden",
@@ -634,6 +633,30 @@
"url": "Server-URL",
"urlHelp": "Kopiere diese URL in die Konfiguration deiner Wallbox. Details findest du im Handbuch des Herstellers. Die Wallbox sollte automatisch ihre eindeutige Kennung (Station-ID) an die URL anhängen. In seltenen Fällen musst du die Kennung manuell angeben. Beispiel: `{url}`"
},
+ "ocppforwarder": {
+ "description": "Leite die OCPP-Nachrichten dieser Wallbox zusätzlich zu evcc an einen externen Server weiter (Abrechnungsplattform, Netzbetreiber usw.).",
+ "editTitle": "OCPP-Weiterleitung",
+ "labelCaCert": "Serverzertifikat (CA)",
+ "labelCheckInsecure": "Selbstsignierte Zertifikate erlauben",
+ "labelInsecure": "Zertifikatsprüfung",
+ "password": "Passwort",
+ "passwordHelp": "Für HTTP-Basic-Auth am Upstream-Server.",
+ "readOnly": {
+ "check": "Befehle vom Upstream-Server blockieren",
+ "help": {
+ "false": "Der Upstream kann über evcc Befehle an die Wallbox senden.",
+ "true": "Der Upstream empfängt alle Wallbox-Nachrichten, kann aber keine Befehle senden. evcc behält die alleinige Kontrolle."
+ },
+ "label": "Upstream-Befehle"
+ },
+ "status": "Status",
+ "upstreamStationId": "Station-ID",
+ "upstreamStationIdHelp": "Wallbox-Kennung, die an den Upstream-Server gesendet wird.",
+ "upstreamUrl": "Upstream-Server-URL",
+ "upstreamUrlHelp": "OCPP-WebSocket-URL des externen Servers.",
+ "username": "Benutzername",
+ "usernameHelp": "Für HTTP-Basic-Auth am Upstream-Server."
+ },
"optimizer": {
"description": "Analysiert Solarprognose, Strompreise und deinen typischen Verbrauch, um Batterie- und Ladestrategie zu optimieren. Daten werden zur Berechnung an den evcc Optimierungsdienst übertragen. Berechnet und visualisiert aktuell nur. Steuert noch keine Geräte.",
"enable": "Optimizer aktivieren",
diff --git a/i18n/en.json b/i18n/en.json
index 3f70071b1..405cd2120 100644
--- a/i18n/en.json
+++ b/i18n/en.json
@@ -617,13 +617,12 @@
"warningUrlPath": "The URL usually doesn't need a path. Are you sure this is correct?"
},
"ocpp": {
- "connectedChargers": "Connected chargers",
- "connectionStatus": "Configured station IDs",
- "connectionStatusHelp": "Connection status of configured chargers.",
- "detectedChargers": "Detected station IDs",
- "detectedHelp": "These chargers have tried to connect to evcc. To use a charger, create a loadpoint with its station ID.",
+ "forwardingConfigured": "Forwarding configured",
+ "forwardingError": "Forwarding error",
+ "forwardingOff": "No forwarding configured",
"noChargers": "No OCPP chargers detected.",
- "noStations": "No stations connected",
+ "stations": "Station IDs",
+ "stationsHelp": "Configured and detected OCPP chargers. In addition to local control, you can forward messages to an external server (e.g. for billing).",
"status": {
"configured": "Not connected",
"connected": "Connected",
@@ -633,6 +632,30 @@
"url": "Server URL",
"urlHelp": "Copy this URL into your charger's configuration. Check the manufacturer's manual for details. The charger is expected to automatically append its unique identifier (station ID) to the url. In rare cases, you may need to manually specify the identifier. Example: `{url}`"
},
+ "ocppforwarder": {
+ "description": "Forward this charger's OCPP messages to an external server (billing platform, DSO, etc.) in addition to evcc.",
+ "editTitle": "OCPP Forward",
+ "labelCaCert": "Server certificate (CA)",
+ "labelCheckInsecure": "Allow self-signed certificates",
+ "labelInsecure": "Certificate validation",
+ "password": "Password",
+ "passwordHelp": "For HTTP Basic Auth to the upstream server.",
+ "readOnly": {
+ "check": "Block commands from the upstream server",
+ "help": {
+ "false": "Upstream can send commands to the charger via evcc.",
+ "true": "Upstream receives all charger messages but cannot send commands. evcc retains exclusive control."
+ },
+ "label": "Upstream commands"
+ },
+ "status": "Status",
+ "upstreamStationId": "Station ID",
+ "upstreamStationIdHelp": "Charger identifier sent to the upstream server.",
+ "upstreamUrl": "Upstream server URL",
+ "upstreamUrlHelp": "OCPP WebSocket URL of the external server.",
+ "username": "Username",
+ "usernameHelp": "For HTTP Basic Auth to the upstream server."
+ },
"optimizer": {
"description": "Analyzes solar forecast, electricity prices, and your consumption patterns to optimize battery and charging strategy. Data is sent to the evcc optimization service for processing. Currently only calculates and visualizes. Does not control devices yet.",
"enable": "Enable Optimizer",
diff --git a/i18n/fr.json b/i18n/fr.json
index b9f934f8a..7dbc2947c 100644
--- a/i18n/fr.json
+++ b/i18n/fr.json
@@ -622,6 +622,44 @@
"url": "URL du Serveur",
"urlHelp": "Copiez cette URL dans la configuration de votre chargeur. Consultez le manuel du fabricant pour plus de détails. Le chargeur devrait ajouter automatiquement son identifiant unique (ID de station) à l'URL. Dans de rares cas, vous devrez peut-être spécifier manuellement l'identifiant. Exemple : `{url}`"
},
+ "ocppforwarder": {
+ "add": "Ajouter une règle",
+ "charger": "Chargeur",
+ "connectionStatus": "État de connexion",
+ "description": "Transmet les appels OCPP à un serveur externe (plateforme de facturation, DSO, etc.) en plus d'evcc. Utilisez \"*\" comme identifiant de station pour correspondre à tous les chargeurs.",
+ "labelCaCert": "Certificat serveur (CA)",
+ "labelCheckInsecure": "Autoriser les certificats auto-signés",
+ "labelInsecure": "Validation du certificat",
+ "noConnections": "Aucune connexion de chargeur active.",
+ "option": {
+ "false": "non",
+ "true": "oui"
+ },
+ "password": "Mot de passe en amont",
+ "passwordHelp": "Mot de passe optionnel pour l'authentification HTTP Basic vers le serveur en amont. L'identifiant de station en amont est utilisé comme nom d'utilisateur.",
+ "readOnly": {
+ "help": {
+ "false": "Le serveur en amont peut envoyer des commandes au chargeur via evcc.",
+ "true": "Le serveur en amont reçoit tous les messages du chargeur mais ne peut pas envoyer de commandes. evcc conserve le contrôle exclusif."
+ },
+ "label": "Lecture seule"
+ },
+ "rule": "Règle {number}",
+ "stationId": "Identifiant de station",
+ "stationIdHelp": "Identifiant de station du chargeur, ou \"*\" pour correspondre à tous les chargeurs sans règle spécifique.",
+ "title": "Transmetteur OCPP",
+ "upstream": "En amont",
+ "upstreamConnected": "Connecté",
+ "upstreamConnectedHelp": "Un chargeur est connecté et la liaison en amont est active.",
+ "upstreamDisconnected": "Déconnecté",
+ "upstreamDisconnectedHelp": "Un chargeur est connecté mais la connexion en amont a échoué.",
+ "upstreamIdle": "Aucun chargeur",
+ "upstreamIdleHelp": "Le transmetteur est actif mais aucun chargeur ne s'est encore connecté via cette règle.",
+ "upstreamStationId": "Identifiant de station en amont",
+ "upstreamStationIdHelp": "Remplace l'identifiant de station envoyé au serveur en amont. Laisser vide pour utiliser l'identifiant du chargeur.",
+ "upstreamUrl": "URL du serveur en amont",
+ "upstreamUrlHelp": "URL WebSocket OCPP du serveur externe (ex. wss://facturation.exemple.com/ocpp)."
+ },
"optimizer": {
"description": "Analyse les prévisions d'ensoleillement, les prix de l'électricité et vos habitudes de consommation afin d'optimiser la stratégie de gestion de la batterie et de recharge. Les données sont transmises au service d'optimisation evcc pour y être traitées. Se contente actuellement de calculer est visualiser les données, il n’y a pas encore de contrôle d’appareils possible.",
"enable": "Activer l'optimiseur",
diff --git a/server/http.go b/server/http.go
index 264530b40..c38531cca 100644
--- a/server/http.go
+++ b/server/http.go
@@ -339,6 +339,9 @@ func (s *HTTPd) RegisterSystemHandler(site *core.Site, pub publisher, cache *uti
routes["delete"+key] = route{Method: "DELETE", Pattern: "/" + key, HandlerFunc: settingsDeleteJsonHandler(key, pub, fun())}
}
+ // ocpp forwarder rules apply at runtime and republish via the ocpp package
+ routes["updateocppforwarder"] = route{Method: "POST", Pattern: "/ocppforwarder", HandlerFunc: updateOcppForwarderHandler}
+
for _, r := range routes {
api.Methods(r.Methods()...).Path(r.Pattern).Handler(r.HandlerFunc)
}
diff --git a/server/http_ocppforwarder_handler.go b/server/http_ocppforwarder_handler.go
new file mode 100644
index 000000000..2ffa4c38c
--- /dev/null
+++ b/server/http_ocppforwarder_handler.go
@@ -0,0 +1,45 @@
+package server
+
+import (
+ "encoding/json"
+ "net/http"
+
+ "github.com/evcc-io/evcc/charger/ocpp"
+ "github.com/evcc-io/evcc/core/keys"
+ "github.com/evcc-io/evcc/server/db/settings"
+)
+
+// updateOcppForwarderHandler persists the OCPP forwarder rules, restoring masked
+// secrets from the stored rules by station id, and applies them at runtime.
+func updateOcppForwarderHandler(w http.ResponseWriter, r *http.Request) {
+ var rules []ocpp.ForwarderRule
+ if err := json.NewDecoder(r.Body).Decode(&rules); err != nil {
+ jsonError(w, http.StatusBadRequest, err)
+ return
+ }
+
+ // restore masked secrets (password, caCert) from stored rules by station id
+ var old []ocpp.ForwarderRule
+ if err := settings.Json(keys.OcppForwarder, &old); err == nil {
+ stored := make(map[string]ocpp.ForwarderRule, len(old))
+ for _, o := range old {
+ stored[o.StationID] = o
+ }
+ for i := range rules {
+ if o, ok := stored[rules[i].StationID]; ok {
+ if err := mergeMaskedAny(&o, &rules[i]); err != nil {
+ jsonError(w, http.StatusInternalServerError, err)
+ return
+ }
+ }
+ }
+ }
+
+ if err := settings.SetJson(keys.OcppForwarder, rules); err != nil {
+ jsonError(w, http.StatusInternalServerError, err)
+ return
+ }
+ ocpp.ApplyForwarderRules(rules)
+
+ jsonWrite(w, true)
+}
diff --git a/tests/config-ocpp.spec.ts b/tests/config-ocpp.spec.ts
index 62b86eef9..afb9bbca8 100644
--- a/tests/config-ocpp.spec.ts
+++ b/tests/config-ocpp.spec.ts
@@ -1,6 +1,6 @@
-import { test, expect } from "@playwright/test";
+import { test, expect, type Page } from "@playwright/test";
import { start, stop, baseUrl } from "./evcc";
-import { startSimulator, stopSimulator, simulatorUrl } from "./simulator";
+import { startSimulator, stopSimulator, simulatorUrl, simulatorHost } from "./simulator";
import { expectModalVisible } from "./utils";
import axios from "axios";
@@ -8,6 +8,66 @@ test.use({ baseURL: baseUrl() });
const OCPP_STATION_ID = "test-station-001";
+// open the OCPP modal on the evcc config page
+async function openOcppModal(page: Page) {
+ await page.goto("/#/config");
+ const ocppModal = page.getByTestId("ocpp-modal");
+ await page.getByTestId("ocpp").getByRole("button", { name: "edit" }).click();
+ await expectModalVisible(ocppModal);
+ return ocppModal;
+}
+
+async function evccServerUrl(page: Page) {
+ const ocppModal = await openOcppModal(page);
+ return ocppModal.getByLabel("Server URL").inputValue();
+}
+
+// connect a charger to evcc via the simulator UI
+async function connectCharger(page: Page, serverUrl: string) {
+ await page.goto(simulatorUrl());
+ const card = page.getByTestId("ocpp-add-client");
+ await card.getByLabel("Server URL").fill(serverUrl);
+ await card.getByLabel("Station ID").fill(OCPP_STATION_ID);
+ await card.getByRole("button", { name: "Connect" }).click();
+}
+
+// enable/disable the mock upstream OCPP server via the simulator UI. the toggle
+// button name reflects the live state, so getByRole auto-waits past the initial fetch.
+async function setMockServer(
+ page: Page,
+ opts: { enabled: boolean; username?: string; password?: string }
+) {
+ await page.goto(simulatorUrl());
+ const server = page.getByTestId("ocpp-server");
+ if (opts.enabled) {
+ if (opts.username) await server.getByLabel("Username").fill(opts.username);
+ if (opts.password) await server.getByLabel("Password").fill(opts.password);
+ await server.getByRole("button", { name: "Enable server" }).click();
+ await expect(server.getByRole("button", { name: "Disable server" })).toBeVisible();
+ } else {
+ await server.getByRole("button", { name: "Disable server" }).click();
+ await expect(server.getByRole("button", { name: "Enable server" })).toBeVisible();
+ }
+}
+
+// open the forwarder editor for the given station (whatever the button state)
+async function openForwarderEditor(page: Page, stationId: string) {
+ const ocppModal = await openOcppModal(page);
+ const station = ocppModal.getByTestId("ocpp-station").filter({ hasText: stationId });
+ await station.getByRole("button").click();
+ const forwarderModal = page.getByTestId("ocppforwarder-modal");
+ await expectModalVisible(forwarderModal);
+ return { ocppModal, station, forwarderModal };
+}
+
+// enable the mock upstream server, connect a charger, open the forwarder editor
+async function startForwarding(page: Page, creds?: { username: string; password: string }) {
+ const serverUrl = await evccServerUrl(page);
+ await setMockServer(page, { enabled: true, ...creds });
+ await connectCharger(page, serverUrl);
+ return openForwarderEditor(page, OCPP_STATION_ID);
+}
+
test.beforeEach(async () => {
await startSimulator();
await start();
@@ -45,8 +105,10 @@ test.describe("ocpp", () => {
await page.goto("/#/config");
await ocppCard.getByRole("button", { name: "edit" }).click();
await expectModalVisible(ocppModal);
- await expect(ocppModal).toContainText("Detected station IDs");
- await expect(ocppModal).toContainText([OCPP_STATION_ID, "Unknown"].join(""));
+ await expect(ocppModal).toContainText("Station IDs");
+ const station = ocppModal.getByTestId("ocpp-station");
+ await expect(station).toContainText(OCPP_STATION_ID);
+ await expect(station.getByRole("img", { name: "Unknown" })).toBeVisible();
await expect(ocppModal).not.toContainText("No OCPP chargers detected.");
});
@@ -91,3 +153,98 @@ test.describe("ocpp", () => {
await expect(testResult).toContainText("No sponsor token configured.");
});
});
+
+test.describe("ocpp forwarder", () => {
+ test("unreachable upstream shows error", async ({ page }) => {
+ await connectCharger(page, await evccServerUrl(page));
+ const { ocppModal, station, forwarderModal } = await openForwarderEditor(page, OCPP_STATION_ID);
+
+ await forwarderModal.getByLabel("Upstream server URL").fill("ws://localhost:1/ocpp");
+ await forwarderModal.getByRole("button", { name: "Save" }).click();
+
+ // modal stays open; upstream is unreachable, so the error surfaces in place
+ await expectModalVisible(forwarderModal);
+ await expect(forwarderModal.getByTestId("ocppforwarder-error")).toBeVisible();
+
+ // remove the rule from the still-open modal
+ await forwarderModal.getByRole("button", { name: "Remove" }).click();
+ await expectModalVisible(ocppModal);
+ await expect(station.getByRole("button", { name: "No forwarding configured" })).toBeVisible();
+ });
+
+ test("connects and drops when upstream stops", async ({ page }) => {
+ const { forwarderModal } = await startForwarding(page);
+
+ await forwarderModal.getByLabel("Upstream server URL").fill(`ws://${simulatorHost()}`);
+ await forwarderModal.getByRole("button", { name: "Save" }).click();
+ await expect(forwarderModal.getByTestId("ocppforwarder-status")).toContainText("Connected");
+
+ // simulator UI shows our station as the active upstream connection
+ await page.goto(simulatorUrl());
+ const lastStation = page.getByTestId("ocpp-server-last-station");
+ await expect(lastStation).toContainText(OCPP_STATION_ID);
+ await expect(lastStation).toContainText("active");
+
+ // disabling the upstream server drops the forwarder connection
+ await setMockServer(page, { enabled: false });
+ const reopened = await openForwarderEditor(page, OCPP_STATION_ID);
+ await expect(reopened.forwarderModal.getByTestId("ocppforwarder-status")).toContainText(
+ "Not connected"
+ );
+ });
+
+ test("basic auth", async ({ page }) => {
+ const { forwarderModal } = await startForwarding(page, {
+ username: "user",
+ password: "secret",
+ });
+
+ await forwarderModal.getByLabel("Upstream server URL").fill(`ws://${simulatorHost()}`);
+ await forwarderModal.getByLabel("Username").fill("user");
+ await forwarderModal.getByLabel("Password").fill("secret");
+ await forwarderModal.getByRole("button", { name: "Save" }).click();
+ await expect(forwarderModal.getByTestId("ocppforwarder-status")).toContainText("Connected");
+
+ // simulator UI shows our station as the active upstream connection
+ await page.goto(simulatorUrl());
+ await expect(page.getByTestId("ocpp-server-last-station")).toContainText(OCPP_STATION_ID);
+ });
+
+ test("param change reconnects", async ({ page }) => {
+ const { forwarderModal } = await startForwarding(page, {
+ username: "user",
+ password: "secret",
+ });
+ const status = forwarderModal.getByTestId("ocppforwarder-status");
+
+ // wrong password is rejected
+ await forwarderModal.getByLabel("Upstream server URL").fill(`ws://${simulatorHost()}`);
+ await forwarderModal.getByLabel("Username").fill("user");
+ await forwarderModal.getByLabel("Password").fill("wrong");
+ await forwarderModal.getByRole("button", { name: "Save" }).click();
+ await expect(forwarderModal.getByTestId("ocppforwarder-error")).toBeVisible();
+ await expect(status).toContainText("Not connected");
+
+ // fixing the password re-establishes the connection
+ await forwarderModal.getByLabel("Password").fill("secret");
+ await forwarderModal.getByRole("button", { name: "Save" }).click();
+ await expect(status).toContainText("Connected");
+ });
+
+ test("removing rule stops forwarding", async ({ page }) => {
+ const { ocppModal, station, forwarderModal } = await startForwarding(page);
+
+ await forwarderModal.getByLabel("Upstream server URL").fill(`ws://${simulatorHost()}`);
+ await forwarderModal.getByRole("button", { name: "Save" }).click();
+ await expect(forwarderModal.getByTestId("ocppforwarder-status")).toContainText("Connected");
+
+ // removing the rule tears down forwarding
+ await forwarderModal.getByRole("button", { name: "Remove" }).click();
+ await expectModalVisible(ocppModal);
+ await expect(station.getByRole("button", { name: "No forwarding configured" })).toBeVisible();
+
+ // simulator UI no longer shows an active upstream connection
+ await page.goto(simulatorUrl());
+ await expect(page.getByTestId("ocpp-server-last-station")).not.toContainText("active");
+ });
+});
diff --git a/tests/simulator/api.ts b/tests/simulator/api.ts
index ae8f105f4..c46395d0f 100644
--- a/tests/simulator/api.ts
+++ b/tests/simulator/api.ts
@@ -2,6 +2,7 @@ import bodyParser from "body-parser";
import type { Connect, ViteDevServer } from "vite";
import type { ServerResponse } from "http";
import { OcppClient } from "./ocppClient";
+import { ocppServer } from "./ocppServer";
const ocppClients = new Map();
@@ -47,6 +48,7 @@ const stateApiMiddleware = (
res.end();
process.exit();
} else if (req.originalUrl === "/api/state") {
+ updateOcppState();
res.end(JSON.stringify(state));
} else {
next();
@@ -169,22 +171,19 @@ const ocppMiddleware = (
const client = new OcppClient(stationId, serverUrl);
ocppClients.set(stationId, client);
-
- client
- .connect()
- .then(() => client.bootNotification())
- .then((response) => {
- console.log("[simulator] OCPP BootNotification response:", response);
- updateOcppState();
- res.end(JSON.stringify({ status: "connected", stationId, response }));
- })
- .catch((error) => {
- console.error("[simulator] OCPP connection error:", error);
- ocppClients.delete(stationId);
- updateOcppState();
- res.statusCode = 500;
- res.end(JSON.stringify({ error: error.message }));
- });
+ // fire-and-forget: the client retries until connected and boots itself on open
+ client.connect();
+ updateOcppState();
+ console.log(`[simulator] OCPP client ${stationId} connecting to ${serverUrl}`);
+ res.end(JSON.stringify({ status: "connecting", stationId }));
+ } else if (req.method === "POST" && req.originalUrl === "/api/ocpp/server") {
+ console.log("[simulator] POST /api/ocpp/server");
+ // @ts-expect-error Property 'body' does not exist on type 'IncomingMessage'
+ const { enabled, username, password } = req.body;
+ ocppServer.configure({ enabled: !!enabled, username, password });
+ res.end(JSON.stringify(ocppServer.status()));
+ } else if (req.method === "GET" && req.originalUrl === "/api/ocpp/server") {
+ res.end(JSON.stringify(ocppServer.status()));
} else if (req.method === "POST" && req.originalUrl === "/api/ocpp/disconnect") {
console.log("[simulator] POST /api/ocpp/disconnect");
// @ts-expect-error Property 'body' does not exist on type 'IncomingMessage'
@@ -216,6 +215,9 @@ export default () => ({
enforce: "pre",
configureServer(server: ViteDevServer) {
console.log("[simulator] configured");
+ if (server.httpServer) {
+ ocppServer.attach(server.httpServer);
+ }
return () => {
server.middlewares.use(loggingMiddleware);
server.middlewares.use(bodyParser.json());
diff --git a/tests/simulator/ocppClient.ts b/tests/simulator/ocppClient.ts
index eae7b22e1..3dddbc8c9 100644
--- a/tests/simulator/ocppClient.ts
+++ b/tests/simulator/ocppClient.ts
@@ -1,5 +1,7 @@
import WebSocket from "ws";
+const RECONNECT_DELAY = 2000;
+
export class OcppClient {
private ws: WebSocket | null = null;
private messageId = 0;
@@ -7,47 +9,66 @@ export class OcppClient {
private connected = false;
private stationId: string;
private serverUrl: string;
+ private shouldReconnect = false;
+ private reconnectTimer: ReturnType | null = null;
constructor(stationId: string, serverUrl: string) {
this.stationId = stationId;
this.serverUrl = serverUrl;
}
- async connect(): Promise {
- return new Promise((resolve, reject) => {
- const url = `${this.serverUrl}${this.stationId}`;
- console.log(`[OCPP Client] Connecting to ${url}`);
+ // connect starts the connection loop and returns immediately. The client keeps
+ // re-dialing every RECONNECT_DELAY until disconnect() is called, so it survives
+ // an unreachable server or a dropped connection (e.g. evcc restart).
+ connect(): void {
+ this.shouldReconnect = true;
+ this.dial();
+ }
- this.ws = new WebSocket(url, ["ocpp1.6"]);
+ private dial(): void {
+ const url = `${this.serverUrl}${this.stationId}`;
+ console.log(`[OCPP Client] Connecting to ${url}`);
- this.ws.on("open", () => {
- console.log(`[OCPP Client] Connected as ${this.stationId}`);
- this.connected = true;
- resolve();
- });
+ const ws = new WebSocket(url, ["ocpp1.6"]);
+ this.ws = ws;
- this.ws.on("message", (data: WebSocket.Data) => {
- const message = JSON.parse(data.toString());
- this.handleMessage(message);
- });
+ // force a close if the handshake never completes, so the retry loop kicks in
+ const handshakeTimer = setTimeout(() => {
+ if (!this.connected) {
+ console.log(`[OCPP Client] ${this.stationId} handshake timeout`);
+ ws.terminate();
+ }
+ }, 5000);
- this.ws.on("error", (error) => {
- console.error(`[OCPP Client] Error:`, error);
- reject(error);
- });
+ ws.on("open", () => {
+ clearTimeout(handshakeTimer);
+ console.log(`[OCPP Client] Connected as ${this.stationId}`);
+ this.connected = true;
+ // boot, then report the connector available, like a real charger would
+ this.bootNotification()
+ .then(() => this.statusNotification(1, "Available"))
+ .catch((error) => console.error(`[OCPP Client] init failed:`, error));
+ });
- this.ws.on("close", () => {
- console.log(`[OCPP Client] Disconnected`);
- this.connected = false;
- });
+ ws.on("message", (data: WebSocket.Data) => {
+ const message = JSON.parse(data.toString());
+ this.handleMessage(message);
+ });
- // Timeout after 5 seconds
- setTimeout(() => {
- if (!this.connected) {
- this.disconnect();
- reject(new Error("Connection timeout"));
- }
- }, 5000);
+ ws.on("error", (error) => {
+ console.error(`[OCPP Client] Error:`, (error as Error).message || error);
+ });
+
+ ws.on("close", () => {
+ clearTimeout(handshakeTimer);
+ this.connected = false;
+ if (this.ws === ws) this.ws = null;
+ if (this.shouldReconnect) {
+ console.log(
+ `[OCPP Client] ${this.stationId} disconnected, retrying in ${RECONNECT_DELAY}ms`
+ );
+ this.reconnectTimer = setTimeout(() => this.dial(), RECONNECT_DELAY);
+ }
});
}
@@ -92,6 +113,9 @@ export class OcppClient {
],
};
break;
+ case "ChangeConfiguration":
+ response = { status: "Accepted" };
+ break;
case "ChangeAvailability":
response = { status: "Accepted" };
break;
@@ -102,9 +126,17 @@ export class OcppClient {
response = { status: "Accepted" };
break;
case "TriggerMessage":
- // Handle trigger and send the requested message
- if (payload.requestedMessage === "StatusNotification") {
- this.statusNotification(1, "Available");
+ // actually send the requested message so evcc's setup wait completes
+ switch (payload.requestedMessage) {
+ case "StatusNotification":
+ this.statusNotification(1, "Available");
+ break;
+ case "MeterValues":
+ this.meterValues(1, 0, 0);
+ break;
+ case "BootNotification":
+ this.bootNotification();
+ break;
}
response = { status: "Accepted" };
break;
@@ -200,11 +232,16 @@ export class OcppClient {
}
disconnect(): void {
+ this.shouldReconnect = false;
+ if (this.reconnectTimer) {
+ clearTimeout(this.reconnectTimer);
+ this.reconnectTimer = null;
+ }
if (this.ws) {
this.ws.close();
this.ws = null;
- this.connected = false;
}
+ this.connected = false;
}
isConnected(): boolean {
diff --git a/tests/simulator/ocppServer.ts b/tests/simulator/ocppServer.ts
new file mode 100644
index 000000000..6bd0b7174
--- /dev/null
+++ b/tests/simulator/ocppServer.ts
@@ -0,0 +1,128 @@
+import { WebSocketServer, WebSocket } from "ws";
+import type { IncomingMessage } from "http";
+import type { Duplex } from "stream";
+
+// minimal surface we need from the shared http(s)/http2 server
+type UpgradableServer = {
+ on(
+ event: "upgrade",
+ listener: (req: IncomingMessage, socket: Duplex, head: Buffer) => void
+ ): unknown;
+};
+
+// OcppServer is a minimal upstream OCPP server used to test the evcc forwarder.
+// It is NOT spec compliant: it accepts a WebSocket connection (optionally behind
+// HTTP Basic Auth) and answers Calls with canned CallResults so the connection
+// stays alive. It shares the simulator's HTTP port via the vite httpServer.
+class OcppServer {
+ // echo the ocpp1.6 subprotocol so clients that require negotiation (evcc's
+ // coder/websocket dialer) complete the handshake
+ private wss = new WebSocketServer({ noServer: true, handleProtocols: () => "ocpp1.6" });
+ private sockets = new Set();
+
+ enabled = false;
+ username = "";
+ password = "";
+ lastStationId: string | null = null;
+
+ // attach hooks the WebSocket upgrade on the shared http server. OCPP upgrades
+ // are recognised by the "ocpp1.6" subprotocol; everything else (e.g. vite HMR)
+ // is left for other listeners.
+ attach(httpServer: UpgradableServer) {
+ httpServer.on("upgrade", (req: IncomingMessage, socket: Duplex, head: Buffer) => {
+ const protocols = String(req.headers["sec-websocket-protocol"] || "");
+ if (!protocols.includes("ocpp1.6")) return; // not an OCPP upgrade (e.g. vite HMR)
+
+ // server off: reject the OCPP upgrade so the client errors out and retries,
+ // rather than leaving the handshake hanging with no response
+ if (!this.enabled) {
+ socket.destroy();
+ return;
+ }
+
+ if ((this.username || this.password) && !this.checkAuth(req)) {
+ console.log("[ocpp-server] rejected: bad credentials");
+ socket.write('HTTP/1.1 401 Unauthorized\r\nWWW-Authenticate: Basic realm="ocpp"\r\n\r\n');
+ socket.destroy();
+ return;
+ }
+
+ this.wss.handleUpgrade(req, socket, head, (ws) => this.onConnection(ws, req));
+ });
+ }
+
+ configure(opts: { enabled: boolean; username?: string; password?: string }) {
+ this.enabled = opts.enabled;
+ this.username = opts.username || "";
+ this.password = opts.password || "";
+ if (!this.enabled) this.closeAll();
+ }
+
+ status() {
+ return {
+ enabled: this.enabled,
+ username: this.username,
+ password: this.password,
+ lastStationId: this.lastStationId,
+ connections: this.sockets.size,
+ };
+ }
+
+ private checkAuth(req: IncomingMessage): boolean {
+ const header = String(req.headers["authorization"] || "");
+ const expected = "Basic " + Buffer.from(`${this.username}:${this.password}`).toString("base64");
+ return header === expected;
+ }
+
+ private onConnection(ws: WebSocket, req: IncomingMessage) {
+ const path = (req.url || "/").split("?")[0];
+ const stationId = decodeURIComponent(path.replace(/^\/+/, "")) || "(root)";
+ this.lastStationId = stationId;
+ this.sockets.add(ws);
+ console.log(`[ocpp-server] ${stationId} connected (${this.sockets.size} active)`);
+
+ ws.on("message", (data) => this.handleMessage(ws, data.toString()));
+ ws.on("close", () => {
+ this.sockets.delete(ws);
+ console.log(`[ocpp-server] ${stationId} disconnected (${this.sockets.size} active)`);
+ });
+ }
+
+ // answer charger Calls (messageType 2) with a CallResult (messageType 3)
+ private handleMessage(ws: WebSocket, raw: string) {
+ let msg: unknown;
+ try {
+ msg = JSON.parse(raw);
+ } catch {
+ return;
+ }
+ if (!Array.isArray(msg) || msg[0] !== 2) return; // only respond to Calls
+ const [, messageId, action] = msg as [number, string, string, unknown];
+ ws.send(JSON.stringify([3, messageId, this.resultFor(action)]));
+ }
+
+ private resultFor(action: string): Record {
+ const now = new Date().toISOString();
+ switch (action) {
+ case "BootNotification":
+ return { status: "Accepted", currentTime: now, interval: 300 };
+ case "Heartbeat":
+ return { currentTime: now };
+ case "Authorize":
+ return { idTagInfo: { status: "Accepted" } };
+ case "StartTransaction":
+ return { transactionId: 1, idTagInfo: { status: "Accepted" } };
+ case "StopTransaction":
+ return { idTagInfo: { status: "Accepted" } };
+ default:
+ return {}; // StatusNotification, MeterValues, DataTransfer, ...
+ }
+ }
+
+ private closeAll() {
+ for (const ws of this.sockets) ws.close();
+ this.sockets.clear();
+ }
+}
+
+export const ocppServer = new OcppServer();
diff --git a/tests/simulator/src/Simulator.vue b/tests/simulator/src/Simulator.vue
index cfb9f5c57..c09b6673e 100644
--- a/tests/simulator/src/Simulator.vue
+++ b/tests/simulator/src/Simulator.vue
@@ -343,6 +343,74 @@
+
+