collector_api.py 32 KB

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