collector_api.py 32 KB

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