collector_api.py 32 KB

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