Skip to content

Commit 7d56c4b

Browse files
authored
Merge pull request #88 from scs/feature/vse-standardization
VSE RL-DSP standardization
2 parents 3f7c72d + 814d075 commit 7d56c4b

33 files changed

Lines changed: 828 additions & 295 deletions

README.md

Lines changed: 7 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -22,15 +22,15 @@ The following smart meters are supported (see [Wiki/Home](https://github.com/scs
2222
* Landis+Gyr E570: \
2323
Data pushed by smart meter over CII interface (wired M-Bus, HDLC, DLMS/COSEM)
2424
* Landis+Gyr E360: \
25-
Data pushed by smart meter over P1 interface (HDLC, DLMS/COSEM only, no DSMR)
25+
Data pushed by smart meter over P1 interface (P1, HDLC, DLMS/COSEM only, no DSMR)
2626
* Iskraemeco AM550: \
27-
Data pushed by smart meter over P1 interface (HDLC, DLMS/COSEM only, no DSMR)
27+
Data pushed by smart meter over P1 interface (P1, HDLC, DLMS/COSEM only, no DSMR)
2828
* Kamstrup OMNIPOWER with HAN-NVE: \
2929
Data pushed by smart meter over inserted [HAN-NVE module](https://www.kamstrup.com/en-en/electricity-solutions/meters-devices/modules/hannve) (wired M-Bus, HDLC, DLMS/COSEM)
3030
* Siemens TD-351x: \
3131
Data fetched over bidirectional IR interface (IEC 62056-21, Mode C, unencrypted)
3232

33-
Note: All smart meters integrated so far push binary data encoded with HDLC (IEC 62056-46) and DLMS/COSEM. Both unencrypted and encrypted DLMS messages are accepted by the software.
33+
Note: All smart meters integrated so far push binary data encoded with HDLC (IEC 62056-46) and DLMS/COSEM. Both unencrypted and encrypted DLMS messages are accepted by the software. The P1-DSMR ASCII-based protocol is not supported.
3434

3535
The following data sinks are implemented:
3636
* MQTT (v3.1.1):
@@ -43,6 +43,8 @@ The following data sinks are implemented:
4343
* Unauthenticated
4444
* Authenticated with username / password
4545
* Authenticated with client certificate
46+
* MQTT RL-DSP: \
47+
like MQTT sink but with topic- and payload-format specified in VSE RL-DSP CH2024 document
4648
* Logger to `stdout`
4749

4850
`smartmeter-datacollector` is fully configurable through a `.ini` configuration file. The [`smartmeter-datacollector-configurator`](https://github.com/scs/smartmeter-datacollector-configurator) web interface can help to create and modify the configuration.
@@ -122,12 +124,12 @@ pipx install smartmeter-datacollector-configurator
122124

123125
## Method 3: Debian package
124126

125-
`smartmeter-datacollector` is also available as a Debian (`.deb`) package from the [releases](https://github.com/scs/smartmeter-datacollector/releases) which installs the application bundled with its Python dependencies in a zipapp / `.pyz` file. The Debian package includes a systemd service file which enables `smartmeter-datacollector` to automatically start after booting the system.
127+
`smartmeter-datacollector` is also available as a Debian (`.deb`) package from [releases](https://github.com/scs/smartmeter-datacollector/releases) which installs the application bundled with its Python dependencies in a zipapp / `.pyz` file. The Debian package includes a systemd service file which enables `smartmeter-datacollector` to automatically start after booting the system.
126128

127129
### Requirements
128130

129131
* Distribution: Debian based (like Debian, Ubuntu, ..)
130-
* Release: bookworm/12 (trixie/13)
132+
* Release: >= bookworm/12
131133
* CPU architecture: independent
132134

133135
### Installation

smartmeter_datacollector/collector.py

Lines changed: 9 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -11,7 +11,7 @@
1111
from typing import List
1212

1313
from smartmeter_datacollector.sinks.data_sink import DataSink
14-
from smartmeter_datacollector.smartmeter.meter_data import MeterDataPoint
14+
from smartmeter_datacollector.smartmeter.meter_data import MeterDataBundle
1515

1616
LOGGER = logging.getLogger("collector")
1717

@@ -25,16 +25,15 @@ def register_sink(self, sink: DataSink) -> None:
2525
assert isinstance(sink, DataSink)
2626
self._data_sinks.append(sink)
2727

28-
def notify(self, reader_data_points: List[MeterDataPoint]) -> None:
29-
for point in reader_data_points:
30-
try:
31-
self._queue.put_nowait(point)
32-
except QueueFull:
33-
LOGGER.warning("Queue is full. Current data points are dropped.")
34-
return
28+
def notify(self, reader_data_bundle: MeterDataBundle) -> None:
29+
try:
30+
self._queue.put_nowait(reader_data_bundle)
31+
except QueueFull:
32+
LOGGER.warning("Queue is full. Current data points are dropped.")
33+
return
3534

3635
async def process_queue(self) -> None:
3736
while True:
38-
data_point: MeterDataPoint = await self._queue.get()
37+
data_bundle: MeterDataBundle = await self._queue.get()
3938
for sink in self._data_sinks:
40-
await sink.send(data_point)
39+
await sink.send(data_bundle)

smartmeter_datacollector/config.py

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -27,7 +27,7 @@ class InvalidConfigError(Exception):
2727
},
2828
'sink1': {
2929
'type': "mqtt",
30-
'host': "localhost",
30+
'host': "127.0.0.1",
3131
'port': 1883,
3232
'tls': False,
3333
'ca_file_path': "",

smartmeter_datacollector/factory.py

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -13,7 +13,7 @@
1313
from smartmeter_datacollector.config import InvalidConfigError
1414
from smartmeter_datacollector.sinks.data_sink import DataSink
1515
from smartmeter_datacollector.sinks.logger_sink import LoggerSink
16-
from smartmeter_datacollector.sinks.mqtt_sink import MqttConfig, MqttDataSink
16+
from smartmeter_datacollector.sinks.mqtt_sink import MqttConfig, MqttDataSink, MqttSinkRlDsp
1717
from smartmeter_datacollector.smartmeter.iskraam550 import IskraAM550
1818
from smartmeter_datacollector.smartmeter.kamstrup_han import KamstrupHAN
1919
from smartmeter_datacollector.smartmeter.lge360 import LGE360
@@ -91,6 +91,9 @@ def build_sinks(config: ConfigParser) -> List[DataSink]:
9191
elif sink_type == "mqtt":
9292
mqtt_config = MqttConfig.from_sink_config(sink_config)
9393
sinks.append(MqttDataSink(mqtt_config))
94+
elif sink_type == "mqttrldsp":
95+
mqtt_config = MqttConfig.from_sink_config(sink_config)
96+
sinks.append(MqttSinkRlDsp(mqtt_config))
9497
else:
9598
raise InvalidConfigError(f"'type' is invalid or missing: {sink_type}")
9699
return sinks

smartmeter_datacollector/sinks/data_sink.py

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -7,7 +7,7 @@
77
#
88
from abc import ABC, abstractmethod
99

10-
from smartmeter_datacollector.smartmeter.meter_data import MeterDataPoint
10+
from smartmeter_datacollector.smartmeter.meter_data import MeterDataBundle
1111

1212

1313
class DataSink(ABC):
@@ -20,5 +20,5 @@ async def stop(self) -> None:
2020
raise NotImplementedError()
2121

2222
@abstractmethod
23-
async def send(self, data_point: MeterDataPoint) -> None:
23+
async def send(self, data_bundle: MeterDataBundle) -> None:
2424
raise NotImplementedError()

smartmeter_datacollector/sinks/logger_sink.py

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -8,7 +8,7 @@
88
import logging
99

1010
from smartmeter_datacollector.sinks.data_sink import DataSink
11-
from smartmeter_datacollector.smartmeter.meter_data import MeterDataPoint
11+
from smartmeter_datacollector.smartmeter.meter_data import MeterDataBundle
1212

1313

1414
class LoggerSink(DataSink):
@@ -22,5 +22,5 @@ async def start(self) -> None:
2222
async def stop(self) -> None:
2323
pass
2424

25-
async def send(self, data_point: MeterDataPoint) -> None:
26-
self._logger.info(str(data_point))
25+
async def send(self, data_bundle: MeterDataBundle) -> None:
26+
self._logger.info(str(data_bundle))

smartmeter_datacollector/sinks/mqtt_sink.py

Lines changed: 36 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -16,7 +16,7 @@
1616
from aiomqtt import Client, MqttCodeError, MqttError
1717

1818
from smartmeter_datacollector.sinks.data_sink import DataSink
19-
from smartmeter_datacollector.smartmeter.meter_data import MeterDataPoint
19+
from smartmeter_datacollector.smartmeter.meter_data import MeterDataBundle, MeterDataPoint
2020

2121
LOGGER = logging.getLogger("sink")
2222

@@ -33,6 +33,7 @@ class MqttConfig:
3333
check_hostname: bool = True
3434
client_cert_path: Optional[str] = None
3535
client_key_path: Optional[str] = None
36+
topic_group: str = "building"
3637

3738
def with_tls(
3839
self, ca_cert_path: Optional[str] = None, check_hostname: bool = True
@@ -72,6 +73,9 @@ def from_sink_config(config: SectionProxy) -> "MqttConfig":
7273
client_key_path = config.get("client_key_path")
7374
if client_cert_path is not None and client_key_path is not None:
7475
mqtt_cfg.with_client_cert_auth(str(client_cert_path), str(client_key_path))
76+
topic_group = config.get("topic_group")
77+
if topic_group and topic_group.strip().isalnum():
78+
mqtt_cfg.topic_group = topic_group.strip()
7579
return mqtt_cfg
7680

7781

@@ -139,11 +143,11 @@ async def stop(self) -> None:
139143
self._client_task = None
140144
LOGGER.info("Disconnected from MQTT broker")
141145

142-
async def send(self, data_point: MeterDataPoint) -> None:
143-
topic = MqttDataSink.get_topic_name_for_datapoint(data_point)
144-
dp_json = self.data_point_to_mqtt_json(data_point)
145-
146-
await self._publish_with_retries(topic, dp_json, retries=self.RETRIES)
146+
async def send(self, data_bundle: MeterDataBundle) -> None:
147+
for data_point in data_bundle.data_points:
148+
topic = MqttDataSink.get_topic_name_for_datapoint(data_point, data_bundle.source)
149+
dp_json = self.data_point_to_mqtt_json(data_point, int(data_bundle.timestamp.timestamp()))
150+
await self._publish_with_retries(topic, dp_json, retries=self.RETRIES)
147151

148152
async def _connection_handler(self) -> None:
149153
while True:
@@ -180,14 +184,36 @@ async def _publish_with_retries(self, topic: str, payload: str, retries: int = 1
180184
topic, retries + 1)
181185

182186
@staticmethod
183-
def get_topic_name_for_datapoint(data_point: MeterDataPoint) -> str:
184-
return f"smartmeter/{data_point.source}/{data_point.type.identifier}"
187+
def get_topic_name_for_datapoint(data_point: MeterDataPoint, source: str) -> str:
188+
return f"smartmeter/{source}/{data_point.type.identifier}"
185189

186190
@staticmethod
187-
def data_point_to_mqtt_json(data_point: MeterDataPoint) -> str:
191+
def data_point_to_mqtt_json(data_point: MeterDataPoint, timestamp: int) -> str:
188192
return json.dumps(
189193
{
190194
"value": data_point.value,
191-
"timestamp": int(data_point.timestamp.timestamp()),
195+
"timestamp": timestamp,
196+
"obis": data_point.obis.to_short_str(),
192197
}
193198
)
199+
200+
201+
class MqttSinkRlDsp(MqttDataSink):
202+
def __init__(self, config: MqttConfig) -> None:
203+
super().__init__(config)
204+
205+
self._group = config.topic_group
206+
207+
async def send(self, data_bundle: MeterDataBundle) -> None:
208+
topic = self.build_topic_name(data_bundle)
209+
payload = self.to_mqtt_payload(data_bundle)
210+
await self._publish_with_retries(topic, payload, retries=self.RETRIES)
211+
212+
def build_topic_name(self, data_bundle: MeterDataBundle) -> str:
213+
return f"dt/{self._group}/{data_bundle.source}/ds"
214+
215+
@staticmethod
216+
def to_mqtt_payload(data_bundle: MeterDataBundle) -> str:
217+
meter: dict[str, str | float] = {"ts": data_bundle.timestamp.isoformat()}
218+
meter.update({dp.obis.to_short_str(): dp.value for dp in data_bundle.data_points})
219+
return json.dumps({"meter": meter})

smartmeter_datacollector/smartmeter/cosem.py

Lines changed: 3 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -27,7 +27,7 @@ class RegisterCosem:
2727
scaling: float = 1.0
2828

2929

30-
DEFAULT_REGISTER_MAPPING = [
30+
DEFAULT_REGISTER_MAP = [
3131
RegisterCosem(OBISCode(1, 0, 1, 7, 0), MeterDataPointTypes.ACTIVE_POWER_P.value),
3232
RegisterCosem(OBISCode(1, 0, 2, 7, 0), MeterDataPointTypes.ACTIVE_POWER_N.value),
3333
RegisterCosem(OBISCode(1, 0, 3, 7, 0), MeterDataPointTypes.REACTIVE_POWER_P.value),
@@ -118,10 +118,9 @@ def __init__(self,
118118
LOGGER.warning("Empty fallback ID. Setting to random UUID %s.", fallback_id)
119119
self._fallback_id = fallback_id
120120
self._id_obis_override = id_obis_override if id_obis_override else []
121-
registers = DEFAULT_REGISTER_MAPPING
121+
self._register_obis = {r.obis: r for r in DEFAULT_REGISTER_MAP}
122122
if register_obis_extended:
123-
registers += register_obis_extended
124-
self._register_obis = {r.obis: r for r in registers}
123+
self._register_obis.update({r.obis: r for r in register_obis_extended})
125124
self._id_detect_countdown = Cosem.OBJECT_DETECT_ATTEMPTS
126125

127126
def retrieve_id(self, dlms_objects: Dict[OBISCode, Any]) -> str:

smartmeter_datacollector/smartmeter/hdlc_dlms_parser.py

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -16,7 +16,7 @@
1616
from gurux_dlms.secure import GXDLMSSecureClient
1717

1818
from smartmeter_datacollector.smartmeter.cosem import Cosem
19-
from smartmeter_datacollector.smartmeter.meter_data import MeterDataPoint
19+
from smartmeter_datacollector.smartmeter.meter_data import MeterDataBundle, MeterDataPoint
2020
from smartmeter_datacollector.smartmeter.obis import OBISCode
2121

2222
LOGGER = logging.getLogger("smartmeter")
@@ -112,8 +112,8 @@ def parse_to_dlms_objects(self) -> List[GXDLMSObject]:
112112
return dlms_objects
113113

114114
def convert_dlms_bundle_to_reader_data(self, dlms_objects: List[GXDLMSObject],
115-
message_time: Optional[datetime] = None) -> List[MeterDataPoint]:
116-
obis_obj_pairs = {}
115+
message_time: Optional[datetime] = None) -> MeterDataBundle:
116+
obis_obj_pairs: dict[OBISCode, GXDLMSObject] = {}
117117
for obj in dlms_objects:
118118
try:
119119
obis = OBISCode.from_string(str(obj.logicalName))
@@ -158,8 +158,8 @@ def convert_dlms_bundle_to_reader_data(self, dlms_objects: List[GXDLMSObject],
158158
except (TypeError, ValueError, OverflowError):
159159
LOGGER.warning("Invalid register value '%s'. Skipping register.", str(raw_value))
160160
continue
161-
data_points.append(MeterDataPoint(data_point_type, value, meter_id, timestamp))
162-
return data_points
161+
data_points.append(MeterDataPoint(data_point_type, value, obis))
162+
return MeterDataBundle(meter_id, timestamp, data_points)
163163

164164
def _parse_dlms_with_push_object_list(self) -> List[GXDLMSObject]:
165165
parsed_objects: List[Tuple[GXDLMSObject, int]] = []

smartmeter_datacollector/smartmeter/iskraam550.py

Lines changed: 9 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -10,8 +10,10 @@
1010

1111
import serial
1212

13-
from smartmeter_datacollector.smartmeter.cosem import Cosem
13+
from smartmeter_datacollector.smartmeter.cosem import Cosem, RegisterCosem
1414
from smartmeter_datacollector.smartmeter.meter import MeterError, SerialHdlcDlmsMeter
15+
from smartmeter_datacollector.smartmeter.meter_data import MeterDataPointTypes
16+
from smartmeter_datacollector.smartmeter.obis import OBISCode
1517
from smartmeter_datacollector.smartmeter.reader import ReaderError
1618
from smartmeter_datacollector.smartmeter.serial_reader import SerialConfig
1719

@@ -33,7 +35,12 @@ def __init__(self, port: str,
3335
stop_bits=serial.STOPBITS_ONE,
3436
termination=SerialHdlcDlmsMeter.HDLC_FLAG
3537
)
36-
cosem = Cosem(fallback_id=port)
38+
voltage_registers = [
39+
RegisterCosem(OBISCode(1, 0, 32, 7, 0), MeterDataPointTypes.VOLTAGE_L1.value, 0.1),
40+
RegisterCosem(OBISCode(1, 0, 52, 7, 0), MeterDataPointTypes.VOLTAGE_L2.value, 0.1),
41+
RegisterCosem(OBISCode(1, 0, 72, 7, 0), MeterDataPointTypes.VOLTAGE_L3.value, 0.1),
42+
]
43+
cosem = Cosem(fallback_id=port, register_obis_extended=voltage_registers)
3744
try:
3845
super().__init__(serial_config, cosem, decryption_key, use_system_time)
3946
except ReaderError as ex:

0 commit comments

Comments
 (0)