collector_api.py 25 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719
  1. from __future__ import annotations
  2. from typing import Any
  3. from .auth import find_project_config, resolve_project_token
  4. from .http_client import request_json
  5. from .protocols import MODBUS_SPEC, S7_SPEC
  6. from .protocols.modbus import MODBUS_POINT_TYPE_ALIASES, MODBUS_REGISTER_TYPE_ALIASES
  7. from .protocols.s7 import S7_POINT_TYPE_ALIASES, S7_REGISTER_TYPE_ALIASES
  8. def _merge_defaults(defaults: dict[str, Any], payload: dict[str, Any]) -> dict[str, Any]:
  9. merged = dict(defaults)
  10. merged.update(payload)
  11. return merged
  12. def _success_response(response: dict[str, Any]) -> bool:
  13. return response.get("state") == 0
  14. def _require_non_empty_text(payload: dict[str, Any], field_name: str) -> str:
  15. value = str(payload.get(field_name) or "").strip()
  16. if not value:
  17. raise ValueError(f"payload.{field_name} is required")
  18. return value
  19. def _require_present(payload: dict[str, Any], field_name: str) -> Any:
  20. if field_name not in payload or payload.get(field_name) is None:
  21. raise ValueError(f"payload.{field_name} is required")
  22. return payload[field_name]
  23. def _normalize_modbus_device_payload(payload: dict[str, Any]) -> dict[str, Any]:
  24. normalized = dict(payload)
  25. normalized["name"] = _require_non_empty_text(normalized, "name")
  26. normalized["ip"] = _require_non_empty_text(normalized, "ip")
  27. _require_present(normalized, "device_type")
  28. _require_present(normalized, "port")
  29. _require_present(normalized, "slave_id")
  30. _require_present(normalized, "word_order")
  31. _require_present(normalized, "byte_order")
  32. address_base = _require_present(normalized, "address_base")
  33. normalized["address_offset"] = address_base
  34. normalized.pop("address_base", None)
  35. return normalized
  36. def _normalize_modbus_device_edit_payload(payload: dict[str, Any]) -> dict[str, Any]:
  37. normalized = dict(payload)
  38. normalized["ori_id"] = _normalize_positive_int(_require_present(normalized, "ori_id"), "payload.ori_id")
  39. normalized["name"] = _require_non_empty_text(normalized, "name")
  40. if normalized.get("type") is None:
  41. normalized["type"] = _require_present(normalized, "device_type")
  42. try:
  43. normalized["type"] = int(normalized["type"])
  44. except Exception as exc:
  45. raise ValueError("payload.type must be one of 1, 2, 3, 4, 5") from exc
  46. if normalized["type"] not in {1, 2, 3, 4, 5}:
  47. raise ValueError("payload.type must be one of 1, 2, 3, 4, 5")
  48. normalized.pop("device_type", None)
  49. if normalized["type"] == 2:
  50. normalized["serial_port"] = _require_non_empty_text(normalized, "serial_port")
  51. else:
  52. normalized["ip"] = _require_non_empty_text(normalized, "ip")
  53. _require_present(normalized, "port")
  54. _require_present(normalized, "slave_id")
  55. _require_present(normalized, "word_order")
  56. _require_present(normalized, "byte_order")
  57. if "address_base" in normalized:
  58. normalized["address_offset"] = normalized["address_base"]
  59. normalized.pop("address_base", None)
  60. if "group_id" in normalized and "device_group_id" not in normalized:
  61. normalized["device_group_id"] = normalized["group_id"]
  62. normalized.pop("group_id", None)
  63. return normalized
  64. def _normalize_modbus_point_payload(payload: dict[str, Any], *, require_device_id: bool = True) -> dict[str, Any]:
  65. normalized = dict(payload)
  66. normalized["name"] = _require_non_empty_text(normalized, "name")
  67. _require_present(normalized, "address")
  68. if require_device_id:
  69. normalized["device_id"] = _normalize_positive_int(
  70. _require_present(normalized, "device_id"),
  71. "payload.device_id",
  72. )
  73. raw_type = _require_non_empty_text(normalized, "type")
  74. normalized_type = MODBUS_POINT_TYPE_ALIASES.get(raw_type)
  75. if normalized_type is None:
  76. normalized_type = MODBUS_POINT_TYPE_ALIASES.get(raw_type.lower())
  77. if normalized_type is None:
  78. raise ValueError(
  79. "payload.type is invalid; use one of bool, int16, uint16, int32, "
  80. "uint32, int64, uint64, float32, float64, or a documented alias "
  81. "such as SHORT, WORD, LONG, DWORD, FLOAT, DOUBLE"
  82. )
  83. normalized["type"] = normalized_type
  84. if normalized.get("func_code") is None:
  85. register_type = _require_non_empty_text(normalized, "register_type")
  86. func_code = MODBUS_REGISTER_TYPE_ALIASES.get(register_type)
  87. if func_code is None:
  88. func_code = MODBUS_REGISTER_TYPE_ALIASES.get(register_type.lower())
  89. if func_code is None:
  90. raise ValueError(
  91. "payload.register_type is invalid; use coil, discrete_input, "
  92. "holding_register, input_register, or func_code 1/2/3/4"
  93. )
  94. normalized["func_code"] = func_code
  95. else:
  96. try:
  97. normalized["func_code"] = int(str(normalized["func_code"]).strip())
  98. except Exception as exc:
  99. raise ValueError("payload.func_code must be one of 1, 2, 3, 4") from exc
  100. if normalized["func_code"] not in {1, 2, 3, 4}:
  101. raise ValueError("payload.func_code must be one of 1, 2, 3, 4")
  102. if "register_type" in normalized:
  103. normalized.pop("register_type")
  104. return normalized
  105. def _normalize_modbus_point_edit_payload(payload: dict[str, Any]) -> dict[str, Any]:
  106. _require_present(payload, "ori_id")
  107. normalized = _normalize_modbus_point_payload(payload, require_device_id=False)
  108. normalized["ori_id"] = _normalize_positive_int(_require_present(normalized, "ori_id"), "payload.ori_id")
  109. return normalized
  110. def _normalize_s7_device_payload(payload: dict[str, Any]) -> dict[str, Any]:
  111. normalized = dict(payload)
  112. normalized["name"] = _require_non_empty_text(normalized, "name")
  113. normalized["ip"] = _require_non_empty_text(normalized, "ip")
  114. _require_present(normalized, "rock")
  115. _require_present(normalized, "slot")
  116. if "group_id" in normalized and "device_group_id" not in normalized:
  117. normalized["device_group_id"] = normalized["group_id"]
  118. normalized.pop("group_id", None)
  119. if "tsap_conn_type" in normalized and normalized["tsap_conn_type"] is not None:
  120. normalized["tsap_conn_type"] = _normalize_s7_tsap_conn_type(normalized["tsap_conn_type"])
  121. elif int(normalized.get("device_type", 1)) == 3:
  122. normalized["tsap_conn_type"] = "OP"
  123. else:
  124. normalized["tsap_conn_type"] = "PG"
  125. return normalized
  126. def _normalize_s7_device_edit_payload(payload: dict[str, Any]) -> dict[str, Any]:
  127. normalized = _normalize_s7_device_payload(payload)
  128. if normalized.get("id") is None:
  129. normalized["id"] = _require_present(normalized, "ori_id")
  130. normalized["id"] = _normalize_positive_int(normalized["id"], "payload.id")
  131. normalized.pop("ori_id", None)
  132. return normalized
  133. def _normalize_s7_point_payload(payload: dict[str, Any], *, require_device_id: bool = True) -> dict[str, Any]:
  134. normalized = dict(payload)
  135. normalized["name"] = _require_non_empty_text(normalized, "name")
  136. normalized["address"] = _require_non_empty_text(normalized, "address")
  137. if require_device_id:
  138. normalized["device_id"] = _normalize_positive_int(
  139. _require_present(normalized, "device_id"),
  140. "payload.device_id",
  141. )
  142. raw_type = normalized.get("data_type")
  143. if raw_type is None:
  144. raw_type = _require_non_empty_text(normalized, "type")
  145. else:
  146. raw_type = str(raw_type).strip()
  147. if not raw_type:
  148. raise ValueError("payload.data_type is required")
  149. normalized_type = S7_POINT_TYPE_ALIASES.get(raw_type)
  150. if normalized_type is None:
  151. normalized_type = S7_POINT_TYPE_ALIASES.get(raw_type.lower())
  152. if normalized_type is None:
  153. raise ValueError(
  154. "payload.data_type is invalid; use one of bool, uint8, int8, "
  155. "uint16, int16, uint32, int32, float32, float64, or a documented "
  156. "alias such as BOOL, BYTE, SINT, WORD, INT, DWORD, DINT, REAL, LREAL"
  157. )
  158. normalized["data_type"] = normalized_type
  159. normalized.pop("type", None)
  160. if normalized.get("register_type") is None:
  161. register_area = _require_non_empty_text(normalized, "register_area")
  162. register_type = S7_REGISTER_TYPE_ALIASES.get(register_area)
  163. if register_type is None:
  164. register_type = S7_REGISTER_TYPE_ALIASES.get(register_area.lower())
  165. if register_type is None:
  166. raise ValueError("payload.register_area is invalid; use I, Q, M, DB, V, AI, or register_type 1..6")
  167. normalized["register_type"] = register_type
  168. else:
  169. try:
  170. normalized["register_type"] = int(str(normalized["register_type"]).strip())
  171. except Exception as exc:
  172. raise ValueError("payload.register_type must be one of 1, 2, 3, 4, 5, 6") from exc
  173. if normalized["register_type"] not in {1, 2, 3, 4, 5, 6}:
  174. raise ValueError("payload.register_type must be one of 1, 2, 3, 4, 5, 6")
  175. normalized.pop("register_area", None)
  176. if "group_id" in normalized and "group_Id" not in normalized:
  177. normalized["group_Id"] = normalized["group_id"]
  178. normalized.pop("group_id", None)
  179. return normalized
  180. def _normalize_s7_point_edit_payload(payload: dict[str, Any]) -> dict[str, Any]:
  181. normalized = _normalize_s7_point_payload(payload)
  182. if normalized.get("id") is None:
  183. normalized["id"] = _require_present(normalized, "ori_id")
  184. normalized["id"] = _normalize_positive_int(normalized["id"], "payload.id")
  185. normalized.pop("ori_id", None)
  186. return normalized
  187. def _request_collector(
  188. project_key: str,
  189. method: str,
  190. path: str,
  191. *,
  192. json_payload: dict[str, Any] | None = None,
  193. ) -> dict[str, Any]:
  194. project = find_project_config(project_key)
  195. authorization = resolve_project_token(project)
  196. response_payload = request_json(
  197. method,
  198. f"{project['data_collector_base_url']}{path}",
  199. authorization,
  200. json_payload=json_payload,
  201. )
  202. if not isinstance(response_payload, dict):
  203. raise ValueError(f"collector API returned invalid payload for {path}: {response_payload}")
  204. return response_payload
  205. def _build_modbus_device_create_payload(payload: dict[str, Any]) -> dict[str, Any]:
  206. return _merge_defaults(
  207. MODBUS_SPEC.device_defaults,
  208. _normalize_modbus_device_payload(payload),
  209. )
  210. def _build_modbus_point_create_payload(payload: dict[str, Any]) -> dict[str, Any]:
  211. return _merge_defaults(
  212. MODBUS_SPEC.point_defaults,
  213. _normalize_modbus_point_payload(payload),
  214. )
  215. def _build_s7_device_create_payload(payload: dict[str, Any]) -> dict[str, Any]:
  216. return _merge_defaults(
  217. S7_SPEC.device_defaults,
  218. _normalize_s7_device_payload(payload),
  219. )
  220. def _build_s7_point_create_payload(payload: dict[str, Any]) -> dict[str, Any]:
  221. return _merge_defaults(
  222. S7_SPEC.point_defaults,
  223. _normalize_s7_point_payload(payload),
  224. )
  225. def create_modbus_device(project_key: str, payload: dict[str, Any]) -> dict[str, Any]:
  226. return _request_collector(
  227. project_key,
  228. "POST",
  229. MODBUS_SPEC.create_device_path,
  230. json_payload=_build_modbus_device_create_payload(payload),
  231. )
  232. def create_modbus_point(project_key: str, payload: dict[str, Any]) -> dict[str, Any]:
  233. return _request_collector(
  234. project_key,
  235. "POST",
  236. MODBUS_SPEC.create_point_path,
  237. json_payload=_build_modbus_point_create_payload(payload),
  238. )
  239. def create_s7_device(project_key: str, payload: dict[str, Any]) -> dict[str, Any]:
  240. return _request_collector(
  241. project_key,
  242. "POST",
  243. S7_SPEC.create_device_path,
  244. json_payload=_build_s7_device_create_payload(payload),
  245. )
  246. def create_s7_point(project_key: str, payload: dict[str, Any]) -> dict[str, Any]:
  247. return _request_collector(
  248. project_key,
  249. "POST",
  250. S7_SPEC.create_point_path,
  251. json_payload=_build_s7_point_create_payload(payload),
  252. )
  253. def create_modbus_devices(project_key: str, devices: list[dict[str, Any]]) -> dict[str, Any]:
  254. return _create_devices_batch(
  255. project_key,
  256. devices,
  257. create_path=MODBUS_SPEC.create_device_path,
  258. build_payload=_build_modbus_device_create_payload,
  259. match_device=_match_modbus_device,
  260. )
  261. def create_modbus_points(project_key: str, points: list[dict[str, Any]]) -> dict[str, Any]:
  262. return _create_points_batch(
  263. project_key,
  264. points,
  265. create_path=MODBUS_SPEC.create_point_path,
  266. build_payload=_build_modbus_point_create_payload,
  267. )
  268. def create_s7_devices(project_key: str, devices: list[dict[str, Any]]) -> dict[str, Any]:
  269. return _create_devices_batch(
  270. project_key,
  271. devices,
  272. create_path=S7_SPEC.create_device_path,
  273. build_payload=_build_s7_device_create_payload,
  274. match_device=_match_s7_device,
  275. )
  276. def create_s7_points(project_key: str, points: list[dict[str, Any]]) -> dict[str, Any]:
  277. return _create_points_batch(
  278. project_key,
  279. points,
  280. create_path=S7_SPEC.create_point_path,
  281. build_payload=_build_s7_point_create_payload,
  282. )
  283. def _create_devices_batch(
  284. project_key: str,
  285. devices: list[dict[str, Any]],
  286. *,
  287. create_path: str,
  288. build_payload: Any,
  289. match_device: Any,
  290. ) -> dict[str, Any]:
  291. results: list[dict[str, Any]] = []
  292. errors: list[dict[str, Any]] = []
  293. created: list[tuple[dict[str, Any], dict[str, Any]]] = []
  294. for index, item in enumerate(devices):
  295. if not isinstance(item, dict):
  296. errors.append({"index": index, "stage": "create_device", "error": "device item must be an object"})
  297. continue
  298. result: dict[str, Any] = {
  299. "index": index,
  300. "name": item.get("name"),
  301. "device_id": None,
  302. "create_response": None,
  303. "matched_device": None,
  304. }
  305. results.append(result)
  306. try:
  307. request_payload = build_payload(item)
  308. response = _request_collector(project_key, "POST", create_path, json_payload=request_payload)
  309. except Exception as exc:
  310. errors.append({"index": index, "name": item.get("name"), "stage": "create_device", "error": str(exc)})
  311. continue
  312. result["create_response"] = response
  313. if _success_response(response):
  314. created.append((request_payload, result))
  315. else:
  316. errors.append(
  317. {
  318. "index": index,
  319. "name": item.get("name"),
  320. "stage": "create_device",
  321. "error": response.get("state_info") or response.get("msg") or str(response),
  322. }
  323. )
  324. if created:
  325. try:
  326. device_list = list_devices(project_key, num_points=False)
  327. flattened_devices = _flatten_device_tree(device_list.get("devices", []))
  328. except Exception as exc:
  329. for _, result in created:
  330. errors.append(
  331. {
  332. "index": result["index"],
  333. "name": result.get("name"),
  334. "stage": "match_device",
  335. "error": str(exc),
  336. }
  337. )
  338. else:
  339. for request_payload, result in created:
  340. candidates = [item for item in flattened_devices if match_device(item, request_payload)]
  341. candidates = [item for item in candidates if _safe_int(item.get("id")) is not None]
  342. if not candidates:
  343. errors.append(
  344. {
  345. "index": result["index"],
  346. "name": result.get("name"),
  347. "stage": "match_device",
  348. "error": "created device was not found in device list",
  349. }
  350. )
  351. continue
  352. selected = max(candidates, key=lambda item: _safe_int(item.get("id")) or 0)
  353. result["device_id"] = _safe_int(selected.get("id"))
  354. result["matched_device"] = selected
  355. matched = sum(1 for item in results if item.get("device_id") is not None)
  356. return {
  357. "state": 0 if not errors else 1,
  358. "summary": {
  359. "total": len(devices),
  360. "created": len(created),
  361. "matched": matched,
  362. "failed": len(errors),
  363. },
  364. "results": results,
  365. "errors": errors,
  366. }
  367. def _create_points_batch(
  368. project_key: str,
  369. points: list[dict[str, Any]],
  370. *,
  371. create_path: str,
  372. build_payload: Any,
  373. ) -> dict[str, Any]:
  374. results: list[dict[str, Any]] = []
  375. errors: list[dict[str, Any]] = []
  376. for index, item in enumerate(points):
  377. if not isinstance(item, dict):
  378. errors.append({"index": index, "stage": "create_point", "error": "point item must be an object"})
  379. continue
  380. result: dict[str, Any] = {
  381. "index": index,
  382. "name": item.get("name"),
  383. "device_id": item.get("device_id"),
  384. "response": None,
  385. }
  386. results.append(result)
  387. try:
  388. request_payload = build_payload(item)
  389. response = _request_collector(project_key, "POST", create_path, json_payload=request_payload)
  390. except Exception as exc:
  391. errors.append({"index": index, "name": item.get("name"), "stage": "create_point", "error": str(exc)})
  392. continue
  393. result["response"] = response
  394. if not _success_response(response):
  395. errors.append(
  396. {
  397. "index": index,
  398. "name": item.get("name"),
  399. "stage": "create_point",
  400. "error": response.get("state_info") or response.get("msg") or str(response),
  401. }
  402. )
  403. success = sum(1 for item in results if isinstance(item.get("response"), dict) and _success_response(item["response"]))
  404. return {
  405. "state": 0 if not errors else 1,
  406. "summary": {
  407. "total": len(points),
  408. "success": success,
  409. "failed": len(errors),
  410. },
  411. "results": results,
  412. "errors": errors,
  413. }
  414. def _flatten_device_tree(items: list[Any]) -> list[dict[str, Any]]:
  415. flattened: list[dict[str, Any]] = []
  416. for item in items:
  417. if not isinstance(item, dict):
  418. continue
  419. flattened.append(item)
  420. groups = item.get("groups", [])
  421. if isinstance(groups, list):
  422. flattened.extend(_flatten_device_tree(groups))
  423. return flattened
  424. def _match_modbus_device(device: dict[str, Any], expected: dict[str, Any]) -> bool:
  425. if device.get("type") != "modbus":
  426. return False
  427. if not _same_text(device.get("name"), expected.get("name")):
  428. return False
  429. if not _same_int(device.get("device_type"), expected.get("device_type")):
  430. return False
  431. if not _same_int(device.get("slave_id"), expected.get("slave_id")):
  432. return False
  433. if not _same_int(device.get("group_id"), expected.get("group_id")):
  434. return False
  435. if not _same_int(device.get("address_offset"), expected.get("address_offset")):
  436. return False
  437. if _safe_int(expected.get("device_type")) == 2:
  438. return _same_text(device.get("serial_port"), expected.get("serial_port"))
  439. return _same_text(device.get("ip"), expected.get("ip")) and _same_int(device.get("port"), expected.get("port"))
  440. def _match_s7_device(device: dict[str, Any], expected: dict[str, Any]) -> bool:
  441. if device.get("type") != "s7":
  442. return False
  443. actual_group_id = device.get("device_group_id", device.get("group_id"))
  444. return (
  445. _same_text(device.get("name"), expected.get("name"))
  446. and _same_text(device.get("ip"), expected.get("ip"))
  447. and _same_int(device.get("port"), expected.get("port"))
  448. and _same_int(device.get("rock"), expected.get("rock"))
  449. and _same_int(device.get("slot"), expected.get("slot"))
  450. and _same_int(device.get("device_type"), expected.get("device_type"))
  451. and _same_text(device.get("tsap_conn_type"), expected.get("tsap_conn_type"))
  452. and _same_int(actual_group_id, expected.get("device_group_id"))
  453. )
  454. def _same_text(left: Any, right: Any) -> bool:
  455. return str(left or "").strip() == str(right or "").strip()
  456. def _same_int(left: Any, right: Any) -> bool:
  457. normalized_left = _safe_int(left)
  458. normalized_right = _safe_int(right)
  459. return normalized_left is not None and normalized_left == normalized_right
  460. def _safe_int(value: Any) -> int | None:
  461. try:
  462. return int(value)
  463. except Exception:
  464. return None
  465. def edit_modbus_device(project_key: str, payload: dict[str, Any]) -> dict[str, Any]:
  466. return _request_collector(
  467. project_key,
  468. "POST",
  469. "/api/collector/modbus/device/edit",
  470. json_payload=_merge_defaults(
  471. {
  472. "timeout": 3,
  473. "is_persistent": True,
  474. "device_group_id": 0,
  475. "alarm_interval": 90,
  476. "collect_interval": 5,
  477. "address_offset": 0,
  478. "retry_times": 0,
  479. "mode": 0,
  480. },
  481. _normalize_modbus_device_edit_payload(payload),
  482. ),
  483. )
  484. def edit_modbus_point(project_key: str, payload: dict[str, Any]) -> dict[str, Any]:
  485. return _request_collector(
  486. project_key,
  487. "POST",
  488. "/api/collector/modbus/point/edit_collect_point",
  489. json_payload=_merge_defaults(
  490. MODBUS_SPEC.point_defaults,
  491. _normalize_modbus_point_edit_payload(payload),
  492. ),
  493. )
  494. def edit_s7_device(project_key: str, payload: dict[str, Any]) -> dict[str, Any]:
  495. return _request_collector(
  496. project_key,
  497. "POST",
  498. "/api/collector/s7/device/update",
  499. json_payload=_merge_defaults(
  500. S7_SPEC.device_defaults,
  501. _normalize_s7_device_edit_payload(payload),
  502. ),
  503. )
  504. def edit_s7_point(project_key: str, payload: dict[str, Any]) -> dict[str, Any]:
  505. return _request_collector(
  506. project_key,
  507. "POST",
  508. "/api/collector/s7/point/update",
  509. json_payload=_merge_defaults(
  510. S7_SPEC.point_defaults,
  511. _normalize_s7_point_edit_payload(payload),
  512. ),
  513. )
  514. def list_devices(project_key: str, num_points: bool = False) -> dict[str, Any]:
  515. num_points_text = "true" if num_points else "false"
  516. return _request_collector(
  517. project_key,
  518. "GET",
  519. f"/api/collector/device?num_points={num_points_text}",
  520. )
  521. def connect_device(project_key: str, device_id: int, device_type: str = "modbus") -> dict[str, Any]:
  522. return _set_device_connect_status(
  523. project_key,
  524. device_id=device_id,
  525. device_type=device_type,
  526. status=2,
  527. )
  528. def disconnect_device(project_key: str, device_id: int, device_type: str = "modbus") -> dict[str, Any]:
  529. return _set_device_connect_status(
  530. project_key,
  531. device_id=device_id,
  532. device_type=device_type,
  533. status=1,
  534. )
  535. def list_device_points(
  536. project_key: str,
  537. device_id: int,
  538. device_type: str = "modbus",
  539. group_id: int = 0,
  540. ) -> dict[str, Any]:
  541. return _request_collector(
  542. project_key,
  543. "POST",
  544. "/api/collector/common/device/get_collect_point",
  545. json_payload={
  546. "id": _normalize_positive_int(device_id, "device_id"),
  547. "type": _normalize_device_type(device_type),
  548. "group_id": _normalize_non_negative_int(group_id, "group_id"),
  549. },
  550. )
  551. def _set_device_connect_status(
  552. project_key: str,
  553. *,
  554. device_id: int,
  555. device_type: str,
  556. status: int,
  557. ) -> dict[str, Any]:
  558. return _request_collector(
  559. project_key,
  560. "POST",
  561. "/api/collector/common/device/set_connect_status",
  562. json_payload={
  563. "id": _normalize_positive_int(device_id, "device_id"),
  564. "type": _normalize_device_type(device_type),
  565. "status": status,
  566. },
  567. )
  568. def _normalize_device_type(device_type: str) -> str:
  569. normalized = str(device_type or "").strip()
  570. if not normalized:
  571. raise ValueError("device_type is required")
  572. return normalized
  573. def _normalize_s7_tsap_conn_type(value: Any) -> str:
  574. normalized = str(value or "").strip().upper()
  575. if normalized not in {"PG", "OP", "BASIC"}:
  576. raise ValueError("payload.tsap_conn_type must be one of PG, OP, BASIC")
  577. return normalized
  578. def _normalize_positive_int(value: int, field_name: str) -> int:
  579. try:
  580. normalized = int(value)
  581. except Exception as exc:
  582. raise ValueError(f"{field_name} must be a positive integer") from exc
  583. if normalized <= 0:
  584. raise ValueError(f"{field_name} must be a positive integer")
  585. return normalized
  586. def _normalize_non_negative_int(value: int, field_name: str) -> int:
  587. try:
  588. normalized = int(value)
  589. except Exception as exc:
  590. raise ValueError(f"{field_name} must be a non-negative integer") from exc
  591. if normalized < 0:
  592. raise ValueError(f"{field_name} must be a non-negative integer")
  593. return normalized