Add OCPP forwarder (#29154)

This commit is contained in:
Alexandre JARDON 2026-06-07 13:17:54 +02:00 • committed by GitHub
parent a8361b575b
commit 3758ce1957
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
28 changed files with 2182 additions and 137 deletions

View file

@ -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 |

View file

@ -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;

View file

@ -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) {

View file

@ -0,0 +1,107 @@
<template>
<div class="forward-group d-flex align-items-center gap-3 ms-auto">
<code v-if="rule" class="station-host text-truncate" :class="hostClass">{{ host }}</code>
<button
type="button"
class="station-forward btn d-flex align-items-center justify-content-center p-2 flex-shrink-0"
:class="buttonClass"
:title="title"
:aria-label="title"
data-bs-toggle="tooltip"
@click="edit"
>
<OcppForwardStatus :status="status" />
</button>
</div>
</template>
<script lang="ts">
import { defineComponent, type PropType } from "vue";
import OcppForwardStatus from "../MaterialIcon/OcppForwardStatus.vue";
import type { OcppForwarderRule } from "@/types/evcc";
import { openModal } from "@/configModal";
type ForwardStatus = "unconfigured" | "configured" | "error";
export default defineComponent({
name: "OcppForwarderButton",
components: { OcppForwardStatus },
props: {
stationId: { type: String, required: true },
rule: { type: Object as PropType<OcppForwarderRule>, default: undefined },
error: { type: String, default: undefined },
},
computed: {
status(): ForwardStatus {
if (!this.rule) return "unconfigured";
return this.error ? "error" : "configured";
},
// hostname of the upstream URL, scheme and path stripped
host(): string {
if (!this.rule) return "";
try {
return new URL(this.rule.upstreamUrl).host;
} catch {
return this.rule.upstreamUrl || "";
}
},
// host color tracks the button: green when configured, red on error
hostClass(): string {
return this.status === "error" ? "text-danger" : "text-success";
},
buttonClass(): string {
switch (this.status) {
case "configured":
return "text-success border border-success forward-bg-success";
case "error":
return "text-danger border border-danger forward-bg-error";
default:
return "forward-muted border border-dashed";
}
},
title(): string {
switch (this.status) {
case "configured":
return this.$t("config.ocpp.forwardingConfigured");
case "error":
return this.$t("config.ocpp.forwardingError");
default:
return this.$t("config.ocpp.forwardingOff");
}
},
},
methods: {
edit() {
openModal("ocppforwarder", { station: this.stationId });
},
},
});
</script>
<style scoped>
.station-host {
min-width: 0;
max-width: 16rem;
text-align: right;
font-size: var(--bs-body-font-size);
}
/* single line: the host (with its button) gives way and truncates before the identity */
.forward-group {
min-width: 0;
flex-shrink: 100;
}
/* button tints, scoped so global subtle colors stay untouched */
.forward-muted {
color: var(--bs-gray-light);
}
.border-dashed {
border-style: dashed !important;
}
.forward-bg-success {
background-color: color-mix(in srgb, var(--evcc-primary) 10%, transparent);
}
.forward-bg-error {
background-color: color-mix(in srgb, var(--evcc-red) 10%, transparent);
}
</style>

View file

@ -0,0 +1,280 @@
<template>
<JsonModal
name="ocppforwarder"
:title="$t('config.ocppforwarder.editTitle')"
:description="$t('config.ocppforwarder.description')"
endpoint="/config/ocppforwarder"
state-key="ocppforwarder.config"
no-buttons
:transform-read-values="transformReadValues"
:transform-write-values="transformWriteValues"
@changed="$emit('changed')"
>
<template
#default="{
values,
changes,
save,
}: {
values: OcppForwarderRule;
changes: boolean;
save: (close?: boolean) => void;
}"
>
<p v-if="ruleExists" class="mb-3" data-testid="ocppforwarder-status">
{{ $t("config.ocppforwarder.status") }}:
<span :class="connectionConnected ? 'text-success' : 'text-danger'">{{
connectionLabel
}}</span>
</p>
<div
v-if="sessionError"
class="alert alert-danger"
role="alert"
data-testid="ocppforwarder-error"
>
{{ sessionError }}
</div>
<FormRow
id="ocppforwarderUpstreamUrl"
:label="$t('config.ocppforwarder.upstreamUrl')"
:help="$t('config.ocppforwarder.upstreamUrlHelp')"
example="wss://billing.example.com/ocpp"
>
<input
id="ocppforwarderUpstreamUrl"
v-model="values.upstreamUrl"
type="text"
class="form-control"
inputmode="url"
spellcheck="false"
autocomplete="off"
required
/>
</FormRow>
<FormRow
id="ocppforwarderUsername"
:label="$t('config.ocppforwarder.username')"
:help="$t('config.ocppforwarder.usernameHelp')"
optional
>
<input
id="ocppforwarderUsername"
v-model="values.username"
type="text"
class="form-control"
spellcheck="false"
autocomplete="off"
/>
</FormRow>
<FormRow
id="ocppforwarderPassword"
:label="$t('config.ocppforwarder.password')"
:help="$t('config.ocppforwarder.passwordHelp')"
optional
>
<input
id="ocppforwarderPassword"
v-model="values.password"
type="password"
class="form-control"
autocomplete="new-password"
/>
</FormRow>
<FormRow
id="ocppforwarderUpstreamStationId"
:label="$t('config.ocppforwarder.upstreamStationId')"
:help="$t('config.ocppforwarder.upstreamStationIdHelp')"
optional
>
<input
id="ocppforwarderUpstreamStationId"
v-model="values.upstreamStationId"
type="text"
class="form-control"
:placeholder="targetStationId"
spellcheck="false"
autocomplete="off"
/>
</FormRow>
<PropertyCollapsible>
<template #advanced>
<FormRow
id="ocppforwarderReadOnly"
:label="$t('config.ocppforwarder.readOnly.label')"
:help="getReadOnlyHelp(values.readOnly)"
>
<div class="d-flex">
<input
id="ocppforwarderReadOnly"
v-model="values.readOnly"
class="form-check-input"
type="checkbox"
/>
<label class="form-check-label ms-2" for="ocppforwarderReadOnly">
{{ $t("config.ocppforwarder.readOnly.check") }}
</label>
</div>
</FormRow>
<FormRow
id="ocppforwarderInsecure"
:label="$t('config.ocppforwarder.labelInsecure')"
>
<div class="d-flex">
<input
id="ocppforwarderInsecure"
v-model="values.insecure"
class="form-check-input"
type="checkbox"
/>
<label class="form-check-label ms-2" for="ocppforwarderInsecure">
{{ $t("config.ocppforwarder.labelCheckInsecure") }}
</label>
</div>
</FormRow>
<FormRow
id="ocppforwarderCaCert"
:label="$t('config.ocppforwarder.labelCaCert')"
optional
>
<PropertyCertField id="ocppforwarderCaCert" v-model="values.caCert" />
</FormRow>
</template>
</PropertyCollapsible>
<div class="mt-4 d-flex justify-content-between gap-2 flex-column flex-sm-row">
<div
class="d-flex justify-content-between order-2 order-sm-1 gap-2 flex-grow-1 flex-sm-grow-0"
>
<button
type="button"
class="btn btn-link text-muted btn-cancel"
data-bs-dismiss="modal"
>
{{ $t("config.general.cancel") }}
</button>
<button
v-if="ruleExists"
type="button"
class="btn btn-link text-danger"
:disabled="removing"
@click="removeRule"
>
{{ $t("config.general.remove") }}
</button>
</div>
<button
v-if="changes"
type="button"
class="btn btn-primary order-1 order-sm-2 flex-grow-1 flex-sm-grow-0 px-4"
:disabled="!values.upstreamUrl"
@click="save(false)"
>
{{ $t("config.general.save") }}
</button>
<button
v-else
type="button"
class="btn btn-outline-primary order-1 order-sm-2 flex-grow-1 flex-sm-grow-0 px-4"
data-bs-dismiss="modal"
>
{{ $t("config.general.close") }}
</button>
</div>
</template>
</JsonModal>
</template>
<script lang="ts">
import { defineComponent } from "vue";
import JsonModal from "./JsonModal.vue";
import FormRow from "./FormRow.vue";
import PropertyCollapsible from "./PropertyCollapsible.vue";
import PropertyCertField from "./PropertyCertField.vue";
import type { OcppForwarderRule, OcppForwarderSession } from "@/types/evcc";
import { getModal, closeModal } from "@/configModal";
import api from "@/api";
import store from "@/store";
export default defineComponent({
name: "OcppForwarderModal",
components: { JsonModal, FormRow, PropertyCollapsible, PropertyCertField },
emits: ["changed"],
data() {
return { removing: false };
},
computed: {
// station id the modal is editing, carried via the config modal stack
targetStationId(): string {
return getModal("ocppforwarder")?.station || "";
},
rules(): OcppForwarderRule[] {
return store.state?.ocppforwarder?.config || [];
},
ruleExists(): boolean {
return this.rules.some((r) => r.stationId === this.targetStationId);
},
session(): OcppForwarderSession | undefined {
return store.state?.ocppforwarder?.status?.find(
(s) => s.chargerId === this.targetStationId
);
},
sessionError(): string | undefined {
return this.session?.error;
},
connectionConnected(): boolean {
return !!this.session?.upstreamConnected;
},
connectionLabel(): string {
return this.$t(
this.connectionConnected
? "config.ocpp.status.connected"
: "config.ocpp.status.configured"
);
},
},
methods: {
getReadOnlyHelp(readOnly?: boolean): string {
return this.$t(`config.ocppforwarder.readOnly.help.${readOnly ? "true" : "false"}`);
},
// pick the rule for the target station, or seed a new one prefilled with the station id
transformReadValues(rules: OcppForwarderRule[]): OcppForwarderRule {
const list = Array.isArray(rules) ? rules : [];
const existing = list.find((r) => r.stationId === this.targetStationId);
return existing
? { ...existing }
: {
stationId: this.targetStationId,
upstreamUrl: "",
};
},
// merge the edited rule back into the complete set that gets persisted
transformWriteValues(rule: OcppForwarderRule): OcppForwarderRule[] {
const list = this.rules.map((r) => ({ ...r }));
const index = list.findIndex((r) => r.stationId === rule.stationId);
if (index >= 0) {
list[index] = rule;
} else {
list.push(rule);
}
return list;
},
async removeRule() {
this.removing = true;
try {
const list = this.rules.filter((r) => r.stationId !== this.targetStationId);
const res = await api.post("/config/ocppforwarder", list, {
validateStatus: (code: number) => [200, 202, 400].includes(code),
});
if (res.status === 200 || res.status === 202) {
this.$emit("changed");
await closeModal();
}
} catch (e) {
console.error(e);
}
this.removing = false;
},
},
});
</script>

View file

@ -23,52 +23,40 @@
<hr class="my-4" />
<!-- No chargers message -->
<div
v-if="connectedStations.length === 0 && detectedStations.length === 0"
class="text-muted"
>
<div v-if="entries.length === 0" class="text-muted">
{{ $t("config.ocpp.noChargers") }}
</div>
<!-- Connection status -->
<div v-if="connectedStations.length > 0">
<h6 class="mb-3">{{ $t("config.ocpp.connectionStatus") }}</h6>
<p class="text-muted small mb-3">{{ $t("config.ocpp.connectionStatusHelp") }}</p>
<!-- Stations -->
<div v-else>
<h6 class="mb-3">{{ $t("config.ocpp.stations") }}</h6>
<p class="text-muted small mb-3">{{ $t("config.ocpp.stationsHelp") }}</p>
<ul class="list-group stations-list mb-4">
<ul class="list-unstyled d-flex flex-column gap-3 mb-0">
<li
v-for="station in connectedStations"
:key="station.id"
class="list-group-item d-flex justify-content-between align-items-center"
v-for="entry in entries"
:key="entry.id"
class="station-bar border rounded d-flex align-items-center gap-3 px-3 py-3"
data-testid="ocpp-station"
>
<code>{{ station.id }}</code>
<span
class="badge"
:class="statusBadgeClass(station.status)"
:title="$t(`config.ocpp.status.${station.status}`)"
data-bs-toggle="tooltip"
>
{{ $t(`config.ocpp.status.${station.status}`) }}
</span>
</li>
</ul>
<StatusIndicator
class="flex-shrink-0"
:variant="statusVariant(entry.status)"
:tooltip="$t(`config.ocpp.status.${entry.status}`)"
/>
<div class="station-identity text-truncate">
<div v-if="entry.title" class="fw-bold fs-6 text-truncate">
{{ entry.title }}
</div>
<!-- Detected chargers -->
<div v-if="detectedStations.length > 0">
<h6 class="mb-3">{{ $t("config.ocpp.detectedChargers") }}</h6>
<p class="text-muted small mb-3">{{ $t("config.ocpp.detectedHelp") }}</p>
<ul class="list-group stations-list">
<li
v-for="station in detectedStations"
:key="station.id"
class="list-group-item d-flex justify-content-between align-items-center"
>
<code>{{ station.id }}</code>
<span class="badge bg-secondary">
{{ $t(`config.ocpp.status.${station.status}`) }}
</span>
<code class="text-muted" :class="entry.title ? 'small' : 'fs-6'">{{
entry.id
}}</code>
</div>
<OcppForwarderButton
:station-id="entry.id"
:rule="entry.rule"
:error="entry.error"
/>
</li>
</ul>
</div>
@ -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<Ocpp>,
default: () => ({ config: { port: 0 }, status: { stations: [] } }),
},
stationTitles: {
type: Object as PropType<Record<string, string>>,
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<string, StationEntry>();
const published = new Set<string>();
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;
}
</style>

View file

@ -0,0 +1,75 @@
<template>
<span
ref="root"
class="d-flex align-items-center gap-2 text-nowrap evcc-gray"
:role="tooltip ? 'img' : undefined"
:aria-label="tooltip || undefined"
:title="tooltip || undefined"
:data-bs-toggle="tooltip ? 'tooltip' : undefined"
>
<span class="d-inline-block rounded-circle status-dot" :class="dotClass"></span>
<slot />
</span>
</template>
<script lang="ts">
import { defineComponent, type PropType } from "vue";
import Tooltip from "bootstrap/js/dist/tooltip";
type Variant = "success" | "warning" | "muted";
export default defineComponent({
name: "StatusIndicator",
props: {
variant: { type: String as PropType<Variant>, default: "muted" },
tooltip: { type: String, default: "" },
},
data() {
return { tooltipInstance: null as Tooltip | null };
},
computed: {
dotClass(): string {
switch (this.variant) {
case "success":
return "bg-success";
case "warning":
return "bg-warning";
default:
return "border border-secondary";
}
},
},
watch: {
tooltip() {
this.initTooltip();
},
},
mounted() {
this.initTooltip();
},
beforeUnmount() {
this.tooltipInstance?.dispose();
},
methods: {
initTooltip() {
this.$nextTick(() => {
this.tooltipInstance?.dispose();
this.tooltipInstance = null;
const el = this.$refs["root"] as Element | undefined;
if (el && this.tooltip) {
// explicit title: bootstrap clears the attr, so :title alone goes stale
this.tooltipInstance = new Tooltip(el, { title: this.tooltip });
}
});
},
},
});
</script>
<style scoped>
.status-dot {
width: 0.8rem;
height: 0.8rem;
box-sizing: border-box;
}
</style>

View file

@ -0,0 +1,37 @@
<template>
<svg :style="svgStyle" viewBox="0 0 24 24">
<!-- configured and connected: cloud with check -->
<path
v-if="status === 'configured'"
fill="currentColor"
d="m10.325 14.125l-1.4-1.4q-.3-.3-.7-.3t-.7.3t-.3.713t.3.712L9.65 16.3q.3.3.7.3t.7-.3l4.225-4.225q.3-.3.3-.725t-.3-.725t-.725-.3t-.725.3zM6.5 20q-2.275 0-3.887-1.575T1 14.575q0-1.95 1.175-3.475T5.25 9.15q.625-2.3 2.5-3.725T12 4q2.925 0 4.963 2.038T19 11q1.725.2 2.863 1.488T23 15.5q0 1.875-1.312 3.188T18.5 20zm0-2h12q1.05 0 1.775-.725T21 15.5t-.725-1.775T18.5 13H17v-2q0-2.075-1.463-3.538T12 6T8.463 7.463T7 11h-.5q-1.45 0-2.475 1.025T3 14.5t1.025 2.475T6.5 18m5.5-6"
/>
<!-- configured but not connected: cloud with alert -->
<path
v-else-if="status === 'error'"
fill="currentColor"
d="M6.5 20q-2.275 0-3.887-1.575T1 14.575q0-1.95 1.175-3.475T5.25 9.15q.625-2.3 2.5-3.725T12 4q2.925 0 4.963 2.038T19 11q1.725.2 2.863 1.488T23 15.5q0 1.875-1.312 3.188T18.5 20zm0-2h12q1.05 0 1.775-.725T21 15.5t-.725-1.775T18.5 13H17v-2q0-2.075-1.463-3.538T12 6T8.463 7.463T7 11h-.5q-1.45 0-2.475 1.025T3 14.5t1.025 2.475T6.5 18m5.5-2q.425 0 .713-.288T13 15t-.288-.712T12 14t-.712.288T11 15t.288.713T12 16m0-3.5q.425 0 .713-.288T13 11.5V9q0-.425-.288-.712T12 8t-.712.288T11 9v2.5q0 .425.288.713T12 12.5"
/>
<!-- no forwarding configured: plain cloud -->
<path
v-else
fill="currentColor"
d="M6.5 20q-2.275 0-3.887-1.575T1 14.575q0-1.95 1.175-3.475T5.25 9.15q.625-2.3 2.5-3.725T12 4q2.925 0 4.963 2.038T19 11q1.725.2 2.863 1.488T23 15.5q0 1.875-1.312 3.188T18.5 20zm0-2h12q1.05 0 1.775-.725T21 15.5t-.725-1.775T18.5 13H17v-2q0-2.075-1.463-3.538T12 6T8.463 7.463T7 11h-.5q-1.45 0-2.475 1.025T3 14.5t1.025 2.475T6.5 18m5.5-6"
/>
</svg>
</template>
<script lang="ts">
import { defineComponent, type PropType } from "vue";
import icon from "@/mixins/icon";
type ForwardStatus = "unconfigured" | "configured" | "error";
export default defineComponent({
name: "OcppForwardStatus",
mixins: [icon],
props: {
status: { type: String as PropType<ForwardStatus>, default: "unconfigured" },
},
});
</script>

View file

@ -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<string, string> {
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<ModalResult> {
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<void> {
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);

View file

@ -121,6 +121,7 @@ export interface State {
config?: string;
database?: string;
ocpp?: Ocpp;
ocppforwarder?: ConfigStatus<OcppForwarderRule[], OcppForwarderSession[]>;
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[];

View file

@ -449,7 +449,8 @@
:yamlSource="eebus?.yamlSource"
@changed="loadDirty"
/>
<OcppModal :ocpp="ocpp" />
<OcppModal :ocpp="ocpp" :stationTitles="stationTitles" />
<OcppForwarderModal @changed="loadDirty" />
<BackupRestoreModal v-bind="backupRestoreProps" />
<SecurityModal :auth-disabled="authDisabled" />
<ApiKeyModal :auth-disabled="authDisabled" />
@ -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<string, string> {
const map: Record<string, string> = {};
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 || [];

View file

@ -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 {

746
charger/ocpp/forwarder.go Normal file
View file

@ -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
}

View file

@ -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())

View file

@ -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/<stationId>", 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,

View file

@ -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

View file

@ -16,6 +16,7 @@ const (
MessagingEvents = "messagingEvents"
ModbusProxy = "modbusproxy"
Ocpp = "ocpp"
OcppForwarder = "ocppforwarder"
Tariffs = "tariffs"
TariffRefs = "tariffRefs"
Version = "version"

View file

@ -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.

View file

@ -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",

View file

@ -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",

View file

@ -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",

View file

@ -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)
}

View file

@ -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)
}

View file

@ -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");
});
});

View file

@ -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<string, OcppClient>();
@ -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);
// fire-and-forget: the client retries until connected and boots itself on open
client.connect();
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 }));
});
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());

View file

@ -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<typeof setTimeout> | null = null;
constructor(stationId: string, serverUrl: string) {
this.stationId = stationId;
this.serverUrl = serverUrl;
}
async connect(): Promise<void> {
return new Promise((resolve, reject) => {
// 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();
}
private dial(): void {
const url = `${this.serverUrl}${this.stationId}`;
console.log(`[OCPP Client] Connecting to ${url}`);
this.ws = new WebSocket(url, ["ocpp1.6"]);
const ws = new WebSocket(url, ["ocpp1.6"]);
this.ws = ws;
this.ws.on("open", () => {
// 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);
ws.on("open", () => {
clearTimeout(handshakeTimer);
console.log(`[OCPP Client] Connected as ${this.stationId}`);
this.connected = true;
resolve();
// 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("message", (data: WebSocket.Data) => {
ws.on("message", (data: WebSocket.Data) => {
const message = JSON.parse(data.toString());
this.handleMessage(message);
});
this.ws.on("error", (error) => {
console.error(`[OCPP Client] Error:`, error);
reject(error);
ws.on("error", (error) => {
console.error(`[OCPP Client] Error:`, (error as Error).message || error);
});
this.ws.on("close", () => {
console.log(`[OCPP Client] Disconnected`);
ws.on("close", () => {
clearTimeout(handshakeTimer);
this.connected = false;
});
// Timeout after 5 seconds
setTimeout(() => {
if (!this.connected) {
this.disconnect();
reject(new Error("Connection timeout"));
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);
}
}, 5000);
});
}
@ -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") {
// 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 {

View file

@ -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<WebSocket>();
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<string, unknown> {
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();

View file

@ -343,6 +343,74 @@
</div>
</div>
<!-- Upstream OCPP Server (for forwarder testing) -->
<div class="card mb-3" data-testid="ocpp-server">
<div class="card-body">
<div class="d-flex justify-content-between align-items-center mb-2">
<h5 class="card-title mb-0">OCPP Server</h5>
<span class="badge" :class="ocppServer.enabled ? 'bg-success' : 'bg-secondary'">
{{ ocppServer.enabled ? "Listening" : "Off" }}
</span>
</div>
<p class="card-text text-muted small mb-3">
Upstream server for forwarder testing. Point a rule at
<code>ws://{{ host }}/&lt;stationId&gt;</code>
</p>
<div class="row">
<label for="ocppServerUsername" class="col-sm-6 col-form-label">
Username <small class="text-muted">(optional)</small>
</label>
<div class="col-sm-6">
<div class="input-group mb-3">
<input
id="ocppServerUsername"
v-model="ocppServer.username"
type="text"
class="form-control"
autocomplete="off"
:disabled="ocppServer.enabled"
/>
</div>
</div>
</div>
<div class="row">
<label for="ocppServerPassword" class="col-sm-6 col-form-label">
Password <small class="text-muted">(optional)</small>
</label>
<div class="col-sm-6">
<div class="input-group mb-3">
<input
id="ocppServerPassword"
v-model="ocppServer.password"
type="text"
class="form-control"
autocomplete="off"
:disabled="ocppServer.enabled"
/>
</div>
</div>
</div>
<div class="row mb-3">
<label class="col-sm-6 col-form-label">Last station</label>
<div class="col-sm-6 col-form-label" data-testid="ocpp-server-last-station">
{{ ocppServer.lastStationId || "—" }}
<span v-if="ocppServer.connections" class="text-muted small">
({{ ocppServer.connections }} active)
</span>
</div>
</div>
<button
type="button"
class="btn w-100"
:class="ocppServer.enabled ? 'btn-danger' : 'btn-primary'"
data-testid="ocpp-server-toggle"
@click="toggleOcppServer"
>
{{ ocppServer.enabled ? "Disable server" : "Enable server" }}
</button>
</div>
</div>
<div class="p-4 text-center fixed-bottom bg-light text-dark bg-opacity-75">
<button type="submit" class="btn btn-primary">Apply changes</button>
</div>
@ -381,14 +449,32 @@ export default defineComponent({
ocppServerUrl: "ws://127.0.0.1:8887/",
ocppStationId: "",
connecting: false,
ocppServer: {
enabled: false,
username: "",
password: "",
lastStationId: null as string | null,
connections: 0,
},
ocppServerPoll: undefined as ReturnType<typeof setInterval> | undefined,
};
},
computed: {
host(): string {
return window.location.host;
},
},
mounted() {
this.checkMockLoginMode();
if (!this.mockLoginMode) {
this.load();
this.refreshOcppServer();
this.ocppServerPoll = setInterval(() => this.refreshOcppServer(), 2000);
}
},
beforeUnmount() {
clearInterval(this.ocppServerPoll);
},
methods: {
checkMockLoginMode() {
const urlParams = new URLSearchParams(window.location.search);
@ -435,6 +521,36 @@ export default defineComponent({
this.connecting = false;
}
},
// refresh server status without clobbering credentials the user is editing
async refreshOcppServer() {
try {
const { data } = await axios.get("/api/ocpp/server");
this.ocppServer.enabled = data.enabled;
this.ocppServer.lastStationId = data.lastStationId;
this.ocppServer.connections = data.connections;
if (data.enabled) {
this.ocppServer.username = data.username;
this.ocppServer.password = data.password;
}
} catch (error) {
console.error("Failed to read OCPP server status:", error);
}
},
async toggleOcppServer() {
try {
const { data } = await axios.post("/api/ocpp/server", {
enabled: !this.ocppServer.enabled,
username: this.ocppServer.username,
password: this.ocppServer.password,
});
this.ocppServer.enabled = data.enabled;
this.ocppServer.lastStationId = data.lastStationId;
this.ocppServer.connections = data.connections;
} catch (error) {
console.error("Failed to toggle OCPP server:", error);
alert("Failed to toggle OCPP server");
}
},
async disconnectOcpp(stationId: string) {
try {
await axios.post("/api/ocpp/disconnect", { stationId });