package com.mes.util; import com.fazecast.jSerialComm.SerialPort; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import java.io.ByteArrayOutputStream; import java.io.IOException; /** * Modbus RTU over RS232/RS485 serial port. */ public class ModbusRtuClient { public static final Logger log = LoggerFactory.getLogger(ModbusRtuClient.class); private static final int INTER_FRAME_MS = 5; private final SerialPort serialPort; private final int slaveId; private final int readTimeoutMs; public ModbusRtuClient(SerialPort serialPort, int slaveId, int readTimeoutMs) { this.serialPort = serialPort; this.slaveId = slaveId & 0xFF; this.readTimeoutMs = readTimeoutMs; } public int[] readHoldingRegisters(int startAddress, int quantity) throws IOException { if (quantity <= 0 || quantity > 125) { throw new IOException("Invalid register quantity: " + quantity); } byte[] request = new byte[8]; request[0] = (byte) slaveId; request[1] = 0x03; request[2] = (byte) ((startAddress >> 8) & 0xFF); request[3] = (byte) (startAddress & 0xFF); request[4] = (byte) ((quantity >> 8) & 0xFF); request[5] = (byte) (quantity & 0xFF); appendCrc(request, 6); byte[] response = transact(request); validateResponse(response, 0x03); int byteCount = response[2] & 0xFF; if (byteCount != quantity * 2) { throw new IOException("Unexpected byte count: " + byteCount); } int[] registers = new int[quantity]; for (int i = 0; i < quantity; i++) { int hi = response[3 + i * 2] & 0xFF; int lo = response[4 + i * 2] & 0xFF; registers[i] = (hi << 8) | lo; } return registers; } public void writeSingleRegister(int address, int value) throws IOException { writeMultipleRegisters(address, new int[]{value & 0xFFFF}); } public void writeMultipleRegisters(int startAddress, int[] values) throws IOException { if (values == null || values.length == 0 || values.length > 123) { throw new IOException("Invalid write values"); } int byteCount = values.length * 2; byte[] request = new byte[9 + byteCount]; request[0] = (byte) slaveId; request[1] = 0x10; request[2] = (byte) ((startAddress >> 8) & 0xFF); request[3] = (byte) (startAddress & 0xFF); request[4] = (byte) ((values.length >> 8) & 0xFF); request[5] = (byte) (values.length & 0xFF); request[6] = (byte) byteCount; for (int i = 0; i < values.length; i++) { request[7 + i * 2] = (byte) ((values[i] >> 8) & 0xFF); request[8 + i * 2] = (byte) (values[i] & 0xFF); } appendCrc(request, 7 + byteCount); byte[] response = transact(request); validateResponse(response, 0x10); } private byte[] transact(byte[] request) throws IOException { drainInput(); int written = serialPort.writeBytes(request, request.length); if (written != request.length) { throw new IOException("Modbus write incomplete"); } log.info("Modbus TX: {}", toHex(request)); sleepQuiet(INTER_FRAME_MS); ByteArrayOutputStream out = new ByteArrayOutputStream(); byte[] buffer = new byte[256]; long deadline = System.currentTimeMillis() + readTimeoutMs; long lastDataTime = 0; while (System.currentTimeMillis() < deadline) { int read = serialPort.readBytes(buffer, buffer.length); if (read > 0) { out.write(buffer, 0, read); lastDataTime = System.currentTimeMillis(); byte[] data = out.toByteArray(); if (isResponseComplete(data)) { log.info("Modbus RX: {}", toHex(data)); return data; } } else if (out.size() > 0 && lastDataTime > 0 && System.currentTimeMillis() - lastDataTime > 80) { break; } else { sleepQuiet(20); } } if (out.size() == 0) { throw new IOException("Modbus response timeout"); } byte[] partial = out.toByteArray(); log.info("Modbus RX(partial): {}", toHex(partial)); return partial; } private static boolean isResponseComplete(byte[] data) { if (data.length < 5) { return false; } int functionCode = data[1] & 0xFF; if ((functionCode & 0x80) != 0) { return data.length >= 5; } if (functionCode == 0x03) { int byteCount = data[2] & 0xFF; return data.length >= 3 + byteCount + 2; } if (functionCode == 0x10) { return data.length >= 8; } return data.length >= 5; } private static void validateResponse(byte[] response, int functionCode) throws IOException { if (response.length < 5) { throw new IOException("Modbus response too short"); } if ((response[1] & 0xFF) == (functionCode | 0x80)) { throw new IOException("Modbus exception code: " + (response[2] & 0xFF)); } if ((response[1] & 0xFF) != functionCode) { throw new IOException("Unexpected function code: " + (response[1] & 0xFF)); } if (!verifyCrc(response)) { throw new IOException("Modbus CRC error"); } } private void drainInput() { byte[] buf = new byte[256]; try { while (serialPort.bytesAvailable() > 0) { serialPort.readBytes(buf, Math.min(buf.length, serialPort.bytesAvailable())); } } catch (Exception ignored) { } } static int crc16(byte[] data, int length) { int crc = 0xFFFF; for (int i = 0; i < length; i++) { crc ^= (data[i] & 0xFF); for (int j = 0; j < 8; j++) { if ((crc & 0x0001) != 0) { crc = (crc >> 1) ^ 0xA001; } else { crc >>= 1; } } } return crc & 0xFFFF; } private static void appendCrc(byte[] frame, int length) { int crc = crc16(frame, length); frame[length] = (byte) (crc & 0xFF); frame[length + 1] = (byte) ((crc >> 8) & 0xFF); } private static boolean verifyCrc(byte[] frame) { if (frame.length < 3) { return false; } int crc = crc16(frame, frame.length - 2); int lo = frame[frame.length - 2] & 0xFF; int hi = frame[frame.length - 1] & 0xFF; return crc == ((hi << 8) | lo); } static float registersToFloat(int regHi, int regLo) { int bits = ((regHi & 0xFFFF) << 16) | (regLo & 0xFFFF); return Float.intBitsToFloat(bits); } static String toHex(byte[] data) { StringBuilder sb = new StringBuilder(); for (byte b : data) { sb.append(String.format("%02X ", b & 0xFF)); } return sb.toString().trim(); } private static void sleepQuiet(int ms) throws IOException { try { Thread.sleep(ms); } catch (InterruptedException e) { Thread.currentThread().interrupt(); throw new IOException("Modbus interrupted", e); } } }