openrvdas.logger.readers.modbus_serial_reader
'''
ModBusSerialReader
This module provides a thread-safe reader for polling Modbus RTU devices over a serial connection. It supports single-slave and multi-slave setups, multiple Modbus function types, and flexible output formats (binary or text).
Key Features
- Reads holding registers, input registers, coils, and discrete inputs.
- Supports configuration via direct register specification or YAML scan file.
- Partial failure tolerance: failed polls return placeholder values without stopping other reads.
- Exponential backoff for serial connection retries (1s → 30s max).
- Thread-safe access via internal Lock.
- Output can be raw bytes (binary) or encoded text with customizable separator.
- Automatic cleanup of the serial port on
stop()or object deletion.
Typical Usage
reader = ModBusSerialReader(
registers="0:5,10:15",
port="/dev/ttyUSB0",
baudrate=19200,
interval=5,
encoding="utf-8"
)
data = reader.read()
print(data)
reader.stop()
Dependencies
- pymodbus (optional; raises RuntimeError if missing)
- PyYAML
- Standard Python libraries: logging, time, threading, sys, os.path '''
1#!/usr/bin/env python3 2""" 3''' 4ModBusSerialReader 5================== 6 7This module provides a thread-safe reader for polling Modbus RTU devices 8over a serial connection. It supports single-slave and multi-slave setups, 9multiple Modbus function types, and flexible output formats (binary or text). 10 11Key Features 12------------ 13- Reads holding registers, input registers, coils, and discrete inputs. 14- Supports configuration via direct register specification or YAML scan file. 15- Partial failure tolerance: failed polls return placeholder values without 16 stopping other reads. 17- Exponential backoff for serial connection retries (1s → 30s max). 18- Thread-safe access via internal Lock. 19- Output can be raw bytes (binary) or encoded text with customizable separator. 20- Automatic cleanup of the serial port on `stop()` or object deletion. 21 22Typical Usage 23------------- 24 reader = ModBusSerialReader( 25 registers="0:5,10:15", 26 port="/dev/ttyUSB0", 27 baudrate=19200, 28 interval=5, 29 encoding="utf-8" 30 ) 31 data = reader.read() 32 print(data) 33 reader.stop() 34 35Dependencies 36------------ 37- pymodbus (optional; raises RuntimeError if missing) 38- PyYAML 39- Standard Python libraries: logging, time, threading, sys, os.path 40''' 41""" 42 43import logging 44import yaml 45import time 46from threading import Lock 47 48try: 49 from pymodbus.client import ModbusSerialClient 50 from pymodbus import ModbusException 51 MODBUS_MODULE_FOUND = True 52except ModuleNotFoundError: 53 MODBUS_MODULE_FOUND = False 54 55from logger.readers.reader import Reader # noqa 56 57 58############################################################################### 59class ModBusSerialReader(Reader): 60 """ 61 Read data from Modbus RTU devices over a serial connection. 62 """ 63 64 ############################ 65 def __init__(self, 66 registers=None, 67 port="/dev/ttyUSB0", 68 baudrate=9600, 69 parity="N", 70 stopbits=1, 71 bytesize=8, 72 scan_file=None, 73 slave=1, 74 function='holding_registers', 75 interval=1.0, 76 sep=" ", 77 encoding="utf-8", 78 encoding_errors="ignore", 79 timeout=None, 80 **kwargs): 81 """ 82 ``` 83 registers - Comma-separated string (e.g. '0:5,10:15') or list of 84 tuples [(start,count), ...] specifying which registers 85 to poll. Required if scan_file is not provided. 86 87 port - Serial port device (e.g. '/dev/ttyUSB0'). 88 89 baudrate - Serial connection baud rate (default 9600). 90 91 parity - Serial parity, one of 'N', 'E', 'O' (default 'N'). 92 93 stopbits - Number of stop bits (default 1). 94 95 bytesize - Number of data bits (default 8). 96 97 scan_file - Optional YAML file specifying multiple slaves, functions, 98 and register blocks. Overrides registers/slave/function 99 parameters if provided. 100 101 slave - Modbus slave ID (default 1). Ignored if scan_file is provided. 102 103 function - Modbus function type to read: 'holding_registers', 104 'input_registers', 'coils', or 'discrete_inputs'. 105 Ignored if scan_file is provided. 106 107 interval - Seconds between consecutive reads. Must be >= 0.1. 108 109 sep - Separator string used when encoding output as text. Ignored 110 if encoding is None. 111 112 encoding - Character encoding for text output (default 'utf-8'). 113 If None, raw bytes are returned. 114 115 encoding_errors - Strategy for handling encoding errors ('ignore' 116 by default). Other options: 'strict', 'replace', 117 'backslashreplace'. 118 119 timeout - Max time in seconds to wait for serial response. Defaults 120 to 2 seconds if None. Minimum timeout is 1s. 121 ``` 122 """ 123 super().__init__(encoding=encoding, encoding_errors=encoding_errors, **kwargs) 124 125 if not MODBUS_MODULE_FOUND: 126 raise RuntimeError( 127 "Modbus functionality not available. Install pymodbus: " 128 'pip install "pymodbus[serial]"' 129 ) 130 131 self._FUNCTION_MAP = { 132 "holding_registers": "read_holding_registers", 133 "input_registers": "read_input_registers", 134 "coils": "read_coils", 135 "discrete_inputs": "read_discrete_inputs", 136 } 137 138 if interval < 0.1: 139 raise ValueError('Interval must be greater or equal to 0.1 seconds (10Hz)') 140 141 timeout = timeout or 2.0 142 if timeout < 1.0: 143 raise ValueError('Timeout must be greater or equal to 1 seconds') 144 145 self.polls = [] 146 if scan_file: 147 self._load_scan_file(scan_file) 148 else: 149 if registers is None: 150 raise ValueError("registers required when scan_file not provided") 151 self.polls = [{ 152 "slave": slave, 153 "function": function, 154 "registers": self._parse_registers(registers) 155 }] 156 157 self.sep = sep if encoding is not None else None 158 self.interval = interval 159 self._next_read_time = 0.0 160 self._connected = False 161 self._reconnect_delay = timeout 162 self._reconnect_delay_max = 30.0 163 self._next_connect_time = 0.0 164 self._lock = Lock() 165 166 self.client = ModbusSerialClient( 167 port=port, 168 baudrate=baudrate, 169 parity=parity, 170 stopbits=stopbits, 171 bytesize=bytesize, 172 timeout=timeout, 173 ) 174 175 ############################ 176 def _reset_client(self): 177 """Hard-reset of the serial client.""" 178 try: 179 self.client.close() 180 except Exception: 181 pass 182 183 self._connected = False 184 185 now = time.monotonic() 186 self._next_connect_time = now + self._reconnect_delay 187 self._reconnect_delay = min( 188 self._reconnect_delay * 2, 189 self._reconnect_delay_max 190 ) 191 logging.info( 192 "Next Modbus reconnect attempt in %.1f seconds", 193 self._reconnect_delay 194 ) 195 196 ############################ 197 def _parse_register_spec(self, spec): 198 """Parse a single register spec into (start, count).""" 199 if ":" in spec: 200 start, end = spec.split(":", 1) 201 if end == "": 202 raise ValueError(f"Ambiguous register range: '{spec}'") 203 start = int(start) if start else 0 204 end = int(end) 205 if start < 0 or end < start: 206 raise ValueError(f"Invalid register range: '{spec}'") 207 return (start, end - start + 1) 208 addr = int(spec) 209 if addr < 0: 210 raise ValueError(f"Invalid register address: '{spec}'") 211 return (addr, 1) 212 213 ############################ 214 def _validate_blocks(self, blocks): 215 """Validate and individual register block""" 216 validated = [] 217 for block in blocks: 218 if not ( 219 isinstance(block, tuple) 220 and len(block) == 2 221 and all(isinstance(x, int) and x >= 0 for x in block) 222 ): 223 raise ValueError(f"Invalid register block: {block}") 224 validated.append(block) 225 return validated 226 227 ############################ 228 def _parse_registers(self, registers): 229 """Parse comma-separated string or list into [(start, count), ...] blocks.""" 230 if registers is None: 231 return [(0, 10)] 232 if isinstance(registers, list): 233 return self._validate_blocks(registers) 234 if not isinstance(registers, str): 235 raise TypeError("registers must be str, list, or None") 236 237 blocks = [] 238 for spec in registers.split(","): 239 spec = spec.strip() 240 if spec: 241 blocks.append(self._parse_register_spec(spec)) 242 if not blocks: 243 raise ValueError("No valid registers specified") 244 return blocks 245 246 ############################ 247 def _client_connected(self) -> bool: 248 """Connect if not already connected. Returns True if OK.""" 249 if self._connected: 250 return True 251 252 now = time.monotonic() 253 if now < self._next_connect_time: 254 return False 255 256 try: 257 self._connected = self.client.connect() 258 except ModbusException as exc: 259 logging.error(f"ModbusException: connect() failed: {exc}") 260 self._reset_client() 261 return False 262 except Exception as exc: 263 logging.error(f"Exception: connect() failed: {exc}") 264 self._reset_client() 265 return False 266 267 if not self._connected: 268 logging.error("ModBusSerialReader: unable to open serial port") 269 self._reset_client() 270 return False 271 272 self._reconnect_delay = 1.0 273 self._next_connect_time = 0.0 274 275 return self._connected 276 277 ############################ 278 def _load_scan_file(self, path): 279 """Load the YAML-formatted scan file""" 280 with open(path, "r") as f: 281 cfg = yaml.safe_load(f) 282 283 if "polls" not in cfg: 284 raise ValueError("scan_file must contain 'polls' list") 285 286 for poll in cfg["polls"]: 287 slave = poll["slave"] 288 function = poll.get("function", "holding_registers") 289 registers = poll["registers"] 290 291 self.polls.append({ 292 "slave": slave, 293 "function": function, 294 "registers": self._parse_registers(registers) 295 }) 296 297 ############################ 298 def _format_record(self, readings: list[list[int | bool | None]], function: str) -> bytes | str: 299 """ 300 Flatten nested register/coil values into a single record. 301 302 Parameters 303 ---------- 304 readings : list[list[int | bool | None]] 305 Nested lists of values returned from Modbus reads. 306 307 function : str 308 Modbus function type. Determines formatting behavior. 309 Expected values: 310 - "holding_registers" 311 - "input_registers" 312 - "coils" 313 - "discrete_inputs" 314 315 Behavior 316 -------- 317 - encoding is None: 318 - Registers → 16-bit unsigned big-endian 319 - Coils/discrete_inputs → packed bits (LSB-first) 320 - None values → 0 321 - Otherwise convert to text and encode with _encode_str, using self.sep 322 """ 323 324 flat_values: list[int | bool | None] = [] 325 for block in readings: 326 if block: 327 flat_values.extend(block) 328 329 is_coils = function in ("coils", "discrete_inputs") 330 331 # ------------------------------------------------------------------ 332 # Binary 333 # ------------------------------------------------------------------ 334 if self.encoding is None: 335 if not flat_values: 336 return b"" 337 338 if is_coils: 339 # Pack bits LSB-first, per Modbus spec 340 byte_vals = [] 341 current_byte = 0 342 bit_index = 0 343 344 for bit in flat_values: 345 bit_val = 1 if bit else 0 346 current_byte |= bit_val << bit_index 347 bit_index += 1 348 349 if bit_index == 8: 350 byte_vals.append(current_byte) 351 current_byte = 0 352 bit_index = 0 353 354 if bit_index: 355 byte_vals.append(current_byte) 356 357 return bytes(byte_vals) 358 359 # Registers: 16-bit unsigned big-endian 360 return b"".join( 361 int(val or 0).to_bytes(2, byteorder="big", signed=False) 362 for val in flat_values 363 ) 364 365 # ------------------------------------------------------------------ 366 # Encoded text 367 # ------------------------------------------------------------------ 368 text_values = [ 369 str(val) if val is not None else "nan" 370 for val in flat_values 371 ] 372 return self._encode_str(self.sep.join(text_values)) 373 374 375 ############################ 376 def read(self): 377 """Read all configured polls and return a list of formatted records per slave. 378 379 Returns: 380 list[bytes|str|None]: None for polls that cannot be read; otherwise formatted output. 381 """ 382 with self._lock: 383 now = time.monotonic() 384 if now < self._next_read_time: 385 time.sleep(self._next_read_time - now) 386 start_time = time.monotonic() 387 388 if not self._client_connected(): 389 return None 390 391 records = [] 392 393 for poll in self.polls: 394 slave = poll["slave"] 395 func_name = self._FUNCTION_MAP.get(poll["function"]) 396 397 if not func_name: 398 logging.warning("Invalid function for slave=%s: %s", slave, poll["function"]) 399 total_count = sum(count for _, count in poll["registers"]) 400 readings = [[None] * total_count] 401 else: 402 func = getattr(self.client, func_name) 403 readings = [] 404 for start, count in poll["registers"]: 405 try: 406 resp = func(address=start, count=count, unit=slave) 407 if resp is None or getattr(resp, "isError", lambda: False)(): 408 logging.warning("Failed to read slave=%d addr=%d:%d", slave, start, count) 409 readings.append([None] * count) 410 else: 411 values = getattr(resp, "registers", getattr(resp, "bits", None)) 412 readings.append(values or [None] * count) 413 except (ModbusException, OSError, ConnectionError, ValueError, AttributeError) as exc: 414 logging.warning("Modbus read error slave=%d addr=%d:%d: %s", slave, start, count, exc) 415 readings.append([None] * count) 416 self._reset_client() 417 break 418 419 record = self._format_record(readings, poll["function"]) 420 if self.encoding is not None and isinstance(record, bytes): 421 record = record.decode(self.encoding, errors=self.encoding_errors) 422 record = self._encode_str(f"slave {slave}:{self.sep}{record}") 423 records.append(record) 424 425 self._next_read_time = start_time + self.interval 426 return records 427 428 ############################ 429 def stop(self): 430 """Explicitly close the serial port.""" 431 with self._lock: 432 try: 433 self.client.close() 434 except Exception: 435 pass 436 self._connected = False 437 438 ############################ 439 def __del__(self): 440 """Destructor — ensures port is closed safely.""" 441 try: 442 self.stop() 443 except Exception: 444 pass
class
ModBusSerialReader(logger.readers.reader.Reader):
60class ModBusSerialReader(Reader): 61 """ 62 Read data from Modbus RTU devices over a serial connection. 63 """ 64 65 ############################ 66 def __init__(self, 67 registers=None, 68 port="/dev/ttyUSB0", 69 baudrate=9600, 70 parity="N", 71 stopbits=1, 72 bytesize=8, 73 scan_file=None, 74 slave=1, 75 function='holding_registers', 76 interval=1.0, 77 sep=" ", 78 encoding="utf-8", 79 encoding_errors="ignore", 80 timeout=None, 81 **kwargs): 82 """ 83 ``` 84 registers - Comma-separated string (e.g. '0:5,10:15') or list of 85 tuples [(start,count), ...] specifying which registers 86 to poll. Required if scan_file is not provided. 87 88 port - Serial port device (e.g. '/dev/ttyUSB0'). 89 90 baudrate - Serial connection baud rate (default 9600). 91 92 parity - Serial parity, one of 'N', 'E', 'O' (default 'N'). 93 94 stopbits - Number of stop bits (default 1). 95 96 bytesize - Number of data bits (default 8). 97 98 scan_file - Optional YAML file specifying multiple slaves, functions, 99 and register blocks. Overrides registers/slave/function 100 parameters if provided. 101 102 slave - Modbus slave ID (default 1). Ignored if scan_file is provided. 103 104 function - Modbus function type to read: 'holding_registers', 105 'input_registers', 'coils', or 'discrete_inputs'. 106 Ignored if scan_file is provided. 107 108 interval - Seconds between consecutive reads. Must be >= 0.1. 109 110 sep - Separator string used when encoding output as text. Ignored 111 if encoding is None. 112 113 encoding - Character encoding for text output (default 'utf-8'). 114 If None, raw bytes are returned. 115 116 encoding_errors - Strategy for handling encoding errors ('ignore' 117 by default). Other options: 'strict', 'replace', 118 'backslashreplace'. 119 120 timeout - Max time in seconds to wait for serial response. Defaults 121 to 2 seconds if None. Minimum timeout is 1s. 122 ``` 123 """ 124 super().__init__(encoding=encoding, encoding_errors=encoding_errors, **kwargs) 125 126 if not MODBUS_MODULE_FOUND: 127 raise RuntimeError( 128 "Modbus functionality not available. Install pymodbus: " 129 'pip install "pymodbus[serial]"' 130 ) 131 132 self._FUNCTION_MAP = { 133 "holding_registers": "read_holding_registers", 134 "input_registers": "read_input_registers", 135 "coils": "read_coils", 136 "discrete_inputs": "read_discrete_inputs", 137 } 138 139 if interval < 0.1: 140 raise ValueError('Interval must be greater or equal to 0.1 seconds (10Hz)') 141 142 timeout = timeout or 2.0 143 if timeout < 1.0: 144 raise ValueError('Timeout must be greater or equal to 1 seconds') 145 146 self.polls = [] 147 if scan_file: 148 self._load_scan_file(scan_file) 149 else: 150 if registers is None: 151 raise ValueError("registers required when scan_file not provided") 152 self.polls = [{ 153 "slave": slave, 154 "function": function, 155 "registers": self._parse_registers(registers) 156 }] 157 158 self.sep = sep if encoding is not None else None 159 self.interval = interval 160 self._next_read_time = 0.0 161 self._connected = False 162 self._reconnect_delay = timeout 163 self._reconnect_delay_max = 30.0 164 self._next_connect_time = 0.0 165 self._lock = Lock() 166 167 self.client = ModbusSerialClient( 168 port=port, 169 baudrate=baudrate, 170 parity=parity, 171 stopbits=stopbits, 172 bytesize=bytesize, 173 timeout=timeout, 174 ) 175 176 ############################ 177 def _reset_client(self): 178 """Hard-reset of the serial client.""" 179 try: 180 self.client.close() 181 except Exception: 182 pass 183 184 self._connected = False 185 186 now = time.monotonic() 187 self._next_connect_time = now + self._reconnect_delay 188 self._reconnect_delay = min( 189 self._reconnect_delay * 2, 190 self._reconnect_delay_max 191 ) 192 logging.info( 193 "Next Modbus reconnect attempt in %.1f seconds", 194 self._reconnect_delay 195 ) 196 197 ############################ 198 def _parse_register_spec(self, spec): 199 """Parse a single register spec into (start, count).""" 200 if ":" in spec: 201 start, end = spec.split(":", 1) 202 if end == "": 203 raise ValueError(f"Ambiguous register range: '{spec}'") 204 start = int(start) if start else 0 205 end = int(end) 206 if start < 0 or end < start: 207 raise ValueError(f"Invalid register range: '{spec}'") 208 return (start, end - start + 1) 209 addr = int(spec) 210 if addr < 0: 211 raise ValueError(f"Invalid register address: '{spec}'") 212 return (addr, 1) 213 214 ############################ 215 def _validate_blocks(self, blocks): 216 """Validate and individual register block""" 217 validated = [] 218 for block in blocks: 219 if not ( 220 isinstance(block, tuple) 221 and len(block) == 2 222 and all(isinstance(x, int) and x >= 0 for x in block) 223 ): 224 raise ValueError(f"Invalid register block: {block}") 225 validated.append(block) 226 return validated 227 228 ############################ 229 def _parse_registers(self, registers): 230 """Parse comma-separated string or list into [(start, count), ...] blocks.""" 231 if registers is None: 232 return [(0, 10)] 233 if isinstance(registers, list): 234 return self._validate_blocks(registers) 235 if not isinstance(registers, str): 236 raise TypeError("registers must be str, list, or None") 237 238 blocks = [] 239 for spec in registers.split(","): 240 spec = spec.strip() 241 if spec: 242 blocks.append(self._parse_register_spec(spec)) 243 if not blocks: 244 raise ValueError("No valid registers specified") 245 return blocks 246 247 ############################ 248 def _client_connected(self) -> bool: 249 """Connect if not already connected. Returns True if OK.""" 250 if self._connected: 251 return True 252 253 now = time.monotonic() 254 if now < self._next_connect_time: 255 return False 256 257 try: 258 self._connected = self.client.connect() 259 except ModbusException as exc: 260 logging.error(f"ModbusException: connect() failed: {exc}") 261 self._reset_client() 262 return False 263 except Exception as exc: 264 logging.error(f"Exception: connect() failed: {exc}") 265 self._reset_client() 266 return False 267 268 if not self._connected: 269 logging.error("ModBusSerialReader: unable to open serial port") 270 self._reset_client() 271 return False 272 273 self._reconnect_delay = 1.0 274 self._next_connect_time = 0.0 275 276 return self._connected 277 278 ############################ 279 def _load_scan_file(self, path): 280 """Load the YAML-formatted scan file""" 281 with open(path, "r") as f: 282 cfg = yaml.safe_load(f) 283 284 if "polls" not in cfg: 285 raise ValueError("scan_file must contain 'polls' list") 286 287 for poll in cfg["polls"]: 288 slave = poll["slave"] 289 function = poll.get("function", "holding_registers") 290 registers = poll["registers"] 291 292 self.polls.append({ 293 "slave": slave, 294 "function": function, 295 "registers": self._parse_registers(registers) 296 }) 297 298 ############################ 299 def _format_record(self, readings: list[list[int | bool | None]], function: str) -> bytes | str: 300 """ 301 Flatten nested register/coil values into a single record. 302 303 Parameters 304 ---------- 305 readings : list[list[int | bool | None]] 306 Nested lists of values returned from Modbus reads. 307 308 function : str 309 Modbus function type. Determines formatting behavior. 310 Expected values: 311 - "holding_registers" 312 - "input_registers" 313 - "coils" 314 - "discrete_inputs" 315 316 Behavior 317 -------- 318 - encoding is None: 319 - Registers → 16-bit unsigned big-endian 320 - Coils/discrete_inputs → packed bits (LSB-first) 321 - None values → 0 322 - Otherwise convert to text and encode with _encode_str, using self.sep 323 """ 324 325 flat_values: list[int | bool | None] = [] 326 for block in readings: 327 if block: 328 flat_values.extend(block) 329 330 is_coils = function in ("coils", "discrete_inputs") 331 332 # ------------------------------------------------------------------ 333 # Binary 334 # ------------------------------------------------------------------ 335 if self.encoding is None: 336 if not flat_values: 337 return b"" 338 339 if is_coils: 340 # Pack bits LSB-first, per Modbus spec 341 byte_vals = [] 342 current_byte = 0 343 bit_index = 0 344 345 for bit in flat_values: 346 bit_val = 1 if bit else 0 347 current_byte |= bit_val << bit_index 348 bit_index += 1 349 350 if bit_index == 8: 351 byte_vals.append(current_byte) 352 current_byte = 0 353 bit_index = 0 354 355 if bit_index: 356 byte_vals.append(current_byte) 357 358 return bytes(byte_vals) 359 360 # Registers: 16-bit unsigned big-endian 361 return b"".join( 362 int(val or 0).to_bytes(2, byteorder="big", signed=False) 363 for val in flat_values 364 ) 365 366 # ------------------------------------------------------------------ 367 # Encoded text 368 # ------------------------------------------------------------------ 369 text_values = [ 370 str(val) if val is not None else "nan" 371 for val in flat_values 372 ] 373 return self._encode_str(self.sep.join(text_values)) 374 375 376 ############################ 377 def read(self): 378 """Read all configured polls and return a list of formatted records per slave. 379 380 Returns: 381 list[bytes|str|None]: None for polls that cannot be read; otherwise formatted output. 382 """ 383 with self._lock: 384 now = time.monotonic() 385 if now < self._next_read_time: 386 time.sleep(self._next_read_time - now) 387 start_time = time.monotonic() 388 389 if not self._client_connected(): 390 return None 391 392 records = [] 393 394 for poll in self.polls: 395 slave = poll["slave"] 396 func_name = self._FUNCTION_MAP.get(poll["function"]) 397 398 if not func_name: 399 logging.warning("Invalid function for slave=%s: %s", slave, poll["function"]) 400 total_count = sum(count for _, count in poll["registers"]) 401 readings = [[None] * total_count] 402 else: 403 func = getattr(self.client, func_name) 404 readings = [] 405 for start, count in poll["registers"]: 406 try: 407 resp = func(address=start, count=count, unit=slave) 408 if resp is None or getattr(resp, "isError", lambda: False)(): 409 logging.warning("Failed to read slave=%d addr=%d:%d", slave, start, count) 410 readings.append([None] * count) 411 else: 412 values = getattr(resp, "registers", getattr(resp, "bits", None)) 413 readings.append(values or [None] * count) 414 except (ModbusException, OSError, ConnectionError, ValueError, AttributeError) as exc: 415 logging.warning("Modbus read error slave=%d addr=%d:%d: %s", slave, start, count, exc) 416 readings.append([None] * count) 417 self._reset_client() 418 break 419 420 record = self._format_record(readings, poll["function"]) 421 if self.encoding is not None and isinstance(record, bytes): 422 record = record.decode(self.encoding, errors=self.encoding_errors) 423 record = self._encode_str(f"slave {slave}:{self.sep}{record}") 424 records.append(record) 425 426 self._next_read_time = start_time + self.interval 427 return records 428 429 ############################ 430 def stop(self): 431 """Explicitly close the serial port.""" 432 with self._lock: 433 try: 434 self.client.close() 435 except Exception: 436 pass 437 self._connected = False 438 439 ############################ 440 def __del__(self): 441 """Destructor — ensures port is closed safely.""" 442 try: 443 self.stop() 444 except Exception: 445 pass
Read data from Modbus RTU devices over a serial connection.
ModBusSerialReader( registers=None, port='/dev/ttyUSB0', baudrate=9600, parity='N', stopbits=1, bytesize=8, scan_file=None, slave=1, function='holding_registers', interval=1.0, sep=' ', encoding='utf-8', encoding_errors='ignore', timeout=None, **kwargs)
66 def __init__(self, 67 registers=None, 68 port="/dev/ttyUSB0", 69 baudrate=9600, 70 parity="N", 71 stopbits=1, 72 bytesize=8, 73 scan_file=None, 74 slave=1, 75 function='holding_registers', 76 interval=1.0, 77 sep=" ", 78 encoding="utf-8", 79 encoding_errors="ignore", 80 timeout=None, 81 **kwargs): 82 """ 83 ``` 84 registers - Comma-separated string (e.g. '0:5,10:15') or list of 85 tuples [(start,count), ...] specifying which registers 86 to poll. Required if scan_file is not provided. 87 88 port - Serial port device (e.g. '/dev/ttyUSB0'). 89 90 baudrate - Serial connection baud rate (default 9600). 91 92 parity - Serial parity, one of 'N', 'E', 'O' (default 'N'). 93 94 stopbits - Number of stop bits (default 1). 95 96 bytesize - Number of data bits (default 8). 97 98 scan_file - Optional YAML file specifying multiple slaves, functions, 99 and register blocks. Overrides registers/slave/function 100 parameters if provided. 101 102 slave - Modbus slave ID (default 1). Ignored if scan_file is provided. 103 104 function - Modbus function type to read: 'holding_registers', 105 'input_registers', 'coils', or 'discrete_inputs'. 106 Ignored if scan_file is provided. 107 108 interval - Seconds between consecutive reads. Must be >= 0.1. 109 110 sep - Separator string used when encoding output as text. Ignored 111 if encoding is None. 112 113 encoding - Character encoding for text output (default 'utf-8'). 114 If None, raw bytes are returned. 115 116 encoding_errors - Strategy for handling encoding errors ('ignore' 117 by default). Other options: 'strict', 'replace', 118 'backslashreplace'. 119 120 timeout - Max time in seconds to wait for serial response. Defaults 121 to 2 seconds if None. Minimum timeout is 1s. 122 ``` 123 """ 124 super().__init__(encoding=encoding, encoding_errors=encoding_errors, **kwargs) 125 126 if not MODBUS_MODULE_FOUND: 127 raise RuntimeError( 128 "Modbus functionality not available. Install pymodbus: " 129 'pip install "pymodbus[serial]"' 130 ) 131 132 self._FUNCTION_MAP = { 133 "holding_registers": "read_holding_registers", 134 "input_registers": "read_input_registers", 135 "coils": "read_coils", 136 "discrete_inputs": "read_discrete_inputs", 137 } 138 139 if interval < 0.1: 140 raise ValueError('Interval must be greater or equal to 0.1 seconds (10Hz)') 141 142 timeout = timeout or 2.0 143 if timeout < 1.0: 144 raise ValueError('Timeout must be greater or equal to 1 seconds') 145 146 self.polls = [] 147 if scan_file: 148 self._load_scan_file(scan_file) 149 else: 150 if registers is None: 151 raise ValueError("registers required when scan_file not provided") 152 self.polls = [{ 153 "slave": slave, 154 "function": function, 155 "registers": self._parse_registers(registers) 156 }] 157 158 self.sep = sep if encoding is not None else None 159 self.interval = interval 160 self._next_read_time = 0.0 161 self._connected = False 162 self._reconnect_delay = timeout 163 self._reconnect_delay_max = 30.0 164 self._next_connect_time = 0.0 165 self._lock = Lock() 166 167 self.client = ModbusSerialClient( 168 port=port, 169 baudrate=baudrate, 170 parity=parity, 171 stopbits=stopbits, 172 bytesize=bytesize, 173 timeout=timeout, 174 )
registers - Comma-separated string (e.g. '0:5,10:15') or list of
tuples [(start,count), ...] specifying which registers
to poll. Required if scan_file is not provided.
port - Serial port device (e.g. '/dev/ttyUSB0').
baudrate - Serial connection baud rate (default 9600).
parity - Serial parity, one of 'N', 'E', 'O' (default 'N').
stopbits - Number of stop bits (default 1).
bytesize - Number of data bits (default 8).
scan_file - Optional YAML file specifying multiple slaves, functions,
and register blocks. Overrides registers/slave/function
parameters if provided.
slave - Modbus slave ID (default 1). Ignored if scan_file is provided.
function - Modbus function type to read: 'holding_registers',
'input_registers', 'coils', or 'discrete_inputs'.
Ignored if scan_file is provided.
interval - Seconds between consecutive reads. Must be >= 0.1.
sep - Separator string used when encoding output as text. Ignored
if encoding is None.
encoding - Character encoding for text output (default 'utf-8').
If None, raw bytes are returned.
encoding_errors - Strategy for handling encoding errors ('ignore'
by default). Other options: 'strict', 'replace',
'backslashreplace'.
timeout - Max time in seconds to wait for serial response. Defaults
to 2 seconds if None. Minimum timeout is 1s.
def
read(self):
377 def read(self): 378 """Read all configured polls and return a list of formatted records per slave. 379 380 Returns: 381 list[bytes|str|None]: None for polls that cannot be read; otherwise formatted output. 382 """ 383 with self._lock: 384 now = time.monotonic() 385 if now < self._next_read_time: 386 time.sleep(self._next_read_time - now) 387 start_time = time.monotonic() 388 389 if not self._client_connected(): 390 return None 391 392 records = [] 393 394 for poll in self.polls: 395 slave = poll["slave"] 396 func_name = self._FUNCTION_MAP.get(poll["function"]) 397 398 if not func_name: 399 logging.warning("Invalid function for slave=%s: %s", slave, poll["function"]) 400 total_count = sum(count for _, count in poll["registers"]) 401 readings = [[None] * total_count] 402 else: 403 func = getattr(self.client, func_name) 404 readings = [] 405 for start, count in poll["registers"]: 406 try: 407 resp = func(address=start, count=count, unit=slave) 408 if resp is None or getattr(resp, "isError", lambda: False)(): 409 logging.warning("Failed to read slave=%d addr=%d:%d", slave, start, count) 410 readings.append([None] * count) 411 else: 412 values = getattr(resp, "registers", getattr(resp, "bits", None)) 413 readings.append(values or [None] * count) 414 except (ModbusException, OSError, ConnectionError, ValueError, AttributeError) as exc: 415 logging.warning("Modbus read error slave=%d addr=%d:%d: %s", slave, start, count, exc) 416 readings.append([None] * count) 417 self._reset_client() 418 break 419 420 record = self._format_record(readings, poll["function"]) 421 if self.encoding is not None and isinstance(record, bytes): 422 record = record.decode(self.encoding, errors=self.encoding_errors) 423 record = self._encode_str(f"slave {slave}:{self.sep}{record}") 424 records.append(record) 425 426 self._next_read_time = start_time + self.interval 427 return records
Read all configured polls and return a list of formatted records per slave.
Returns: list[bytes|str|None]: None for polls that cannot be read; otherwise formatted output.