openrvdas.logger.readers.modbus_reader

No module-level documentation available.
  1#!/usr/bin/env python3
  2
  3import logging
  4import time
  5import yaml
  6from threading import Lock
  7
  8# Optional dependency
  9try:
 10    from pyModbusTCP.client import ModbusClient
 11    MODBUS_MODULE_FOUND = True
 12except ModuleNotFoundError:
 13    MODBUS_MODULE_FOUND = False
 14
 15from logger.readers.reader import Reader  # noqa
 16
 17
 18################################################################################
 19class ModBusTCPReader(Reader):
 20    """
 21    Read data from Modbus TCP devices.
 22
 23    Supports:
 24      - single-slave (legacy) mode
 25      - multi-slave / multi-function scan files
 26      - holding/input registers, coils, discrete inputs
 27      - text or raw-bytes output
 28      - partial failure tolerance
 29      - non-fatal connection failures
 30    """
 31
 32    _FUNCTION_MAP = {
 33        "holding_registers": "read_holding_registers",
 34        "input_registers": "read_input_registers",
 35        "coils": "read_coils",
 36        "discrete_inputs": "read_discrete_inputs",
 37    }
 38
 39    ############################
 40    def __init__(self,
 41                 registers=None,
 42                 host="localhost",
 43                 port=502,
 44                 scan_file=None,
 45                 slave=1,
 46                 function="holding_registers",
 47                 interval=10,
 48                 sep=" ",
 49                 encoding="utf-8",
 50                 encoding_errors="ignore",
 51                 **kwargs):
 52
 53        super().__init__(encoding=encoding, encoding_errors=encoding_errors, **kwargs)
 54
 55        if not MODBUS_MODULE_FOUND:
 56            raise RuntimeError(
 57                "Modbus functionality not available. Please install Python module pyModbusTCP."
 58            )
 59
 60        self.sep = sep if encoding is not None else None
 61        self.interval = interval
 62        self._lock = Lock()
 63        self._next_read_time = 0.0
 64
 65        self.client = ModbusClient(
 66            host=host,
 67            port=port,
 68            auto_open=True,
 69            auto_close=True,
 70        )
 71
 72        self.polls = []
 73
 74        if scan_file:
 75            self._load_scan_file(scan_file)
 76        else:
 77            # Legacy single-poll behavior
 78            if registers is None:
 79                raise ValueError("registers required when scan_file not provided")
 80
 81            self.polls.append({
 82                "slave": slave,
 83                "function": function,
 84                "registers": self._parse_registers(registers)
 85            })
 86
 87    ############################
 88    def _parse_register_spec(self, spec):
 89        """Parse a single register spec into (start, count) tuple."""
 90        if ":" in spec:
 91            start, end = spec.split(":", 1)
 92            if end == "":
 93                raise ValueError(f"Ambiguous register range: '{spec}'")
 94
 95            start = int(start) if start else 0
 96            end = int(end)
 97
 98            if start < 0 or end < start:
 99                raise ValueError(f"Invalid register range: '{spec}'")
100
101            return (start, end - start + 1)
102
103        addr = int(spec)
104        if addr < 0:
105            raise ValueError(f"Invalid register address: '{spec}'")
106
107        return (addr, 1)
108
109    ############################
110    def _validate_blocks(self, blocks):
111        validated = []
112        for block in blocks:
113            if not (
114                isinstance(block, tuple)
115                and len(block) == 2
116                and all(isinstance(x, int) and x >= 0 for x in block)
117            ):
118                raise ValueError(f"Invalid register block: {block}")
119            validated.append(block)
120        return validated
121
122    ############################
123    def _parse_registers(self, registers):
124        """Parse register specification into [(start, count), ...] blocks."""
125        if isinstance(registers, list):
126            return self._validate_blocks(registers)
127
128        if not isinstance(registers, str):
129            raise TypeError("registers must be a string or list")
130
131        blocks = []
132        for spec in registers.split(","):
133            spec = spec.strip()
134            if spec:
135                blocks.append(self._parse_register_spec(spec))
136
137        if not blocks:
138            raise ValueError("No valid registers specified")
139
140        return blocks
141
142    ############################
143    def _load_scan_file(self, path):
144        with open(path, "r") as f:
145            cfg = yaml.safe_load(f)
146
147        if "polls" not in cfg:
148            raise ValueError("scan_file must contain 'polls' list")
149
150        for poll in cfg["polls"]:
151            self.polls.append({
152                "slave": poll["slave"],
153                "function": poll.get("function", "holding_registers"),
154                "registers": self._parse_registers(poll["registers"]),
155            })
156
157    ############################
158    def _format_record(self, readings):
159        """
160        Flatten nested register/coil values into a single record.
161
162        - encoding=None:
163            - registers -> 16-bit unsigned big-endian
164            - coils -> packed bits
165        - otherwise:
166            - text with sep, 'nan' placeholders
167        """
168        if readings is None:
169            return None
170
171        flat_values = [val for block in readings if block for val in block]
172
173        if self.encoding is None:
174            if not flat_values:
175                return b""
176
177            first_val = next((v for v in flat_values if v is not None), 0)
178
179            # Coils / discrete inputs
180            if isinstance(first_val, (bool, int)) and all(
181                v in (0, 1, True, False, None) for v in flat_values
182            ):
183                byte_vals = []
184                current = 0
185                for i, bit in enumerate(flat_values):
186                    current |= (1 if bit else 0) << (i % 8)
187                    if i % 8 == 7:
188                        byte_vals.append(current)
189                        current = 0
190                if len(flat_values) % 8:
191                    byte_vals.append(current)
192                return bytes(byte_vals)
193
194            # Registers
195            return b"".join(
196                int(val or 0).to_bytes(2, byteorder="big", signed=False)
197                for val in flat_values
198            )
199
200        # Text mode
201        text_values = [str(val) if val is not None else "nan" for val in flat_values]
202        return self._encode_str(self.sep.join(text_values))
203
204    ############################
205    def read(self):
206        """
207        Read all configured polls.
208
209        Returns:
210          - list[str|bytes|None] for each poll
211          - None in place of any poll that cannot be read
212        """
213        with self._lock:
214            now = time.monotonic()
215            if now < self._next_read_time:
216                time.sleep(self._next_read_time - now)
217
218            start_time = time.monotonic()
219            results = []
220
221            for poll in self.polls:
222                slave = poll["slave"]
223                func_name = self._FUNCTION_MAP.get(poll["function"])
224
225                if not func_name:
226                    logging.warning(
227                        "Invalid function for slave=%s: %s",
228                        slave, poll["function"]
229                    )
230                    total = sum(c for _, c in poll["registers"])
231                    results.append([[None] * total])
232                    continue
233
234                func = getattr(self.client, func_name)
235                poll_readings = []
236
237                try:
238                    for start, count in poll["registers"]:
239                        block = func(start, count, unit_id=slave)
240                        if block is None:
241                            logging.warning(
242                                "Failed to read slave=%d addr=%d:%d",
243                                slave, start, count
244                            )
245                            poll_readings.append([None] * count)
246                        else:
247                            poll_readings.append(block)
248
249                    results.append(poll_readings)
250
251                except (OSError, AttributeError, ValueError, ConnectionError) as exc:
252                    logging.warning(
253                        "Modbus TCP connection error for slave=%d: %s", slave, exc
254                    )
255                    results.append(None)  # Entire poll failed, not fatal
256
257            self._next_read_time = start_time + self.interval
258
259            # Format records, preserving None for failed polls
260            formatted = [self._format_record(r) if r is not None else None for r in results]
261
262            # Prepend slave info for text mode
263            if self.encoding is not None:
264                formatted = [
265                    self._encode_str(f"slave {poll['slave']}:{self.sep}{line.decode(self.encoding) if isinstance(line, bytes) else line}")
266                    if line is not None else None
267                    for poll, line in zip(self.polls, formatted)
268                ]
269
270            return formatted
class ModBusTCPReader(logger.readers.reader.Reader):
 20class ModBusTCPReader(Reader):
 21    """
 22    Read data from Modbus TCP devices.
 23
 24    Supports:
 25      - single-slave (legacy) mode
 26      - multi-slave / multi-function scan files
 27      - holding/input registers, coils, discrete inputs
 28      - text or raw-bytes output
 29      - partial failure tolerance
 30      - non-fatal connection failures
 31    """
 32
 33    _FUNCTION_MAP = {
 34        "holding_registers": "read_holding_registers",
 35        "input_registers": "read_input_registers",
 36        "coils": "read_coils",
 37        "discrete_inputs": "read_discrete_inputs",
 38    }
 39
 40    ############################
 41    def __init__(self,
 42                 registers=None,
 43                 host="localhost",
 44                 port=502,
 45                 scan_file=None,
 46                 slave=1,
 47                 function="holding_registers",
 48                 interval=10,
 49                 sep=" ",
 50                 encoding="utf-8",
 51                 encoding_errors="ignore",
 52                 **kwargs):
 53
 54        super().__init__(encoding=encoding, encoding_errors=encoding_errors, **kwargs)
 55
 56        if not MODBUS_MODULE_FOUND:
 57            raise RuntimeError(
 58                "Modbus functionality not available. Please install Python module pyModbusTCP."
 59            )
 60
 61        self.sep = sep if encoding is not None else None
 62        self.interval = interval
 63        self._lock = Lock()
 64        self._next_read_time = 0.0
 65
 66        self.client = ModbusClient(
 67            host=host,
 68            port=port,
 69            auto_open=True,
 70            auto_close=True,
 71        )
 72
 73        self.polls = []
 74
 75        if scan_file:
 76            self._load_scan_file(scan_file)
 77        else:
 78            # Legacy single-poll behavior
 79            if registers is None:
 80                raise ValueError("registers required when scan_file not provided")
 81
 82            self.polls.append({
 83                "slave": slave,
 84                "function": function,
 85                "registers": self._parse_registers(registers)
 86            })
 87
 88    ############################
 89    def _parse_register_spec(self, spec):
 90        """Parse a single register spec into (start, count) tuple."""
 91        if ":" in spec:
 92            start, end = spec.split(":", 1)
 93            if end == "":
 94                raise ValueError(f"Ambiguous register range: '{spec}'")
 95
 96            start = int(start) if start else 0
 97            end = int(end)
 98
 99            if start < 0 or end < start:
100                raise ValueError(f"Invalid register range: '{spec}'")
101
102            return (start, end - start + 1)
103
104        addr = int(spec)
105        if addr < 0:
106            raise ValueError(f"Invalid register address: '{spec}'")
107
108        return (addr, 1)
109
110    ############################
111    def _validate_blocks(self, blocks):
112        validated = []
113        for block in blocks:
114            if not (
115                isinstance(block, tuple)
116                and len(block) == 2
117                and all(isinstance(x, int) and x >= 0 for x in block)
118            ):
119                raise ValueError(f"Invalid register block: {block}")
120            validated.append(block)
121        return validated
122
123    ############################
124    def _parse_registers(self, registers):
125        """Parse register specification into [(start, count), ...] blocks."""
126        if isinstance(registers, list):
127            return self._validate_blocks(registers)
128
129        if not isinstance(registers, str):
130            raise TypeError("registers must be a string or list")
131
132        blocks = []
133        for spec in registers.split(","):
134            spec = spec.strip()
135            if spec:
136                blocks.append(self._parse_register_spec(spec))
137
138        if not blocks:
139            raise ValueError("No valid registers specified")
140
141        return blocks
142
143    ############################
144    def _load_scan_file(self, path):
145        with open(path, "r") as f:
146            cfg = yaml.safe_load(f)
147
148        if "polls" not in cfg:
149            raise ValueError("scan_file must contain 'polls' list")
150
151        for poll in cfg["polls"]:
152            self.polls.append({
153                "slave": poll["slave"],
154                "function": poll.get("function", "holding_registers"),
155                "registers": self._parse_registers(poll["registers"]),
156            })
157
158    ############################
159    def _format_record(self, readings):
160        """
161        Flatten nested register/coil values into a single record.
162
163        - encoding=None:
164            - registers -> 16-bit unsigned big-endian
165            - coils -> packed bits
166        - otherwise:
167            - text with sep, 'nan' placeholders
168        """
169        if readings is None:
170            return None
171
172        flat_values = [val for block in readings if block for val in block]
173
174        if self.encoding is None:
175            if not flat_values:
176                return b""
177
178            first_val = next((v for v in flat_values if v is not None), 0)
179
180            # Coils / discrete inputs
181            if isinstance(first_val, (bool, int)) and all(
182                v in (0, 1, True, False, None) for v in flat_values
183            ):
184                byte_vals = []
185                current = 0
186                for i, bit in enumerate(flat_values):
187                    current |= (1 if bit else 0) << (i % 8)
188                    if i % 8 == 7:
189                        byte_vals.append(current)
190                        current = 0
191                if len(flat_values) % 8:
192                    byte_vals.append(current)
193                return bytes(byte_vals)
194
195            # Registers
196            return b"".join(
197                int(val or 0).to_bytes(2, byteorder="big", signed=False)
198                for val in flat_values
199            )
200
201        # Text mode
202        text_values = [str(val) if val is not None else "nan" for val in flat_values]
203        return self._encode_str(self.sep.join(text_values))
204
205    ############################
206    def read(self):
207        """
208        Read all configured polls.
209
210        Returns:
211          - list[str|bytes|None] for each poll
212          - None in place of any poll that cannot be read
213        """
214        with self._lock:
215            now = time.monotonic()
216            if now < self._next_read_time:
217                time.sleep(self._next_read_time - now)
218
219            start_time = time.monotonic()
220            results = []
221
222            for poll in self.polls:
223                slave = poll["slave"]
224                func_name = self._FUNCTION_MAP.get(poll["function"])
225
226                if not func_name:
227                    logging.warning(
228                        "Invalid function for slave=%s: %s",
229                        slave, poll["function"]
230                    )
231                    total = sum(c for _, c in poll["registers"])
232                    results.append([[None] * total])
233                    continue
234
235                func = getattr(self.client, func_name)
236                poll_readings = []
237
238                try:
239                    for start, count in poll["registers"]:
240                        block = func(start, count, unit_id=slave)
241                        if block is None:
242                            logging.warning(
243                                "Failed to read slave=%d addr=%d:%d",
244                                slave, start, count
245                            )
246                            poll_readings.append([None] * count)
247                        else:
248                            poll_readings.append(block)
249
250                    results.append(poll_readings)
251
252                except (OSError, AttributeError, ValueError, ConnectionError) as exc:
253                    logging.warning(
254                        "Modbus TCP connection error for slave=%d: %s", slave, exc
255                    )
256                    results.append(None)  # Entire poll failed, not fatal
257
258            self._next_read_time = start_time + self.interval
259
260            # Format records, preserving None for failed polls
261            formatted = [self._format_record(r) if r is not None else None for r in results]
262
263            # Prepend slave info for text mode
264            if self.encoding is not None:
265                formatted = [
266                    self._encode_str(f"slave {poll['slave']}:{self.sep}{line.decode(self.encoding) if isinstance(line, bytes) else line}")
267                    if line is not None else None
268                    for poll, line in zip(self.polls, formatted)
269                ]
270
271            return formatted

Read data from Modbus TCP devices.

Supports:

  • single-slave (legacy) mode
  • multi-slave / multi-function scan files
  • holding/input registers, coils, discrete inputs
  • text or raw-bytes output
  • partial failure tolerance
  • non-fatal connection failures
ModBusTCPReader( registers=None, host='localhost', port=502, scan_file=None, slave=1, function='holding_registers', interval=10, sep=' ', encoding='utf-8', encoding_errors='ignore', **kwargs)
41    def __init__(self,
42                 registers=None,
43                 host="localhost",
44                 port=502,
45                 scan_file=None,
46                 slave=1,
47                 function="holding_registers",
48                 interval=10,
49                 sep=" ",
50                 encoding="utf-8",
51                 encoding_errors="ignore",
52                 **kwargs):
53
54        super().__init__(encoding=encoding, encoding_errors=encoding_errors, **kwargs)
55
56        if not MODBUS_MODULE_FOUND:
57            raise RuntimeError(
58                "Modbus functionality not available. Please install Python module pyModbusTCP."
59            )
60
61        self.sep = sep if encoding is not None else None
62        self.interval = interval
63        self._lock = Lock()
64        self._next_read_time = 0.0
65
66        self.client = ModbusClient(
67            host=host,
68            port=port,
69            auto_open=True,
70            auto_close=True,
71        )
72
73        self.polls = []
74
75        if scan_file:
76            self._load_scan_file(scan_file)
77        else:
78            # Legacy single-poll behavior
79            if registers is None:
80                raise ValueError("registers required when scan_file not provided")
81
82            self.polls.append({
83                "slave": slave,
84                "function": function,
85                "registers": self._parse_registers(registers)
86            })
quiet - if type checking should log type errors or operate silently.

Two additional arguments govern how records will be encoded/decoded
from bytes, if desired by the Writer subclass when it calls
_encode_str() or _decode_bytes:

encoding - 'utf-8' by default. If empty or None, do not attempt any
        decoding and return raw bytes. Other possible encodings are
        listed in online documentation here:
        https://docs.python.org/3/library/codecs.html#standard-encodings

encoding_errors - 'ignore' by default. Other error strategies are
        'strict', 'replace', and 'backslashreplace', described here:
        https://docs.python.org/3/howto/unicode.html#encodings

mirror_to - Optional Writer to which all records read or transformed
        by this module (if it is a Reader or Transform) will be
        "mirrored" (copied). Mirroring happens asynchronously via
        a queue and background thread to minimize impact on the
        primary data flow. Writers cannot be mirrored.
sep
interval
client
polls
def read(self):
206    def read(self):
207        """
208        Read all configured polls.
209
210        Returns:
211          - list[str|bytes|None] for each poll
212          - None in place of any poll that cannot be read
213        """
214        with self._lock:
215            now = time.monotonic()
216            if now < self._next_read_time:
217                time.sleep(self._next_read_time - now)
218
219            start_time = time.monotonic()
220            results = []
221
222            for poll in self.polls:
223                slave = poll["slave"]
224                func_name = self._FUNCTION_MAP.get(poll["function"])
225
226                if not func_name:
227                    logging.warning(
228                        "Invalid function for slave=%s: %s",
229                        slave, poll["function"]
230                    )
231                    total = sum(c for _, c in poll["registers"])
232                    results.append([[None] * total])
233                    continue
234
235                func = getattr(self.client, func_name)
236                poll_readings = []
237
238                try:
239                    for start, count in poll["registers"]:
240                        block = func(start, count, unit_id=slave)
241                        if block is None:
242                            logging.warning(
243                                "Failed to read slave=%d addr=%d:%d",
244                                slave, start, count
245                            )
246                            poll_readings.append([None] * count)
247                        else:
248                            poll_readings.append(block)
249
250                    results.append(poll_readings)
251
252                except (OSError, AttributeError, ValueError, ConnectionError) as exc:
253                    logging.warning(
254                        "Modbus TCP connection error for slave=%d: %s", slave, exc
255                    )
256                    results.append(None)  # Entire poll failed, not fatal
257
258            self._next_read_time = start_time + self.interval
259
260            # Format records, preserving None for failed polls
261            formatted = [self._format_record(r) if r is not None else None for r in results]
262
263            # Prepend slave info for text mode
264            if self.encoding is not None:
265                formatted = [
266                    self._encode_str(f"slave {poll['slave']}:{self.sep}{line.decode(self.encoding) if isinstance(line, bytes) else line}")
267                    if line is not None else None
268                    for poll, line in zip(self.polls, formatted)
269                ]
270
271            return formatted

Read all configured polls.

Returns:

  • list[str|bytes|None] for each poll
  • None in place of any poll that cannot be read