collector_api.py 32 KB

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