collector_api.py 32 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919
  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. if "register_type" in normalized:
  105. normalized.pop("register_type")
  106. return normalized
  107. def _normalize_modbus_point_edit_payload(payload: dict[str, Any]) -> dict[str, Any]:
  108. _require_present(payload, "ori_id")
  109. normalized = _normalize_modbus_point_payload(payload, require_device_id=False)
  110. normalized["ori_id"] = _normalize_positive_int(_require_present(normalized, "ori_id"), "payload.ori_id")
  111. return normalized
  112. def _normalize_s7_device_payload(payload: dict[str, Any]) -> dict[str, Any]:
  113. normalized = dict(payload)
  114. normalized["name"] = _require_non_empty_text(normalized, "name")
  115. normalized["ip"] = _require_non_empty_text(normalized, "ip")
  116. _require_present(normalized, "rock")
  117. _require_present(normalized, "slot")
  118. if "group_id" in normalized and "device_group_id" not in normalized:
  119. normalized["device_group_id"] = normalized["group_id"]
  120. normalized.pop("group_id", None)
  121. if "tsap_conn_type" in normalized and normalized["tsap_conn_type"] is not None:
  122. normalized["tsap_conn_type"] = _normalize_s7_tsap_conn_type(normalized["tsap_conn_type"])
  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 _normalize_bacnet_device_payload(payload: dict[str, Any]) -> dict[str, Any]:
  188. normalized = dict(payload)
  189. normalized["name"] = _require_non_empty_text(normalized, "name")
  190. normalized["ip"] = _require_non_empty_text(normalized, "ip")
  191. normalized["bacnet_device_id"] = _normalize_bacnet_id(
  192. _require_present(normalized, "bacnet_device_id"),
  193. "payload.bacnet_device_id",
  194. )
  195. normalized["device_type"] = 1
  196. normalized["type"] = "bacnet"
  197. if "device_id" in normalized:
  198. normalized.pop("device_id", None)
  199. if "device_group_id" in normalized and "group_id" not in normalized:
  200. normalized["group_id"] = normalized["device_group_id"]
  201. normalized.pop("device_group_id", None)
  202. if "bacnet_net" in normalized:
  203. normalized["bacnet_net"] = _normalize_non_negative_int(normalized["bacnet_net"], "payload.bacnet_net")
  204. elif "net" in normalized:
  205. normalized["bacnet_net"] = _normalize_non_negative_int(normalized["net"], "payload.net")
  206. normalized.pop("net", None)
  207. return normalized
  208. def _normalize_bacnet_device_edit_payload(payload: dict[str, Any]) -> dict[str, Any]:
  209. normalized = dict(payload)
  210. normalized["ori_id"] = _normalize_positive_int(_require_present(normalized, "ori_id"), "payload.ori_id")
  211. normalized["name"] = _require_non_empty_text(normalized, "name")
  212. normalized["ip"] = _require_non_empty_text(normalized, "ip")
  213. bacnet_device_id = _normalize_bacnet_id(
  214. _require_present(normalized, "bacnet_device_id"),
  215. "payload.bacnet_device_id",
  216. )
  217. normalized["device_id"] = str(bacnet_device_id)
  218. normalized["type"] = 1
  219. normalized.pop("bacnet_device_id", None)
  220. normalized.pop("device_type", None)
  221. if "group_id" in normalized and "device_group_id" not in normalized:
  222. normalized["device_group_id"] = normalized["group_id"]
  223. normalized.pop("group_id", None)
  224. if "bacnet_net" in normalized:
  225. normalized["net"] = _normalize_non_negative_int(normalized["bacnet_net"], "payload.bacnet_net")
  226. elif "net" in normalized:
  227. normalized["net"] = _normalize_non_negative_int(normalized["net"], "payload.net")
  228. normalized.pop("bacnet_net", None)
  229. return normalized
  230. def _normalize_bacnet_point_payload(payload: dict[str, Any], *, require_device_id: bool = True) -> dict[str, Any]:
  231. normalized = dict(payload)
  232. if require_device_id:
  233. normalized["device_id"] = _normalize_positive_int(
  234. _require_present(normalized, "device_id"),
  235. "payload.device_id",
  236. )
  237. normalized["object_type"] = normalize_bacnet_object_type(_require_non_empty_text(normalized, "object_type"))
  238. normalized["object_id"] = _normalize_bacnet_id(_require_present(normalized, "object_id"), "payload.object_id")
  239. object_name = str(normalized.get("object_name") or normalized.get("name") or "").strip()
  240. if not object_name:
  241. raise ValueError("payload.object_name is required")
  242. normalized["object_name"] = object_name
  243. normalized["name"] = str(normalized.get("name") or object_name).strip() or object_name
  244. if normalized.get("priority") is not None:
  245. normalized["priority"] = _normalize_positive_int(normalized["priority"], "payload.priority")
  246. return normalized
  247. def _normalize_bacnet_point_create_payload(payload: dict[str, Any]) -> dict[str, Any]:
  248. normalized = _merge_defaults(
  249. BACNET_SPEC.point_defaults,
  250. _normalize_bacnet_point_payload(payload),
  251. )
  252. device_id = normalized.pop("device_id")
  253. normalized.pop("invalid_values", None)
  254. normalized.pop("valid_range_start", None)
  255. normalized.pop("valid_range_end", None)
  256. return {"device_id": device_id, "points": [normalized]}
  257. def _normalize_bacnet_point_edit_payload(payload: dict[str, Any]) -> dict[str, Any]:
  258. normalized = _merge_defaults(
  259. BACNET_SPEC.point_defaults,
  260. _normalize_bacnet_point_payload(payload, require_device_id=False),
  261. )
  262. normalized["id"] = _normalize_positive_int(_require_present(normalized, "ori_id"), "payload.ori_id")
  263. normalized.pop("ori_id", None)
  264. normalized.pop("device_id", None)
  265. return normalized
  266. def _request_collector(
  267. project_key: str,
  268. method: str,
  269. path: str,
  270. *,
  271. json_payload: dict[str, Any] | None = None,
  272. ) -> dict[str, Any]:
  273. project = find_project_config(project_key)
  274. authorization = resolve_project_token(project)
  275. response_payload = request_json(
  276. method,
  277. f"{project['data_collector_base_url']}{path}",
  278. authorization,
  279. json_payload=json_payload,
  280. )
  281. if not isinstance(response_payload, dict):
  282. raise ValueError(f"collector API returned invalid payload for {path}: {response_payload}")
  283. return response_payload
  284. def _build_modbus_device_create_payload(payload: dict[str, Any]) -> dict[str, Any]:
  285. return _merge_defaults(
  286. MODBUS_SPEC.device_defaults,
  287. _normalize_modbus_device_payload(payload),
  288. )
  289. def _build_modbus_point_create_payload(payload: dict[str, Any]) -> dict[str, Any]:
  290. return _merge_defaults(
  291. MODBUS_SPEC.point_defaults,
  292. _normalize_modbus_point_payload(payload),
  293. )
  294. def _build_s7_device_create_payload(payload: dict[str, Any]) -> dict[str, Any]:
  295. return _merge_defaults(
  296. S7_SPEC.device_defaults,
  297. _normalize_s7_device_payload(payload),
  298. )
  299. def _build_s7_point_create_payload(payload: dict[str, Any]) -> dict[str, Any]:
  300. return _merge_defaults(
  301. S7_SPEC.point_defaults,
  302. _normalize_s7_point_payload(payload),
  303. )
  304. def _build_bacnet_device_create_payload(payload: dict[str, Any]) -> dict[str, Any]:
  305. return _merge_defaults(
  306. BACNET_SPEC.device_defaults,
  307. _normalize_bacnet_device_payload(payload),
  308. )
  309. def _build_bacnet_point_create_payload(payload: dict[str, Any]) -> dict[str, Any]:
  310. return _normalize_bacnet_point_create_payload(payload)
  311. def create_modbus_device(project_key: str, payload: dict[str, Any]) -> dict[str, Any]:
  312. return _request_collector(
  313. project_key,
  314. "POST",
  315. MODBUS_SPEC.create_device_path,
  316. json_payload=_build_modbus_device_create_payload(payload),
  317. )
  318. def create_modbus_point(project_key: str, payload: dict[str, Any]) -> dict[str, Any]:
  319. return _request_collector(
  320. project_key,
  321. "POST",
  322. MODBUS_SPEC.create_point_path,
  323. json_payload=_build_modbus_point_create_payload(payload),
  324. )
  325. def create_s7_device(project_key: str, payload: dict[str, Any]) -> dict[str, Any]:
  326. return _request_collector(
  327. project_key,
  328. "POST",
  329. S7_SPEC.create_device_path,
  330. json_payload=_build_s7_device_create_payload(payload),
  331. )
  332. def create_s7_point(project_key: str, payload: dict[str, Any]) -> dict[str, Any]:
  333. return _request_collector(
  334. project_key,
  335. "POST",
  336. S7_SPEC.create_point_path,
  337. json_payload=_build_s7_point_create_payload(payload),
  338. )
  339. def create_bacnet_device(project_key: str, payload: dict[str, Any]) -> dict[str, Any]:
  340. return _request_collector(
  341. project_key,
  342. "POST",
  343. BACNET_SPEC.create_device_path,
  344. json_payload=_build_bacnet_device_create_payload(payload),
  345. )
  346. def create_bacnet_point(project_key: str, payload: dict[str, Any]) -> dict[str, Any]:
  347. return _request_collector(
  348. project_key,
  349. "POST",
  350. BACNET_SPEC.create_point_path,
  351. json_payload=_build_bacnet_point_create_payload(payload),
  352. )
  353. def create_modbus_devices(project_key: str, devices: list[dict[str, Any]]) -> dict[str, Any]:
  354. return _create_devices_batch(
  355. project_key,
  356. devices,
  357. create_path=MODBUS_SPEC.create_device_path,
  358. build_payload=_build_modbus_device_create_payload,
  359. match_device=_match_modbus_device,
  360. )
  361. def create_modbus_points(project_key: str, points: list[dict[str, Any]]) -> dict[str, Any]:
  362. return _create_points_batch(
  363. project_key,
  364. points,
  365. create_path=MODBUS_SPEC.create_point_path,
  366. build_payload=_build_modbus_point_create_payload,
  367. )
  368. def create_s7_devices(project_key: str, devices: list[dict[str, Any]]) -> dict[str, Any]:
  369. return _create_devices_batch(
  370. project_key,
  371. devices,
  372. create_path=S7_SPEC.create_device_path,
  373. build_payload=_build_s7_device_create_payload,
  374. match_device=_match_s7_device,
  375. )
  376. def create_s7_points(project_key: str, points: list[dict[str, Any]]) -> dict[str, Any]:
  377. return _create_points_batch(
  378. project_key,
  379. points,
  380. create_path=S7_SPEC.create_point_path,
  381. build_payload=_build_s7_point_create_payload,
  382. )
  383. def create_bacnet_devices(project_key: str, devices: list[dict[str, Any]]) -> dict[str, Any]:
  384. return _create_devices_batch(
  385. project_key,
  386. devices,
  387. create_path=BACNET_SPEC.create_device_path,
  388. build_payload=_build_bacnet_device_create_payload,
  389. match_device=_match_bacnet_device,
  390. )
  391. def create_bacnet_points(project_key: str, points: list[dict[str, Any]]) -> dict[str, Any]:
  392. return _create_points_batch(
  393. project_key,
  394. points,
  395. create_path=BACNET_SPEC.create_point_path,
  396. build_payload=_build_bacnet_point_create_payload,
  397. )
  398. def _create_devices_batch(
  399. project_key: str,
  400. devices: list[dict[str, Any]],
  401. *,
  402. create_path: str,
  403. build_payload: Any,
  404. match_device: Any,
  405. ) -> dict[str, Any]:
  406. results: list[dict[str, Any]] = []
  407. errors: list[dict[str, Any]] = []
  408. created: list[tuple[dict[str, Any], dict[str, Any]]] = []
  409. for index, item in enumerate(devices):
  410. if not isinstance(item, dict):
  411. errors.append({"index": index, "stage": "create_device", "error": "device item must be an object"})
  412. continue
  413. result: dict[str, Any] = {
  414. "index": index,
  415. "name": item.get("name"),
  416. "device_id": None,
  417. "create_response": None,
  418. "matched_device": None,
  419. }
  420. results.append(result)
  421. try:
  422. request_payload = build_payload(item)
  423. response = _request_collector(project_key, "POST", create_path, json_payload=request_payload)
  424. except Exception as exc:
  425. errors.append({"index": index, "name": item.get("name"), "stage": "create_device", "error": str(exc)})
  426. continue
  427. result["create_response"] = response
  428. if _success_response(response):
  429. created.append((request_payload, result))
  430. else:
  431. errors.append(
  432. {
  433. "index": index,
  434. "name": item.get("name"),
  435. "stage": "create_device",
  436. "error": response.get("state_info") or response.get("msg") or str(response),
  437. }
  438. )
  439. if created:
  440. try:
  441. device_list = list_devices(project_key, num_points=False)
  442. flattened_devices = _flatten_device_tree(device_list.get("devices", []))
  443. except Exception as exc:
  444. for _, result in created:
  445. errors.append(
  446. {
  447. "index": result["index"],
  448. "name": result.get("name"),
  449. "stage": "match_device",
  450. "error": str(exc),
  451. }
  452. )
  453. else:
  454. for request_payload, result in created:
  455. candidates = [item for item in flattened_devices if match_device(item, request_payload)]
  456. candidates = [item for item in candidates if _safe_int(item.get("id")) is not None]
  457. if not candidates:
  458. errors.append(
  459. {
  460. "index": result["index"],
  461. "name": result.get("name"),
  462. "stage": "match_device",
  463. "error": "created device was not found in device list",
  464. }
  465. )
  466. continue
  467. selected = max(candidates, key=lambda item: _safe_int(item.get("id")) or 0)
  468. result["device_id"] = _safe_int(selected.get("id"))
  469. result["matched_device"] = selected
  470. matched = sum(1 for item in results if item.get("device_id") is not None)
  471. return {
  472. "state": 0 if not errors else 1,
  473. "summary": {
  474. "total": len(devices),
  475. "created": len(created),
  476. "matched": matched,
  477. "failed": len(errors),
  478. },
  479. "results": results,
  480. "errors": errors,
  481. }
  482. def _create_points_batch(
  483. project_key: str,
  484. points: list[dict[str, Any]],
  485. *,
  486. create_path: str,
  487. build_payload: Any,
  488. ) -> dict[str, Any]:
  489. results: list[dict[str, Any]] = []
  490. errors: list[dict[str, Any]] = []
  491. for index, item in enumerate(points):
  492. if not isinstance(item, dict):
  493. errors.append({"index": index, "stage": "create_point", "error": "point item must be an object"})
  494. continue
  495. result: dict[str, Any] = {
  496. "index": index,
  497. "name": item.get("name"),
  498. "device_id": item.get("device_id"),
  499. "response": None,
  500. }
  501. results.append(result)
  502. try:
  503. request_payload = build_payload(item)
  504. response = _request_collector(project_key, "POST", create_path, json_payload=request_payload)
  505. except Exception as exc:
  506. errors.append({"index": index, "name": item.get("name"), "stage": "create_point", "error": str(exc)})
  507. continue
  508. result["response"] = response
  509. if not _success_response(response):
  510. errors.append(
  511. {
  512. "index": index,
  513. "name": item.get("name"),
  514. "stage": "create_point",
  515. "error": response.get("state_info") or response.get("msg") or str(response),
  516. }
  517. )
  518. success = sum(1 for item in results if isinstance(item.get("response"), dict) and _success_response(item["response"]))
  519. return {
  520. "state": 0 if not errors else 1,
  521. "summary": {
  522. "total": len(points),
  523. "success": success,
  524. "failed": len(errors),
  525. },
  526. "results": results,
  527. "errors": errors,
  528. }
  529. def _flatten_device_tree(items: list[Any]) -> list[dict[str, Any]]:
  530. flattened: list[dict[str, Any]] = []
  531. for item in items:
  532. if not isinstance(item, dict):
  533. continue
  534. flattened.append(item)
  535. groups = item.get("groups", [])
  536. if isinstance(groups, list):
  537. flattened.extend(_flatten_device_tree(groups))
  538. return flattened
  539. def _match_modbus_device(device: dict[str, Any], expected: dict[str, Any]) -> bool:
  540. if device.get("type") != "modbus":
  541. return False
  542. if not _same_text(device.get("name"), expected.get("name")):
  543. return False
  544. if not _same_int(device.get("device_type"), expected.get("device_type")):
  545. return False
  546. if not _same_int(device.get("slave_id"), expected.get("slave_id")):
  547. return False
  548. if not _same_int(device.get("group_id"), expected.get("group_id")):
  549. return False
  550. if not _same_int(device.get("address_offset"), expected.get("address_offset")):
  551. return False
  552. if _safe_int(expected.get("device_type")) == 2:
  553. return _same_text(device.get("serial_port"), expected.get("serial_port"))
  554. return _same_text(device.get("ip"), expected.get("ip")) and _same_int(device.get("port"), expected.get("port"))
  555. def _match_s7_device(device: dict[str, Any], expected: dict[str, Any]) -> bool:
  556. if device.get("type") != "s7":
  557. return False
  558. actual_group_id = device.get("device_group_id", device.get("group_id"))
  559. return (
  560. _same_text(device.get("name"), expected.get("name"))
  561. and _same_text(device.get("ip"), expected.get("ip"))
  562. and _same_int(device.get("port"), expected.get("port"))
  563. and _same_int(device.get("rock"), expected.get("rock"))
  564. and _same_int(device.get("slot"), expected.get("slot"))
  565. and _same_int(device.get("device_type"), expected.get("device_type"))
  566. and _same_text(device.get("tsap_conn_type"), expected.get("tsap_conn_type"))
  567. and _same_int(actual_group_id, expected.get("device_group_id"))
  568. )
  569. def _match_bacnet_device(device: dict[str, Any], expected: dict[str, Any]) -> bool:
  570. if device.get("type") != "bacnet":
  571. return False
  572. actual_group_id = device.get("group_id", device.get("device_group_id"))
  573. return (
  574. _same_text(device.get("name"), expected.get("name"))
  575. and _same_text(device.get("ip"), expected.get("ip"))
  576. and _same_int(device.get("port"), expected.get("port"))
  577. and _same_int(device.get("device_type"), expected.get("device_type"))
  578. and _same_int(device.get("bacnet_device_id"), expected.get("bacnet_device_id"))
  579. and _same_int(device.get("bacnet_net"), expected.get("bacnet_net"))
  580. and _same_int(actual_group_id, expected.get("group_id"))
  581. )
  582. def _same_text(left: Any, right: Any) -> bool:
  583. return str(left or "").strip() == str(right or "").strip()
  584. def _same_int(left: Any, right: Any) -> bool:
  585. normalized_left = _safe_int(left)
  586. normalized_right = _safe_int(right)
  587. return normalized_left is not None and normalized_left == normalized_right
  588. def _safe_int(value: Any) -> int | None:
  589. try:
  590. return int(value)
  591. except Exception:
  592. return None
  593. def edit_modbus_device(project_key: str, payload: dict[str, Any]) -> dict[str, Any]:
  594. return _request_collector(
  595. project_key,
  596. "POST",
  597. "/api/collector/modbus/device/edit",
  598. json_payload=_merge_defaults(
  599. {
  600. "timeout": 3,
  601. "is_persistent": True,
  602. "device_group_id": 0,
  603. "alarm_interval": 90,
  604. "collect_interval": 5,
  605. "address_offset": 0,
  606. "retry_times": 0,
  607. "mode": 0,
  608. },
  609. _normalize_modbus_device_edit_payload(payload),
  610. ),
  611. )
  612. def edit_modbus_point(project_key: str, payload: dict[str, Any]) -> dict[str, Any]:
  613. return _request_collector(
  614. project_key,
  615. "POST",
  616. "/api/collector/modbus/point/edit_collect_point",
  617. json_payload=_merge_defaults(
  618. MODBUS_SPEC.point_defaults,
  619. _normalize_modbus_point_edit_payload(payload),
  620. ),
  621. )
  622. def edit_s7_device(project_key: str, payload: dict[str, Any]) -> dict[str, Any]:
  623. return _request_collector(
  624. project_key,
  625. "POST",
  626. "/api/collector/s7/device/update",
  627. json_payload=_merge_defaults(
  628. S7_SPEC.device_defaults,
  629. _normalize_s7_device_edit_payload(payload),
  630. ),
  631. )
  632. def edit_s7_point(project_key: str, payload: dict[str, Any]) -> dict[str, Any]:
  633. return _request_collector(
  634. project_key,
  635. "POST",
  636. "/api/collector/s7/point/update",
  637. json_payload=_merge_defaults(
  638. S7_SPEC.point_defaults,
  639. _normalize_s7_point_edit_payload(payload),
  640. ),
  641. )
  642. def edit_bacnet_device(project_key: str, payload: dict[str, Any]) -> dict[str, Any]:
  643. return _request_collector(
  644. project_key,
  645. "POST",
  646. "/api/collector/bacnet/device/edit",
  647. json_payload=_merge_defaults(
  648. {
  649. "type": 1,
  650. "port": 47808,
  651. "net": 0,
  652. "asp_ip": "",
  653. "timeout": 3,
  654. "is_persistent": False,
  655. "device_group_id": 0,
  656. "alarm_interval": 90,
  657. "collect_interval": 5,
  658. },
  659. _normalize_bacnet_device_edit_payload(payload),
  660. ),
  661. )
  662. def edit_bacnet_point(project_key: str, payload: dict[str, Any]) -> dict[str, Any]:
  663. return _request_collector(
  664. project_key,
  665. "POST",
  666. "/api/collector/bacnet/point/edit",
  667. json_payload=_normalize_bacnet_point_edit_payload(payload),
  668. )
  669. def list_devices(project_key: str, num_points: bool = False) -> dict[str, Any]:
  670. num_points_text = "true" if num_points else "false"
  671. return _request_collector(
  672. project_key,
  673. "GET",
  674. f"/api/collector/device?num_points={num_points_text}",
  675. )
  676. def connect_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=2,
  682. )
  683. def disconnect_device(project_key: str, device_id: int, device_type: str = "modbus") -> dict[str, Any]:
  684. return _set_device_connect_status(
  685. project_key,
  686. device_id=device_id,
  687. device_type=device_type,
  688. status=1,
  689. )
  690. def list_device_points(
  691. project_key: str,
  692. device_id: int,
  693. device_type: str = "modbus",
  694. group_id: int = 0,
  695. ) -> dict[str, Any]:
  696. return _request_collector(
  697. project_key,
  698. "POST",
  699. "/api/collector/common/device/get_collect_point",
  700. json_payload={
  701. "id": _normalize_positive_int(device_id, "device_id"),
  702. "type": _normalize_device_type(device_type),
  703. "group_id": _normalize_non_negative_int(group_id, "group_id"),
  704. },
  705. )
  706. def _set_device_connect_status(
  707. project_key: str,
  708. *,
  709. device_id: int,
  710. device_type: str,
  711. status: int,
  712. ) -> dict[str, Any]:
  713. return _request_collector(
  714. project_key,
  715. "POST",
  716. "/api/collector/common/device/set_connect_status",
  717. json_payload={
  718. "id": _normalize_positive_int(device_id, "device_id"),
  719. "type": _normalize_device_type(device_type),
  720. "status": status,
  721. },
  722. )
  723. def _normalize_device_type(device_type: str) -> str:
  724. normalized = str(device_type or "").strip()
  725. if not normalized:
  726. raise ValueError("device_type is required")
  727. return normalized
  728. def _normalize_s7_tsap_conn_type(value: Any) -> str:
  729. normalized = str(value or "").strip().upper()
  730. if normalized not in {"PG", "OP", "BASIC"}:
  731. raise ValueError("payload.tsap_conn_type must be one of PG, OP, BASIC")
  732. return normalized
  733. def _normalize_positive_int(value: int, field_name: str) -> int:
  734. try:
  735. normalized = int(value)
  736. except Exception as exc:
  737. raise ValueError(f"{field_name} must be a positive integer") from exc
  738. if normalized <= 0:
  739. raise ValueError(f"{field_name} must be a positive integer")
  740. return normalized
  741. def _normalize_non_negative_int(value: int, field_name: str) -> int:
  742. try:
  743. normalized = int(value)
  744. except Exception as exc:
  745. raise ValueError(f"{field_name} must be a non-negative integer") from exc
  746. if normalized < 0:
  747. raise ValueError(f"{field_name} must be a non-negative integer")
  748. return normalized
  749. def _normalize_bacnet_id(value: Any, field_name: str) -> int:
  750. normalized = _normalize_non_negative_int(value, field_name)
  751. if normalized > BACNET_OBJECT_ID_MAX:
  752. raise ValueError(f"{field_name} must be between 0 and {BACNET_OBJECT_ID_MAX}")
  753. return normalized