| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719 |
- from __future__ import annotations
- from typing import Any
- from .auth import find_project_config, resolve_project_token
- from .http_client import request_json
- from .protocols import MODBUS_SPEC, S7_SPEC
- from .protocols.modbus import MODBUS_POINT_TYPE_ALIASES, MODBUS_REGISTER_TYPE_ALIASES
- from .protocols.s7 import S7_POINT_TYPE_ALIASES, S7_REGISTER_TYPE_ALIASES
- def _merge_defaults(defaults: dict[str, Any], payload: dict[str, Any]) -> dict[str, Any]:
- merged = dict(defaults)
- merged.update(payload)
- return merged
- def _success_response(response: dict[str, Any]) -> bool:
- return response.get("state") == 0
- def _require_non_empty_text(payload: dict[str, Any], field_name: str) -> str:
- value = str(payload.get(field_name) or "").strip()
- if not value:
- raise ValueError(f"payload.{field_name} is required")
- return value
- def _require_present(payload: dict[str, Any], field_name: str) -> Any:
- if field_name not in payload or payload.get(field_name) is None:
- raise ValueError(f"payload.{field_name} is required")
- return payload[field_name]
- def _normalize_modbus_device_payload(payload: dict[str, Any]) -> dict[str, Any]:
- normalized = dict(payload)
- normalized["name"] = _require_non_empty_text(normalized, "name")
- normalized["ip"] = _require_non_empty_text(normalized, "ip")
- _require_present(normalized, "device_type")
- _require_present(normalized, "port")
- _require_present(normalized, "slave_id")
- _require_present(normalized, "word_order")
- _require_present(normalized, "byte_order")
- address_base = _require_present(normalized, "address_base")
- normalized["address_offset"] = address_base
- normalized.pop("address_base", None)
- return normalized
- def _normalize_modbus_device_edit_payload(payload: dict[str, Any]) -> dict[str, Any]:
- normalized = dict(payload)
- normalized["ori_id"] = _normalize_positive_int(_require_present(normalized, "ori_id"), "payload.ori_id")
- normalized["name"] = _require_non_empty_text(normalized, "name")
- if normalized.get("type") is None:
- normalized["type"] = _require_present(normalized, "device_type")
- try:
- normalized["type"] = int(normalized["type"])
- except Exception as exc:
- raise ValueError("payload.type must be one of 1, 2, 3, 4, 5") from exc
- if normalized["type"] not in {1, 2, 3, 4, 5}:
- raise ValueError("payload.type must be one of 1, 2, 3, 4, 5")
- normalized.pop("device_type", None)
- if normalized["type"] == 2:
- normalized["serial_port"] = _require_non_empty_text(normalized, "serial_port")
- else:
- normalized["ip"] = _require_non_empty_text(normalized, "ip")
- _require_present(normalized, "port")
- _require_present(normalized, "slave_id")
- _require_present(normalized, "word_order")
- _require_present(normalized, "byte_order")
- if "address_base" in normalized:
- normalized["address_offset"] = normalized["address_base"]
- normalized.pop("address_base", None)
- if "group_id" in normalized and "device_group_id" not in normalized:
- normalized["device_group_id"] = normalized["group_id"]
- normalized.pop("group_id", None)
- return normalized
- def _normalize_modbus_point_payload(payload: dict[str, Any], *, require_device_id: bool = True) -> dict[str, Any]:
- normalized = dict(payload)
- normalized["name"] = _require_non_empty_text(normalized, "name")
- _require_present(normalized, "address")
- if require_device_id:
- normalized["device_id"] = _normalize_positive_int(
- _require_present(normalized, "device_id"),
- "payload.device_id",
- )
- raw_type = _require_non_empty_text(normalized, "type")
- normalized_type = MODBUS_POINT_TYPE_ALIASES.get(raw_type)
- if normalized_type is None:
- normalized_type = MODBUS_POINT_TYPE_ALIASES.get(raw_type.lower())
- if normalized_type is None:
- raise ValueError(
- "payload.type is invalid; use one of bool, int16, uint16, int32, "
- "uint32, int64, uint64, float32, float64, or a documented alias "
- "such as SHORT, WORD, LONG, DWORD, FLOAT, DOUBLE"
- )
- normalized["type"] = normalized_type
- if normalized.get("func_code") is None:
- register_type = _require_non_empty_text(normalized, "register_type")
- func_code = MODBUS_REGISTER_TYPE_ALIASES.get(register_type)
- if func_code is None:
- func_code = MODBUS_REGISTER_TYPE_ALIASES.get(register_type.lower())
- if func_code is None:
- raise ValueError(
- "payload.register_type is invalid; use coil, discrete_input, "
- "holding_register, input_register, or func_code 1/2/3/4"
- )
- normalized["func_code"] = func_code
- else:
- try:
- normalized["func_code"] = int(str(normalized["func_code"]).strip())
- except Exception as exc:
- raise ValueError("payload.func_code must be one of 1, 2, 3, 4") from exc
- if normalized["func_code"] not in {1, 2, 3, 4}:
- raise ValueError("payload.func_code must be one of 1, 2, 3, 4")
- if "register_type" in normalized:
- normalized.pop("register_type")
- return normalized
- def _normalize_modbus_point_edit_payload(payload: dict[str, Any]) -> dict[str, Any]:
- _require_present(payload, "ori_id")
- normalized = _normalize_modbus_point_payload(payload, require_device_id=False)
- normalized["ori_id"] = _normalize_positive_int(_require_present(normalized, "ori_id"), "payload.ori_id")
- return normalized
- def _normalize_s7_device_payload(payload: dict[str, Any]) -> dict[str, Any]:
- normalized = dict(payload)
- normalized["name"] = _require_non_empty_text(normalized, "name")
- normalized["ip"] = _require_non_empty_text(normalized, "ip")
- _require_present(normalized, "rock")
- _require_present(normalized, "slot")
- if "group_id" in normalized and "device_group_id" not in normalized:
- normalized["device_group_id"] = normalized["group_id"]
- normalized.pop("group_id", None)
- if "tsap_conn_type" in normalized and normalized["tsap_conn_type"] is not None:
- normalized["tsap_conn_type"] = _normalize_s7_tsap_conn_type(normalized["tsap_conn_type"])
- elif int(normalized.get("device_type", 1)) == 3:
- normalized["tsap_conn_type"] = "OP"
- else:
- normalized["tsap_conn_type"] = "PG"
- return normalized
- def _normalize_s7_device_edit_payload(payload: dict[str, Any]) -> dict[str, Any]:
- normalized = _normalize_s7_device_payload(payload)
- if normalized.get("id") is None:
- normalized["id"] = _require_present(normalized, "ori_id")
- normalized["id"] = _normalize_positive_int(normalized["id"], "payload.id")
- normalized.pop("ori_id", None)
- return normalized
- def _normalize_s7_point_payload(payload: dict[str, Any], *, require_device_id: bool = True) -> dict[str, Any]:
- normalized = dict(payload)
- normalized["name"] = _require_non_empty_text(normalized, "name")
- normalized["address"] = _require_non_empty_text(normalized, "address")
- if require_device_id:
- normalized["device_id"] = _normalize_positive_int(
- _require_present(normalized, "device_id"),
- "payload.device_id",
- )
- raw_type = normalized.get("data_type")
- if raw_type is None:
- raw_type = _require_non_empty_text(normalized, "type")
- else:
- raw_type = str(raw_type).strip()
- if not raw_type:
- raise ValueError("payload.data_type is required")
- normalized_type = S7_POINT_TYPE_ALIASES.get(raw_type)
- if normalized_type is None:
- normalized_type = S7_POINT_TYPE_ALIASES.get(raw_type.lower())
- if normalized_type is None:
- raise ValueError(
- "payload.data_type is invalid; use one of bool, uint8, int8, "
- "uint16, int16, uint32, int32, float32, float64, or a documented "
- "alias such as BOOL, BYTE, SINT, WORD, INT, DWORD, DINT, REAL, LREAL"
- )
- normalized["data_type"] = normalized_type
- normalized.pop("type", None)
- if normalized.get("register_type") is None:
- register_area = _require_non_empty_text(normalized, "register_area")
- register_type = S7_REGISTER_TYPE_ALIASES.get(register_area)
- if register_type is None:
- register_type = S7_REGISTER_TYPE_ALIASES.get(register_area.lower())
- if register_type is None:
- raise ValueError("payload.register_area is invalid; use I, Q, M, DB, V, AI, or register_type 1..6")
- normalized["register_type"] = register_type
- else:
- try:
- normalized["register_type"] = int(str(normalized["register_type"]).strip())
- except Exception as exc:
- raise ValueError("payload.register_type must be one of 1, 2, 3, 4, 5, 6") from exc
- if normalized["register_type"] not in {1, 2, 3, 4, 5, 6}:
- raise ValueError("payload.register_type must be one of 1, 2, 3, 4, 5, 6")
- normalized.pop("register_area", None)
- if "group_id" in normalized and "group_Id" not in normalized:
- normalized["group_Id"] = normalized["group_id"]
- normalized.pop("group_id", None)
- return normalized
- def _normalize_s7_point_edit_payload(payload: dict[str, Any]) -> dict[str, Any]:
- normalized = _normalize_s7_point_payload(payload)
- if normalized.get("id") is None:
- normalized["id"] = _require_present(normalized, "ori_id")
- normalized["id"] = _normalize_positive_int(normalized["id"], "payload.id")
- normalized.pop("ori_id", None)
- return normalized
- def _request_collector(
- project_key: str,
- method: str,
- path: str,
- *,
- json_payload: dict[str, Any] | None = None,
- ) -> dict[str, Any]:
- project = find_project_config(project_key)
- authorization = resolve_project_token(project)
- response_payload = request_json(
- method,
- f"{project['data_collector_base_url']}{path}",
- authorization,
- json_payload=json_payload,
- )
- if not isinstance(response_payload, dict):
- raise ValueError(f"collector API returned invalid payload for {path}: {response_payload}")
- return response_payload
- def _build_modbus_device_create_payload(payload: dict[str, Any]) -> dict[str, Any]:
- return _merge_defaults(
- MODBUS_SPEC.device_defaults,
- _normalize_modbus_device_payload(payload),
- )
- def _build_modbus_point_create_payload(payload: dict[str, Any]) -> dict[str, Any]:
- return _merge_defaults(
- MODBUS_SPEC.point_defaults,
- _normalize_modbus_point_payload(payload),
- )
- def _build_s7_device_create_payload(payload: dict[str, Any]) -> dict[str, Any]:
- return _merge_defaults(
- S7_SPEC.device_defaults,
- _normalize_s7_device_payload(payload),
- )
- def _build_s7_point_create_payload(payload: dict[str, Any]) -> dict[str, Any]:
- return _merge_defaults(
- S7_SPEC.point_defaults,
- _normalize_s7_point_payload(payload),
- )
- def create_modbus_device(project_key: str, payload: dict[str, Any]) -> dict[str, Any]:
- return _request_collector(
- project_key,
- "POST",
- MODBUS_SPEC.create_device_path,
- json_payload=_build_modbus_device_create_payload(payload),
- )
- def create_modbus_point(project_key: str, payload: dict[str, Any]) -> dict[str, Any]:
- return _request_collector(
- project_key,
- "POST",
- MODBUS_SPEC.create_point_path,
- json_payload=_build_modbus_point_create_payload(payload),
- )
- def create_s7_device(project_key: str, payload: dict[str, Any]) -> dict[str, Any]:
- return _request_collector(
- project_key,
- "POST",
- S7_SPEC.create_device_path,
- json_payload=_build_s7_device_create_payload(payload),
- )
- def create_s7_point(project_key: str, payload: dict[str, Any]) -> dict[str, Any]:
- return _request_collector(
- project_key,
- "POST",
- S7_SPEC.create_point_path,
- json_payload=_build_s7_point_create_payload(payload),
- )
- def create_modbus_devices(project_key: str, devices: list[dict[str, Any]]) -> dict[str, Any]:
- return _create_devices_batch(
- project_key,
- devices,
- create_path=MODBUS_SPEC.create_device_path,
- build_payload=_build_modbus_device_create_payload,
- match_device=_match_modbus_device,
- )
- def create_modbus_points(project_key: str, points: list[dict[str, Any]]) -> dict[str, Any]:
- return _create_points_batch(
- project_key,
- points,
- create_path=MODBUS_SPEC.create_point_path,
- build_payload=_build_modbus_point_create_payload,
- )
- def create_s7_devices(project_key: str, devices: list[dict[str, Any]]) -> dict[str, Any]:
- return _create_devices_batch(
- project_key,
- devices,
- create_path=S7_SPEC.create_device_path,
- build_payload=_build_s7_device_create_payload,
- match_device=_match_s7_device,
- )
- def create_s7_points(project_key: str, points: list[dict[str, Any]]) -> dict[str, Any]:
- return _create_points_batch(
- project_key,
- points,
- create_path=S7_SPEC.create_point_path,
- build_payload=_build_s7_point_create_payload,
- )
- def _create_devices_batch(
- project_key: str,
- devices: list[dict[str, Any]],
- *,
- create_path: str,
- build_payload: Any,
- match_device: Any,
- ) -> dict[str, Any]:
- results: list[dict[str, Any]] = []
- errors: list[dict[str, Any]] = []
- created: list[tuple[dict[str, Any], dict[str, Any]]] = []
- for index, item in enumerate(devices):
- if not isinstance(item, dict):
- errors.append({"index": index, "stage": "create_device", "error": "device item must be an object"})
- continue
- result: dict[str, Any] = {
- "index": index,
- "name": item.get("name"),
- "device_id": None,
- "create_response": None,
- "matched_device": None,
- }
- results.append(result)
- try:
- request_payload = build_payload(item)
- response = _request_collector(project_key, "POST", create_path, json_payload=request_payload)
- except Exception as exc:
- errors.append({"index": index, "name": item.get("name"), "stage": "create_device", "error": str(exc)})
- continue
- result["create_response"] = response
- if _success_response(response):
- created.append((request_payload, result))
- else:
- errors.append(
- {
- "index": index,
- "name": item.get("name"),
- "stage": "create_device",
- "error": response.get("state_info") or response.get("msg") or str(response),
- }
- )
- if created:
- try:
- device_list = list_devices(project_key, num_points=False)
- flattened_devices = _flatten_device_tree(device_list.get("devices", []))
- except Exception as exc:
- for _, result in created:
- errors.append(
- {
- "index": result["index"],
- "name": result.get("name"),
- "stage": "match_device",
- "error": str(exc),
- }
- )
- else:
- for request_payload, result in created:
- candidates = [item for item in flattened_devices if match_device(item, request_payload)]
- candidates = [item for item in candidates if _safe_int(item.get("id")) is not None]
- if not candidates:
- errors.append(
- {
- "index": result["index"],
- "name": result.get("name"),
- "stage": "match_device",
- "error": "created device was not found in device list",
- }
- )
- continue
- selected = max(candidates, key=lambda item: _safe_int(item.get("id")) or 0)
- result["device_id"] = _safe_int(selected.get("id"))
- result["matched_device"] = selected
- matched = sum(1 for item in results if item.get("device_id") is not None)
- return {
- "state": 0 if not errors else 1,
- "summary": {
- "total": len(devices),
- "created": len(created),
- "matched": matched,
- "failed": len(errors),
- },
- "results": results,
- "errors": errors,
- }
- def _create_points_batch(
- project_key: str,
- points: list[dict[str, Any]],
- *,
- create_path: str,
- build_payload: Any,
- ) -> dict[str, Any]:
- results: list[dict[str, Any]] = []
- errors: list[dict[str, Any]] = []
- for index, item in enumerate(points):
- if not isinstance(item, dict):
- errors.append({"index": index, "stage": "create_point", "error": "point item must be an object"})
- continue
- result: dict[str, Any] = {
- "index": index,
- "name": item.get("name"),
- "device_id": item.get("device_id"),
- "response": None,
- }
- results.append(result)
- try:
- request_payload = build_payload(item)
- response = _request_collector(project_key, "POST", create_path, json_payload=request_payload)
- except Exception as exc:
- errors.append({"index": index, "name": item.get("name"), "stage": "create_point", "error": str(exc)})
- continue
- result["response"] = response
- if not _success_response(response):
- errors.append(
- {
- "index": index,
- "name": item.get("name"),
- "stage": "create_point",
- "error": response.get("state_info") or response.get("msg") or str(response),
- }
- )
- success = sum(1 for item in results if isinstance(item.get("response"), dict) and _success_response(item["response"]))
- return {
- "state": 0 if not errors else 1,
- "summary": {
- "total": len(points),
- "success": success,
- "failed": len(errors),
- },
- "results": results,
- "errors": errors,
- }
- def _flatten_device_tree(items: list[Any]) -> list[dict[str, Any]]:
- flattened: list[dict[str, Any]] = []
- for item in items:
- if not isinstance(item, dict):
- continue
- flattened.append(item)
- groups = item.get("groups", [])
- if isinstance(groups, list):
- flattened.extend(_flatten_device_tree(groups))
- return flattened
- def _match_modbus_device(device: dict[str, Any], expected: dict[str, Any]) -> bool:
- if device.get("type") != "modbus":
- return False
- if not _same_text(device.get("name"), expected.get("name")):
- return False
- if not _same_int(device.get("device_type"), expected.get("device_type")):
- return False
- if not _same_int(device.get("slave_id"), expected.get("slave_id")):
- return False
- if not _same_int(device.get("group_id"), expected.get("group_id")):
- return False
- if not _same_int(device.get("address_offset"), expected.get("address_offset")):
- return False
- if _safe_int(expected.get("device_type")) == 2:
- return _same_text(device.get("serial_port"), expected.get("serial_port"))
- return _same_text(device.get("ip"), expected.get("ip")) and _same_int(device.get("port"), expected.get("port"))
- def _match_s7_device(device: dict[str, Any], expected: dict[str, Any]) -> bool:
- if device.get("type") != "s7":
- return False
- actual_group_id = device.get("device_group_id", device.get("group_id"))
- return (
- _same_text(device.get("name"), expected.get("name"))
- and _same_text(device.get("ip"), expected.get("ip"))
- and _same_int(device.get("port"), expected.get("port"))
- and _same_int(device.get("rock"), expected.get("rock"))
- and _same_int(device.get("slot"), expected.get("slot"))
- and _same_int(device.get("device_type"), expected.get("device_type"))
- and _same_text(device.get("tsap_conn_type"), expected.get("tsap_conn_type"))
- and _same_int(actual_group_id, expected.get("device_group_id"))
- )
- def _same_text(left: Any, right: Any) -> bool:
- return str(left or "").strip() == str(right or "").strip()
- def _same_int(left: Any, right: Any) -> bool:
- normalized_left = _safe_int(left)
- normalized_right = _safe_int(right)
- return normalized_left is not None and normalized_left == normalized_right
- def _safe_int(value: Any) -> int | None:
- try:
- return int(value)
- except Exception:
- return None
- def edit_modbus_device(project_key: str, payload: dict[str, Any]) -> dict[str, Any]:
- return _request_collector(
- project_key,
- "POST",
- "/api/collector/modbus/device/edit",
- json_payload=_merge_defaults(
- {
- "timeout": 3,
- "is_persistent": True,
- "device_group_id": 0,
- "alarm_interval": 90,
- "collect_interval": 5,
- "address_offset": 0,
- "retry_times": 0,
- "mode": 0,
- },
- _normalize_modbus_device_edit_payload(payload),
- ),
- )
- def edit_modbus_point(project_key: str, payload: dict[str, Any]) -> dict[str, Any]:
- return _request_collector(
- project_key,
- "POST",
- "/api/collector/modbus/point/edit_collect_point",
- json_payload=_merge_defaults(
- MODBUS_SPEC.point_defaults,
- _normalize_modbus_point_edit_payload(payload),
- ),
- )
- def edit_s7_device(project_key: str, payload: dict[str, Any]) -> dict[str, Any]:
- return _request_collector(
- project_key,
- "POST",
- "/api/collector/s7/device/update",
- json_payload=_merge_defaults(
- S7_SPEC.device_defaults,
- _normalize_s7_device_edit_payload(payload),
- ),
- )
- def edit_s7_point(project_key: str, payload: dict[str, Any]) -> dict[str, Any]:
- return _request_collector(
- project_key,
- "POST",
- "/api/collector/s7/point/update",
- json_payload=_merge_defaults(
- S7_SPEC.point_defaults,
- _normalize_s7_point_edit_payload(payload),
- ),
- )
- def list_devices(project_key: str, num_points: bool = False) -> dict[str, Any]:
- num_points_text = "true" if num_points else "false"
- return _request_collector(
- project_key,
- "GET",
- f"/api/collector/device?num_points={num_points_text}",
- )
- def connect_device(project_key: str, device_id: int, device_type: str = "modbus") -> dict[str, Any]:
- return _set_device_connect_status(
- project_key,
- device_id=device_id,
- device_type=device_type,
- status=2,
- )
- def disconnect_device(project_key: str, device_id: int, device_type: str = "modbus") -> dict[str, Any]:
- return _set_device_connect_status(
- project_key,
- device_id=device_id,
- device_type=device_type,
- status=1,
- )
- def list_device_points(
- project_key: str,
- device_id: int,
- device_type: str = "modbus",
- group_id: int = 0,
- ) -> dict[str, Any]:
- return _request_collector(
- project_key,
- "POST",
- "/api/collector/common/device/get_collect_point",
- json_payload={
- "id": _normalize_positive_int(device_id, "device_id"),
- "type": _normalize_device_type(device_type),
- "group_id": _normalize_non_negative_int(group_id, "group_id"),
- },
- )
- def _set_device_connect_status(
- project_key: str,
- *,
- device_id: int,
- device_type: str,
- status: int,
- ) -> dict[str, Any]:
- return _request_collector(
- project_key,
- "POST",
- "/api/collector/common/device/set_connect_status",
- json_payload={
- "id": _normalize_positive_int(device_id, "device_id"),
- "type": _normalize_device_type(device_type),
- "status": status,
- },
- )
- def _normalize_device_type(device_type: str) -> str:
- normalized = str(device_type or "").strip()
- if not normalized:
- raise ValueError("device_type is required")
- return normalized
- def _normalize_s7_tsap_conn_type(value: Any) -> str:
- normalized = str(value or "").strip().upper()
- if normalized not in {"PG", "OP", "BASIC"}:
- raise ValueError("payload.tsap_conn_type must be one of PG, OP, BASIC")
- return normalized
- def _normalize_positive_int(value: int, field_name: str) -> int:
- try:
- normalized = int(value)
- except Exception as exc:
- raise ValueError(f"{field_name} must be a positive integer") from exc
- if normalized <= 0:
- raise ValueError(f"{field_name} must be a positive integer")
- return normalized
- def _normalize_non_negative_int(value: int, field_name: str) -> int:
- try:
- normalized = int(value)
- except Exception as exc:
- raise ValueError(f"{field_name} must be a non-negative integer") from exc
- if normalized < 0:
- raise ValueError(f"{field_name} must be a non-negative integer")
- return normalized
|