s7_service.py 6.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215
  1. from __future__ import annotations
  2. from itertools import product
  3. from typing import Any
  4. from snap7.client import Client
  5. try:
  6. from snap7.type import Area
  7. except ImportError: # pragma: no cover - compatibility for older python-snap7 releases
  8. Area = None
  9. from app.response import response_payload
  10. from app.schemas.s7 import S7ConnectScanRequest, S7PointReadRequest, S7PointSpec, S7RawReadRequest, S7ReadSpec
  11. from app.services.s7_codec import byte_size_for_type, bytes_to_hex, convert_s7_value
  12. CONNECTION_TYPES = {
  13. "PG": 0x01,
  14. "OP": 0x02,
  15. "BASIC": 0x03,
  16. }
  17. AREA_CODES: dict[str, Any] = {
  18. "I": 0x81,
  19. "Q": 0x82,
  20. "M": 0x83,
  21. "DB": 0x84,
  22. }
  23. if Area is not None:
  24. AREA_CODES.update({"I": Area.PE, "Q": Area.PA, "M": Area.MK, "DB": Area.DB})
  25. class S7CommunicationError(RuntimeError):
  26. def __init__(self, message: str) -> None:
  27. super().__init__(message)
  28. class S7CommunicationTrace:
  29. def __init__(self) -> None:
  30. self._messages: list[str] = []
  31. def append(self, message: str) -> None:
  32. self._messages.append(message)
  33. def as_list(self) -> list[str]:
  34. return list(self._messages)
  35. def device_payload(request: S7RawReadRequest | S7PointReadRequest) -> dict[str, Any]:
  36. return {
  37. "device_type": request.device_type,
  38. "ip": request.ip,
  39. "port": request.port,
  40. "rock": request.rock,
  41. "slot": request.slot,
  42. "tsap_conn_type": request.tsap_conn_type,
  43. }
  44. def connect_scan_device_payload(request: S7ConnectScanRequest) -> dict[str, Any]:
  45. return {"device_type": request.device_type, "ip": request.ip, "port": 102}
  46. def connect_client(ip: str, port: int, rock: int, slot: int, tsap_conn_type: str) -> Client:
  47. client = Client()
  48. try:
  49. client.set_connection_type(CONNECTION_TYPES[tsap_conn_type])
  50. client.connect(ip, rock, slot, tcp_port=port)
  51. except TypeError:
  52. try:
  53. client.connect(ip, rock, slot, port)
  54. except Exception:
  55. close_client(client)
  56. raise
  57. except Exception:
  58. close_client(client)
  59. raise
  60. if not client.get_connected():
  61. close_client(client)
  62. raise S7CommunicationError(f"failed to connect to {ip}:{port}")
  63. return client
  64. def close_client(client: Client | None) -> None:
  65. if client is None:
  66. return
  67. try:
  68. if client.get_connected():
  69. client.disconnect()
  70. finally:
  71. destroy = getattr(client, "destroy", None)
  72. if callable(destroy):
  73. destroy()
  74. def pdu_length_message(client: Client) -> str:
  75. try:
  76. return f" pdu_length={client.get_pdu_length()}"
  77. except Exception:
  78. return ""
  79. def read_s7_bytes(client: Client, read: S7ReadSpec | S7PointSpec, size: int) -> bytearray:
  80. if read.area == "V":
  81. return client.db_read(1, read.start, size)
  82. if read.area == "DB":
  83. return client.db_read(read.db, read.start, size)
  84. return client.read_area(AREA_CODES[read.area], 0, read.start, size)
  85. def trace_connect(trace: S7CommunicationTrace, request: S7RawReadRequest | S7PointReadRequest) -> None:
  86. trace.append(
  87. "Tx:S7_CONNECT "
  88. f"ip={request.ip} port={request.port} rock={request.rock} slot={request.slot} "
  89. f"tsap_conn_type={request.tsap_conn_type}"
  90. )
  91. def trace_read(trace: S7CommunicationTrace, read: S7ReadSpec | S7PointSpec, size: int) -> None:
  92. trace.append(f"Tx:S7_READ area={read.area} db={read.db} start={read.start} size={size}")
  93. def read_raw(request: S7RawReadRequest) -> dict[str, Any]:
  94. trace = S7CommunicationTrace()
  95. client: Client | None = None
  96. error: Exception | None = None
  97. try:
  98. trace_connect(trace, request)
  99. client = connect_client(request.ip, request.port, request.rock, request.slot, request.tsap_conn_type)
  100. trace.append(f"Rx:S7_CONNECTED{pdu_length_message(client)}")
  101. trace_read(trace, request.read, request.read.size)
  102. payload = read_s7_bytes(client, request.read, request.read.size)
  103. trace.append(f"Rx:{bytes_to_hex(payload)}")
  104. except Exception as exc:
  105. error = exc
  106. trace.append(f"Rx:S7_ERROR message={exc}")
  107. finally:
  108. if client is not None:
  109. trace.append("Tx:S7_DISCONNECT")
  110. close_client(client)
  111. trace.append("Rx:S7_DISCONNECTED")
  112. if error is not None:
  113. return response_payload(1, str(error), {"device": device_payload(request), "communication": trace.as_list()})
  114. return response_payload(0, "success", {"device": device_payload(request), "communication": trace.as_list()})
  115. def read_points(request: S7PointReadRequest) -> dict[str, Any]:
  116. trace = S7CommunicationTrace()
  117. client: Client | None = None
  118. points: list[dict[str, Any]] = []
  119. error: Exception | None = None
  120. try:
  121. trace_connect(trace, request)
  122. client = connect_client(request.ip, request.port, request.rock, request.slot, request.tsap_conn_type)
  123. trace.append(f"Rx:S7_CONNECTED{pdu_length_message(client)}")
  124. for point in request.points:
  125. size = byte_size_for_type(point.type)
  126. trace_read(trace, point, size)
  127. payload = read_s7_bytes(client, point, size)
  128. trace.append(f"Rx:{bytes_to_hex(payload)}")
  129. result = {
  130. "area": point.area,
  131. "db": point.db,
  132. "start": point.start,
  133. "type": point.type,
  134. "value": convert_s7_value(payload, point.type, point.bit),
  135. }
  136. if point.bit is not None:
  137. result["bit"] = point.bit
  138. points.append(result)
  139. except Exception as exc:
  140. error = exc
  141. trace.append(f"Rx:S7_ERROR message={exc}")
  142. finally:
  143. if client is not None:
  144. trace.append("Tx:S7_DISCONNECT")
  145. close_client(client)
  146. trace.append("Rx:S7_DISCONNECTED")
  147. if error is not None:
  148. return response_payload(
  149. 1,
  150. str(error),
  151. {"device": device_payload(request), "points": points},
  152. )
  153. return response_payload(
  154. 0,
  155. "success",
  156. {"device": device_payload(request), "points": points},
  157. )
  158. def connect_scan(request: S7ConnectScanRequest) -> dict[str, Any]:
  159. available: list[dict[str, Any]] = []
  160. for rock, slot, tsap_conn_type in product(range(3), range(3), CONNECTION_TYPES):
  161. client: Client | None = None
  162. try:
  163. client = connect_client(request.ip, 102, rock, slot, tsap_conn_type)
  164. available.append({"rock": rock, "slot": slot, "tsap_conn_type": tsap_conn_type})
  165. except Exception:
  166. pass
  167. finally:
  168. close_client(client)
  169. return response_payload(0, "success", {"device": connect_scan_device_payload(request), "available": available})