import { BootNotificationResponse, RegistrationStatus } from '../types/ocpp/Responses';
import ChargingStationConfiguration, { ConfigurationKey } from '../types/ChargingStationConfiguration';
import ChargingStationTemplate, { CurrentType, PowerUnits, Voltage } from '../types/ChargingStationTemplate';
-import { ConnectorPhaseRotation, StandardParametersKey, SupportedFeatureProfiles } from '../types/ocpp/Configuration';
+import { ConnectorPhaseRotation, StandardParametersKey, SupportedFeatureProfiles, VendorDefaultParametersKey } from '../types/ocpp/Configuration';
import { ConnectorStatus, SampledValueTemplate } from '../types/Connectors';
import { MeterValueMeasurand, MeterValuePhase } from '../types/ocpp/MeterValues';
import { WSError, WebSocketCloseEventStatusCode } from '../types/WebSocket';
import { ChargePointStatus } from '../types/ocpp/ChargePointStatus';
import { ChargingProfile } from '../types/ocpp/ChargingProfile';
import ChargingStationInfo from '../types/ChargingStationInfo';
+import { ChargingStationWorkerMessageEvents } from '../types/ChargingStationWorker';
import { ClientRequestArgs } from 'http';
import Configuration from '../utils/Configuration';
import Constants from '../utils/Constants';
import crypto from 'crypto';
import fs from 'fs';
import logger from '../utils/Logger';
+import { parentPort } from 'worker_threads';
import path from 'path';
export default class ChargingStation {
private bootNotificationResponse!: BootNotificationResponse | null;
private connectorsConfigurationHash!: string;
private ocppIncomingRequestService!: OCPPIncomingRequestService;
- private readonly messageQueue: string[];
- private wsConnectionUrl!: URL;
+ private readonly messageBuffer: Set<string>;
+ private wsConfiguredConnectionUrl!: URL;
private wsConnectionRestarted: boolean;
private stopped: boolean;
private autoReconnectRetryCount: number;
this.autoReconnectRetryCount = 0;
this.requests = new Map<string, CachedRequest>();
- this.messageQueue = new Array<string>();
+ this.messageBuffer = new Set<string>();
this.authorizedTags = this.getAuthorizedTags();
}
+ get wsConnectionUrl(): URL {
+ return this.getSupervisionUrlOcppConfiguration() ? new URL(this.getConfigurationKey(this.stationInfo.supervisionUrlOcppKey ?? VendorDefaultParametersKey.ConnectionUrl).value + '/' + this.stationInfo.chargingStationId) : this.wsConfiguredConnectionUrl;
+ }
+
public logPrefix(): string {
return Utils.logPrefix(` ${this.stationInfo.chargingStationId} |`);
}
this.wsConnection.on('ping', this.onPing.bind(this));
// Handle WebSocket pong
this.wsConnection.on('pong', this.onPong.bind(this));
+ parentPort.postMessage({ id: ChargingStationWorkerMessageEvents.STARTED, data: { id: this.stationInfo.chargingStationId } });
}
public async stop(reason: StopTransactionReason = StopTransactionReason.NONE): Promise<void> {
this.performanceStatistics.stop();
}
this.bootNotificationResponse = null;
+ parentPort.postMessage({ id: ChargingStationWorkerMessageEvents.STOPPED, data: { id: this.stationInfo.chargingStationId } });
this.stopped = true;
}
});
}
- public addConfigurationKey(key: string | StandardParametersKey, value: string, readonly = false, visible = true, reboot = false): void {
+ public addConfigurationKey(key: string | StandardParametersKey, value: string, options: { readonly?: boolean, visible?: boolean, reboot?: boolean } = { readonly: false, visible: true, reboot: false }): void {
const keyFound = this.getConfigurationKey(key);
+ const readonly = options.readonly;
+ const visible = options.visible;
+ const reboot = options.reboot;
if (!keyFound) {
this.configuration.configurationKey.push({
key,
!cpReplaced && this.getConnectorStatus(connectorId).chargingProfiles?.push(cp);
}
- public resetTransactionOnConnector(connectorId: number): void {
- this.getConnectorStatus(connectorId).authorized = false;
+ public resetConnectorStatus(connectorId: number): void {
+ this.getConnectorStatus(connectorId).idTagLocalAuthorized = false;
+ this.getConnectorStatus(connectorId).idTagAuthorized = false;
+ this.getConnectorStatus(connectorId).transactionRemoteStarted = false;
this.getConnectorStatus(connectorId).transactionStarted = false;
+ delete this.getConnectorStatus(connectorId).localAuthorizeIdTag;
delete this.getConnectorStatus(connectorId).authorizeIdTag;
delete this.getConnectorStatus(connectorId).transactionId;
delete this.getConnectorStatus(connectorId).transactionIdTag;
this.stopMeterValues(connectorId);
}
- public addToMessageQueue(message: string): void {
- let dups = false;
- // Handle dups in message queue
- for (const bufferedMessage of this.messageQueue) {
- // Message already in the queue
- if (message === bufferedMessage) {
- dups = true;
- break;
- }
- }
- if (!dups) {
- // Queue message
- this.messageQueue.push(message);
- }
+ public bufferMessage(message: string): void {
+ this.messageBuffer.add(message);
}
- private flushMessageQueue() {
- if (!Utils.isEmptyArray(this.messageQueue)) {
- this.messageQueue.forEach((message, index) => {
- this.messageQueue.splice(index, 1);
+ private flushMessageBuffer() {
+ if (this.messageBuffer.size > 0) {
+ this.messageBuffer.forEach((message) => {
// TODO: evaluate the need to track performance
this.wsConnection.send(message);
+ this.messageBuffer.delete(message);
});
}
}
+ private getSupervisionUrlOcppConfiguration(): boolean {
+ return this.stationInfo.supervisionUrlOcppConfiguration ?? false;
+ }
+
private getChargingStationId(stationTemplate: ChargingStationTemplate): string {
// In case of multiple instances: add instance index to charging station id
const instanceIndex = process.env.CF_INSTANCE_INDEX ?? 0;
stationTemplateFromFile = JSON.parse(fs.readFileSync(fileDescriptor, 'utf8')) as ChargingStationTemplate;
fs.closeSync(fileDescriptor);
} catch (error) {
- FileUtils.handleFileException(this.logPrefix(), 'Template', this.stationTemplateFile, error);
+ FileUtils.handleFileException(this.logPrefix(), 'Template', this.stationTemplateFile, error as NodeJS.ErrnoException);
}
const stationInfo: ChargingStationInfo = stationTemplateFromFile ?? {} as ChargingStationInfo;
+ stationInfo.wsOptions = stationTemplateFromFile?.wsOptions ?? {};
if (!Utils.isEmptyArray(stationTemplateFromFile.power)) {
stationTemplateFromFile.power = stationTemplateFromFile.power as number[];
const powerArrayRandomIndex = Math.floor(Utils.secureRandom() * stationTemplateFromFile.power.length);
return stationInfo;
}
- private getOCPPVersion(): OCPPVersion {
+ private getOcppVersion(): OCPPVersion {
return this.stationInfo.ocppVersion ? this.stationInfo.ocppVersion : OCPPVersion.VERSION_16;
}
private initialize(): void {
this.stationInfo = this.buildStationInfo();
+ this.configuration = this.getTemplateChargingStationConfiguration();
+ delete this.stationInfo.Configuration;
this.bootNotificationRequest = {
chargePointModel: this.stationInfo.chargePointModel,
chargePointVendor: this.stationInfo.chargePointVendor,
...!Utils.isUndefined(this.stationInfo.chargeBoxSerialNumberPrefix) && { chargeBoxSerialNumber: this.stationInfo.chargeBoxSerialNumberPrefix },
...!Utils.isUndefined(this.stationInfo.firmwareVersion) && { firmwareVersion: this.stationInfo.firmwareVersion },
};
- this.configuration = this.getTemplateChargingStationConfiguration();
- this.wsConnectionUrl = new URL(this.getSupervisionURL().href + '/' + this.stationInfo.chargingStationId);
// Build connectors if needed
const maxConnectors = this.getMaxNumberOfConnectors();
if (maxConnectors <= 0) {
// Initialize transaction attributes on connectors
for (const connectorId of this.connectors.keys()) {
if (connectorId > 0 && !this.getConnectorStatus(connectorId)?.transactionStarted) {
- this.initTransactionAttributesOnConnector(connectorId);
+ this.initializeConnectorStatus(connectorId);
}
}
- switch (this.getOCPPVersion()) {
+ this.wsConfiguredConnectionUrl = new URL(this.getConfiguredSupervisionUrl().href + '/' + this.stationInfo.chargingStationId);
+ switch (this.getOcppVersion()) {
case OCPPVersion.VERSION_16:
this.ocppIncomingRequestService = new OCPP16IncomingRequestService(this);
this.ocppRequestService = new OCPP16RequestService(this, new OCPP16ResponseService(this));
break;
default:
- this.handleUnsupportedVersion(this.getOCPPVersion());
+ this.handleUnsupportedVersion(this.getOcppVersion());
break;
}
// OCPP parameters
- this.initOCPPParameters();
+ this.initOcppParameters();
if (this.stationInfo.autoRegister) {
this.bootNotificationResponse = {
currentTime: new Date().toISOString(),
}
}
- private initOCPPParameters(): void {
+ private initOcppParameters(): void {
+ if (this.getSupervisionUrlOcppConfiguration() && !this.getConfigurationKey(this.stationInfo.supervisionUrlOcppKey ?? VendorDefaultParametersKey.ConnectionUrl)) {
+ this.addConfigurationKey(VendorDefaultParametersKey.ConnectionUrl, this.getConfiguredSupervisionUrl().href, { reboot: true });
+ }
if (!this.getConfigurationKey(StandardParametersKey.SupportedFeatureProfiles)) {
this.addConfigurationKey(StandardParametersKey.SupportedFeatureProfiles, `${SupportedFeatureProfiles.Core},${SupportedFeatureProfiles.Local_Auth_List_Management},${SupportedFeatureProfiles.Smart_Charging}`);
}
- this.addConfigurationKey(StandardParametersKey.NumberOfConnectors, this.getNumberOfConnectors().toString(), true);
+ this.addConfigurationKey(StandardParametersKey.NumberOfConnectors, this.getNumberOfConnectors().toString(), { readonly: true });
if (!this.getConfigurationKey(StandardParametersKey.MeterValuesSampledData)) {
this.addConfigurationKey(StandardParametersKey.MeterValuesSampledData, MeterValueMeasurand.ENERGY_ACTIVE_IMPORT_REGISTER);
}
await this.startMessageSequence();
this.stopped && (this.stopped = false);
if (this.wsConnectionRestarted && this.isWebSocketConnectionOpened()) {
- this.flushMessageQueue();
+ this.flushMessageBuffer();
}
} else {
logger.error(`${this.logPrefix()} Registration failure: max retries reached (${this.getRegistrationMaxRetries()}) or retry disabled (${this.getRegistrationMaxRetries()})`);
}
} catch (error) {
// Log
- logger.error('%s Incoming OCPP message %j matching cached request %j processing error %j', this.logPrefix(), data, this.requests.get(messageId), error);
+ logger.error('%s Incoming OCPP message %j matching cached request %j processing error %j', this.logPrefix(), data.toString(), this.requests.get(messageId), error);
// Send error
- messageType === MessageType.CALL_MESSAGE && await this.ocppRequestService.sendError(messageId, error, commandName);
+ messageType === MessageType.CALL_MESSAGE && await this.ocppRequestService.sendError(messageId, error as OCPPError, commandName);
}
}
authorizedTags = JSON.parse(fs.readFileSync(fileDescriptor, 'utf8')) as string[];
fs.closeSync(fileDescriptor);
} catch (error) {
- FileUtils.handleFileException(this.logPrefix(), 'Authorization', authorizationFile, error);
+ FileUtils.handleFileException(this.logPrefix(), 'Authorization', authorizationFile, error as NodeJS.ErrnoException);
}
} else {
logger.info(this.logPrefix() + ' No authorization file given in template file ' + this.stationTemplateFile);
}
}
- private getSupervisionURL(): URL {
- const supervisionUrls = Utils.cloneObject<string | string[]>(this.stationInfo.supervisionURL ? this.stationInfo.supervisionURL : Configuration.getSupervisionURLs());
+ private getConfiguredSupervisionUrl(): URL {
+ const supervisionUrls = Utils.cloneObject<string | string[]>(this.stationInfo.supervisionUrl ?? Configuration.getSupervisionUrls());
let indexUrl = 0;
if (!Utils.isEmptyArray(supervisionUrls)) {
if (Configuration.getDistributeStationsToTenantsEqually()) {
}
}
- private openWSConnection(options?: ClientOptions & ClientRequestArgs, forceCloseOpened = false): void {
- options = options ?? {};
+ private openWSConnection(options: ClientOptions & ClientRequestArgs = this.stationInfo.wsOptions, forceCloseOpened = false): void {
options.handshakeTimeout = options?.handshakeTimeout ?? this.getConnectionTimeout() * 1000;
if (!Utils.isNullOrUndefined(this.stationInfo.supervisionUser) && !Utils.isNullOrUndefined(this.stationInfo.supervisionPassword)) {
options.auth = `${this.stationInfo.supervisionUser}:${this.stationInfo.supervisionPassword}`;
if (this.isWebSocketConnectionOpened() && forceCloseOpened) {
this.wsConnection.close();
}
- let protocol;
- switch (this.getOCPPVersion()) {
+ let protocol: string;
+ switch (this.getOcppVersion()) {
case OCPPVersion.VERSION_16:
protocol = 'ocpp' + OCPPVersion.VERSION_16;
break;
default:
- this.handleUnsupportedVersion(this.getOCPPVersion());
+ this.handleUnsupportedVersion(this.getOcppVersion());
break;
}
this.wsConnection = new WebSocket(this.wsConnectionUrl, protocol, options);
}
});
} catch (error) {
- FileUtils.handleFileException(this.logPrefix(), 'Authorization', authorizationFile, error);
+ FileUtils.handleFileException(this.logPrefix(), 'Authorization', authorizationFile, error as NodeJS.ErrnoException);
}
} else {
logger.info(this.logPrefix() + ' No authorization file given in template file ' + this.stationTemplateFile + '. Not monitoring changes');
}
});
} catch (error) {
- FileUtils.handleFileException(this.logPrefix(), 'Template', this.stationTemplateFile, error);
+ FileUtils.handleFileException(this.logPrefix(), 'Template', this.stationTemplateFile, error as NodeJS.ErrnoException);
}
}
if (this.autoReconnectRetryCount < this.getAutoReconnectMaxRetries() || this.getAutoReconnectMaxRetries() === -1) {
this.autoReconnectRetryCount++;
const reconnectDelay = (this.getReconnectExponentialDelay() ? Utils.exponentialDelay(this.autoReconnectRetryCount) : this.getConnectionTimeout() * 1000);
- const reconnectTimeout = reconnectDelay - 100;
+ const reconnectTimeout = (reconnectDelay - 100) > 0 && reconnectDelay;
logger.error(`${this.logPrefix()} WebSocket: connection retry in ${Utils.roundTo(reconnectDelay, 2)}ms, timeout ${reconnectTimeout}ms`);
await Utils.sleep(reconnectDelay);
logger.error(this.logPrefix() + ' WebSocket: reconnecting try #' + this.autoReconnectRetryCount.toString());
- this.openWSConnection({ handshakeTimeout: reconnectTimeout }, true);
+ this.openWSConnection({ ...this.stationInfo.wsOptions, handshakeTimeout: reconnectTimeout }, true);
this.wsConnectionRestarted = true;
} else if (this.getAutoReconnectMaxRetries() !== -1) {
logger.error(`${this.logPrefix()} WebSocket reconnect failure: max retries reached (${this.autoReconnectRetryCount}) or retry disabled (${this.getAutoReconnectMaxRetries()})`);
}
}
- private initTransactionAttributesOnConnector(connectorId: number): void {
- this.getConnectorStatus(connectorId).authorized = false;
+ private initializeConnectorStatus(connectorId: number): void {
+ this.getConnectorStatus(connectorId).idTagLocalAuthorized = false;
+ this.getConnectorStatus(connectorId).idTagAuthorized = false;
+ this.getConnectorStatus(connectorId).transactionRemoteStarted = false;
this.getConnectorStatus(connectorId).transactionStarted = false;
this.getConnectorStatus(connectorId).energyActiveImportRegisterValue = 0;
this.getConnectorStatus(connectorId).transactionEnergyActiveImportRegisterValue = 0;