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