s7_service.py 6.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213
  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": "S7TCP", "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 == "DB":
  81. return client.db_read(read.db, read.start, size)
  82. return client.read_area(AREA_CODES[read.area], 0, read.start, size)
  83. def trace_connect(trace: S7CommunicationTrace, request: S7RawReadRequest | S7PointReadRequest) -> None:
  84. trace.append(
  85. "Tx:S7_CONNECT "
  86. f"ip={request.ip} port={request.port} rock={request.rock} slot={request.slot} "
  87. f"tsap_conn_type={request.tsap_conn_type}"
  88. )
  89. def trace_read(trace: S7CommunicationTrace, read: S7ReadSpec | S7PointSpec, size: int) -> None:
  90. trace.append(f"Tx:S7_READ area={read.area} db={read.db} start={read.start} size={size}")
  91. def read_raw(request: S7RawReadRequest) -> dict[str, Any]:
  92. trace = S7CommunicationTrace()
  93. client: Client | None = None
  94. error: Exception | None = None
  95. try:
  96. trace_connect(trace, request)
  97. client = connect_client(request.ip, request.port, request.rock, request.slot, request.tsap_conn_type)
  98. trace.append(f"Rx:S7_CONNECTED{pdu_length_message(client)}")
  99. trace_read(trace, request.read, request.read.size)
  100. payload = read_s7_bytes(client, request.read, request.read.size)
  101. trace.append(f"Rx:{bytes_to_hex(payload)}")
  102. except Exception as exc:
  103. error = exc
  104. trace.append(f"Rx:S7_ERROR message={exc}")
  105. finally:
  106. if client is not None:
  107. trace.append("Tx:S7_DISCONNECT")
  108. close_client(client)
  109. trace.append("Rx:S7_DISCONNECTED")
  110. if error is not None:
  111. return response_payload(1, str(error), {"device": device_payload(request), "communication": trace.as_list()})
  112. return response_payload(0, "success", {"device": device_payload(request), "communication": trace.as_list()})
  113. def read_points(request: S7PointReadRequest) -> dict[str, Any]:
  114. trace = S7CommunicationTrace()
  115. client: Client | None = None
  116. points: list[dict[str, Any]] = []
  117. error: Exception | None = None
  118. try:
  119. trace_connect(trace, request)
  120. client = connect_client(request.ip, request.port, request.rock, request.slot, request.tsap_conn_type)
  121. trace.append(f"Rx:S7_CONNECTED{pdu_length_message(client)}")
  122. for point in request.points:
  123. size = byte_size_for_type(point.type)
  124. trace_read(trace, point, size)
  125. payload = read_s7_bytes(client, point, size)
  126. trace.append(f"Rx:{bytes_to_hex(payload)}")
  127. result = {
  128. "area": point.area,
  129. "db": point.db,
  130. "start": point.start,
  131. "type": point.type,
  132. "value": convert_s7_value(payload, point.type, point.bit),
  133. }
  134. if point.bit is not None:
  135. result["bit"] = point.bit
  136. points.append(result)
  137. except Exception as exc:
  138. error = exc
  139. trace.append(f"Rx:S7_ERROR message={exc}")
  140. finally:
  141. if client is not None:
  142. trace.append("Tx:S7_DISCONNECT")
  143. close_client(client)
  144. trace.append("Rx:S7_DISCONNECTED")
  145. if error is not None:
  146. return response_payload(
  147. 1,
  148. str(error),
  149. {"device": device_payload(request), "points": points},
  150. )
  151. return response_payload(
  152. 0,
  153. "success",
  154. {"device": device_payload(request), "points": points},
  155. )
  156. def connect_scan(request: S7ConnectScanRequest) -> dict[str, Any]:
  157. available: list[dict[str, Any]] = []
  158. for rock, slot, tsap_conn_type in product(range(3), range(3), CONNECTION_TYPES):
  159. client: Client | None = None
  160. try:
  161. client = connect_client(request.ip, 102, rock, slot, tsap_conn_type)
  162. available.append({"rock": rock, "slot": slot, "tsap_conn_type": tsap_conn_type})
  163. except Exception:
  164. pass
  165. finally:
  166. close_client(client)
  167. return response_payload(0, "success", {"device": connect_scan_device_payload(request), "available": available})