from __future__ import annotations import asyncio import logging import os import socket import threading import warnings from contextlib import closing from contextlib import contextmanager from typing import Any BAC0_REGISTERED_WARNING = r"object type .* for vendor identifier 842 already registered.*" warnings.filterwarnings("ignore", message=BAC0_REGISTERED_WARNING) import BAC0 import BAC0.scripts.Base as bac0_base import bacpypes3.app as bacpypes_app from bacpypes3.apdu import ErrorRejectAbortNack from bacpypes3.basetypes import PropertyIdentifier from bacpypes3.comm import bind from bacpypes3.ipv4 import IPv4DatagramServer from bacpypes3.ipv4.bvll import BVLLCodec from bacpypes3.ipv4.service import BIPNormal, UDPMultiplexer from bacpypes3.pdu import Address from bacpypes3.primitivedata import ObjectIdentifier from app.response import response_payload from app.schemas.bacnet import ( BacnetPointReadRequest, BacnetPointSearchRequest, BacnetPointSpec, bac0_object_type, normalize_object_type, ) POINT_OBJECT_TYPES = { "AnalogInput", "AnalogOutput", "AnalogValue", "BinaryInput", "BinaryOutput", "BinaryValue", "MultiStateInput", "MultiStateOutput", "MultiStateValue", } _BACNET_LOCK = threading.Lock() logger = logging.getLogger("uvicorn.error") try: BAC0.log_level("silence") except Exception: pass class BacnetCommunicationError(RuntimeError): def __init__(self, message: str) -> None: super().__init__(message) def device_payload(request: BacnetPointReadRequest | BacnetPointSearchRequest) -> dict[str, Any]: return { "device_type": "BACnet/IP", "ip": request.ip, "port": request.port, "bacnet_device_id": request.bacnet_device_id, } def bacnet_address(ip: str, port: int) -> str: return f"{ip}:{port}" async def read_property_direct( bacnet: Any, address: str, object_type: str, object_id: int, property_name: str, array_index: int | None = None, ) -> Any: logger.info( "BACnet ReadProperty target=%s object_type=%s object_id=%s property=%s array_index=%s", address, object_type, object_id, property_name, array_index, ) value = await bacnet.this_application.app.read_property( Address(address), ObjectIdentifier((bac0_object_type(object_type), object_id)), PropertyIdentifier(property_name), array_index, ) if isinstance(value, ErrorRejectAbortNack): logger.error( "BACnet ReadProperty returned error target=%s object_type=%s object_id=%s property=%s array_index=%s error=%s", address, object_type, object_id, property_name, array_index, value, ) raise BacnetCommunicationError(str(value)) return value def get_local_ip(remote_ip: str, remote_port: int) -> str: configured_ip = os.getenv("BACNET_LOCAL_IP") if configured_ip: return configured_ip with closing(socket.socket(socket.AF_INET, socket.SOCK_DGRAM)) as sock: sock.connect((remote_ip, remote_port)) return str(sock.getsockname()[0]) def get_free_udp_port() -> int: configured_port = os.getenv("BACNET_LOCAL_PORT") if configured_port: return int(configured_port) with closing(socket.socket(socket.AF_INET, socket.SOCK_DGRAM)) as sock: sock.bind(("", 0)) return int(sock.getsockname()[1]) def create_bacnet_bind_socket(port: int) -> socket.socket: bind_ip = os.getenv("BACNET_BIND_IP", "0.0.0.0") sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM) sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) if hasattr(socket, "SO_REUSEPORT"): try: sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEPORT, 1) except OSError: pass sock.bind((bind_ip, port)) logger.info("BACnet UDP socket bound bind_ip=%s bind_port=%s", bind_ip, port) return sock @contextmanager def go_style_bacnet_binding(port: int): bind_socket = create_bacnet_bind_socket(port) original_normal_link_layer = bacpypes_app.NormalLinkLayer_ipv4 original_validate_ip_address = bac0_base.validate_ip_address socket_transferred = False class GoStyleNormalLinkLayer(BIPNormal): def __init__(self, local_address, **kwargs) -> None: nonlocal socket_transferred BIPNormal.__init__(self, **kwargs) self.codec = BVLLCodec() self.multiplexer = UDPMultiplexer() self.server = IPv4DatagramServer(local_address, no_broadcast=True, bind_socket=bind_socket) socket_transferred = True bind(self, self.codec, self.multiplexer.annexJ) bind(self.multiplexer, self.server) def close(self): self.server.close() try: bacpypes_app.NormalLinkLayer_ipv4 = GoStyleNormalLinkLayer bac0_base.validate_ip_address = lambda _ip: True yield finally: bacpypes_app.NormalLinkLayer_ipv4 = original_normal_link_layer bac0_base.validate_ip_address = original_validate_ip_address if not socket_transferred: bind_socket.close() def local_bacnet_ip(remote_ip: str, remote_port: int) -> str: local_ip = get_local_ip(remote_ip, remote_port) local_mask = int(os.getenv("BACNET_LOCAL_MASK", "24")) return f"{local_ip}/{local_mask}" def json_value(value: Any) -> Any: if value is None or isinstance(value, bool | int | float | str): return value if isinstance(value, list | tuple): return [json_value(item) for item in value] if isinstance(value, dict): return {str(key): json_value(item) for key, item in value.items()} if hasattr(value, "value") and isinstance(value.value, bool | int | float | str): return value.value return str(value) async def read_property_or_none(bacnet: Any, address: str, object_type: str, object_id: int, property_name: str) -> Any: try: value = await read_property_direct(bacnet, address, object_type, object_id, property_name) except Exception: return None return json_value(value) async def read_present_value(bacnet: Any, address: str, point: BacnetPointSpec) -> Any: try: value = await read_property_direct(bacnet, address, point.object_type, point.object_id, "presentValue") except Exception as exc: raise BacnetCommunicationError( f"no response from BACnet device {address} while reading " f"{point.object_type} {point.object_id} presentValue" ) from exc return json_value(value) def ignore_cancelled_bacnet_broadcast_endpoint(loop: asyncio.AbstractEventLoop, context: dict[str, Any]) -> None: message = str(context.get("message", "")) exception = context.get("exception") if isinstance(exception, asyncio.CancelledError) and "set_broadcast_transport_protocol" in message: return if isinstance(exception, RuntimeError) and str(exception) == "no broadcast": return loop.default_exception_handler(context) def run_bacnet(coro: Any) -> Any: loop = asyncio.new_event_loop() loop.set_exception_handler(ignore_cancelled_bacnet_broadcast_endpoint) try: asyncio.set_event_loop(loop) return loop.run_until_complete(coro) finally: pending = asyncio.all_tasks(loop) for task in pending: task.cancel() if pending: loop.run_until_complete(asyncio.gather(*pending, return_exceptions=True)) loop.close() asyncio.set_event_loop(None) def split_object_identifier(value: Any) -> tuple[str, int] | None: if isinstance(value, list | tuple) and len(value) == 2: raw_type, raw_id = value try: return normalize_object_type(str(raw_type).replace("ObjectType.", "")), int(raw_id) except (TypeError, ValueError): return None text = str(value).strip().strip("()") if "," in text: raw_type, raw_id = text.split(",", 1) elif ":" in text: raw_type, raw_id = text.split(":", 1) else: return None raw_type = raw_type.strip().strip("'\"").replace("ObjectType.", "") raw_id = raw_id.strip().strip("'\"") try: return normalize_object_type(raw_type), int(raw_id) except (TypeError, ValueError): return None async def read_points_async(request: BacnetPointReadRequest) -> list[dict[str, Any]]: address = bacnet_address(request.ip, request.port) local_ip = local_bacnet_ip(request.ip, request.port) local_port = get_free_udp_port() points: list[dict[str, Any]] = [] logger.info( "BACnet read_points start target=%s local_ip=%s local_port=%s points=%s", address, local_ip, local_port, len(request.points), ) with warnings.catch_warnings(): warnings.filterwarnings("ignore", message=BAC0_REGISTERED_WARNING) with go_style_bacnet_binding(local_port): async with BAC0.start(ip=local_ip, port=local_port, ping=False) as bacnet: for point in request.points: points.append( { "object_type": point.object_type, "object_id": point.object_id, "present_value": await read_present_value(bacnet, address, point), } ) return points async def search_points_async(request: BacnetPointSearchRequest) -> list[dict[str, Any]]: address = bacnet_address(request.ip, request.port) local_ip = local_bacnet_ip(request.ip, request.port) local_port = get_free_udp_port() points: list[dict[str, Any]] = [] logger.info( "BACnet search_points start target=%s local_ip=%s local_port=%s device_id=%s", address, local_ip, local_port, request.bacnet_device_id, ) with warnings.catch_warnings(): warnings.filterwarnings("ignore", message=BAC0_REGISTERED_WARNING) with go_style_bacnet_binding(local_port): async with BAC0.start(ip=local_ip, port=local_port, ping=False) as bacnet: try: object_count = await read_property_direct( bacnet, address, "Device", request.bacnet_device_id, "objectList", array_index=0, ) logger.info("BACnet objectList count target=%s count=%s", address, object_count) except Exception as exc: raise BacnetCommunicationError( f"no response from BACnet device {address} while reading " f"Device {request.bacnet_device_id} objectList" ) from exc for index in range(1, int(object_count) + 1): item = await read_property_direct( bacnet, address, "Device", request.bacnet_device_id, "objectList", array_index=index, ) parsed = split_object_identifier(item) if parsed is None: continue object_type, object_id = parsed if object_type not in POINT_OBJECT_TYPES: continue points.append( { "name": await read_property_or_none(bacnet, address, object_type, object_id, "objectName"), "description": await read_property_or_none(bacnet, address, object_type, object_id, "description"), "object_type": object_type, "object_id": object_id, "present_value": await read_property_or_none(bacnet, address, object_type, object_id, "presentValue"), } ) return points def read_points(request: BacnetPointReadRequest) -> dict[str, Any]: try: with _BACNET_LOCK: points = run_bacnet(read_points_async(request)) except BacnetCommunicationError as exc: logger.exception("BACnet read_points communication failed device=%s", device_payload(request)) return response_payload(1, str(exc), {"device": device_payload(request), "points": []}) except Exception as exc: logger.exception("BACnet read_points failed device=%s", device_payload(request)) return response_payload(1, str(exc), {"device": device_payload(request), "points": []}) return response_payload(0, "success", {"device": device_payload(request), "points": points}) def search_points(request: BacnetPointSearchRequest) -> dict[str, Any]: try: with _BACNET_LOCK: points = run_bacnet(search_points_async(request)) except BacnetCommunicationError as exc: logger.exception("BACnet search_points communication failed device=%s", device_payload(request)) return response_payload(1, str(exc), {"device": device_payload(request), "points": []}) except Exception as exc: logger.exception("BACnet search_points failed device=%s", device_payload(request)) return response_payload(1, str(exc), {"device": device_payload(request), "points": []}) return response_payload(0, "success", {"device": device_payload(request), "points": points})