from __future__ import annotations import asyncio import logging import os import socket import struct import threading import time 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 ( BacnetBBMDWhoIsRequest, BacnetPointReadRequest, BacnetPointSearchRequest, BacnetPointSpec, bac0_object_type, normalize_object_type, ) POINT_OBJECT_TYPES = { "AnalogInput", "AnalogOutput", "AnalogValue", "BinaryInput", "BinaryOutput", "BinaryValue", "MultiStateInput", "MultiStateOutput", "MultiStateValue", } BVLC_TYPE_BACNET_IP = 0x81 BVLC_RESULT = 0x00 BVLC_FORWARDED_NPDU = 0x04 BVLC_REGISTER_FOREIGN_DEVICE = 0x05 BVLC_DISTRIBUTE_BROADCAST_TO_NETWORK = 0x09 BVLC_ORIGINAL_UNICAST_NPDU = 0x0A BVLC_ORIGINAL_BROADCAST_NPDU = 0x0B BVLC_RESULT_SUCCESS = 0x0000 APDU_UNCONFIRMED_REQUEST = 0x10 SERVICE_UNCONFIRMED_I_AM = 0x00 SERVICE_UNCONFIRMED_WHO_IS = 0x08 OBJECT_TYPE_DEVICE = 8 OBJECT_ID_INSTANCE_MASK = 0x3FFFFF SEGMENTATION_SUPPORT = { 0: "segmentedBoth", 1: "segmentedTransmit", 2: "segmentedReceive", 3: "noSegmentation", } _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 bbmd_payload(request: BacnetBBMDWhoIsRequest) -> dict[str, Any]: return { "bbmd_ip": request.bbmd_ip, "bbmd_port": request.bbmd_port, "ttl": request.ttl, "timeout": request.timeout, "low_limit": request.low_limit, "high_limit": request.high_limit, } def bbmd_error_payload(request: BacnetBBMDWhoIsRequest | None) -> dict[str, Any]: if request is not None: return bbmd_payload(request) return { "bbmd_ip": os.getenv("BACNET_BBMD_IP"), "bbmd_port": os.getenv("BACNET_BBMD_PORT", "47808"), } def env_int(name: str, default: int | None = None) -> int | None: value = os.getenv(name) if value is None or value == "": return default try: return int(value) except ValueError as exc: raise BacnetCommunicationError(f"{name} must be an integer") from exc def env_float(name: str, default: float) -> float: value = os.getenv(name) if value is None or value == "": return default try: return float(value) except ValueError as exc: raise BacnetCommunicationError(f"{name} must be a number") from exc def bbmd_whois_request_from_env() -> BacnetBBMDWhoIsRequest: bbmd_ip = os.getenv("BACNET_BBMD_IP") if not bbmd_ip: raise BacnetCommunicationError("BACNET_BBMD_IP environment variable is required") return BacnetBBMDWhoIsRequest( bbmd_ip=bbmd_ip, bbmd_port=env_int("BACNET_BBMD_PORT", 47808), ttl=env_int("BACNET_BBMD_TTL", 60), timeout=env_float("BACNET_BBMD_WHOIS_TIMEOUT", 5), low_limit=env_int("BACNET_BBMD_LOW_LIMIT", 0), high_limit=env_int("BACNET_BBMD_HIGH_LIMIT", OBJECT_ID_INSTANCE_MASK), local_device_id=env_int("BACNET_BBMD_LOCAL_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_bbmd_socket(local_ip: str, local_port: int) -> socket.socket: bind_ip = os.getenv("BACNET_BIND_IP", local_ip) 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, local_port)) return sock def bvlc_packet(function: int, payload: bytes) -> bytes: return struct.pack(">BBH", BVLC_TYPE_BACNET_IP, function, 4 + len(payload)) + payload def register_foreign_device_packet(ttl: int) -> bytes: return bvlc_packet(BVLC_REGISTER_FOREIGN_DEVICE, struct.pack(">H", ttl)) def context_unsigned(tag: int, value: int) -> bytes: if not 0 <= tag <= 14: raise ValueError("BACnet context tag must be between 0 and 14") if not 0 <= value <= 0xFFFFFFFF: raise ValueError("BACnet unsigned context value out of range") if value <= 0xFF: data = value.to_bytes(1, "big") elif value <= 0xFFFF: data = value.to_bytes(2, "big") else: data = value.to_bytes(4, "big") return bytes([(tag << 4) | 0x08 | len(data)]) + data def whois_apdu(low_limit: int | None, high_limit: int | None) -> bytes: payload = bytearray([APDU_UNCONFIRMED_REQUEST, SERVICE_UNCONFIRMED_WHO_IS]) if low_limit is not None and high_limit is not None: payload.extend(context_unsigned(0, low_limit)) payload.extend(context_unsigned(1, high_limit)) return bytes(payload) def bbmd_whois_packet(low_limit: int | None, high_limit: int | None) -> bytes: # Match bacpypes LocalBroadcast(): BIPForeign wraps this NPDU in Distribute-Broadcast-To-Network. npdu = b"\x01\x00" + whois_apdu(low_limit, high_limit) return bvlc_packet(BVLC_DISTRIBUTE_BROADCAST_TO_NETWORK, npdu) def parse_bvlc_result(packet: bytes) -> int | None: if len(packet) < 6 or packet[0] != BVLC_TYPE_BACNET_IP or packet[1] != BVLC_RESULT: return None return int.from_bytes(packet[4:6], "big") def udp_endpoint_from_bytes(data: bytes) -> tuple[str, int]: return socket.inet_ntoa(data[:4]), int.from_bytes(data[4:6], "big") def npdu_apdu(packet: bytes) -> bytes | None: if len(packet) < 2 or packet[0] != 0x01: return None control = packet[1] offset = 2 has_destination = bool(control & 0x20) if has_destination: if len(packet) < offset + 3: return None dlen = packet[offset + 2] offset += 3 + dlen if control & 0x08: if len(packet) < offset + 3: return None slen = packet[offset + 2] offset += 3 + slen if has_destination: offset += 1 if control & 0x80: return None if offset >= len(packet): return None return packet[offset:] def read_application_tag(packet: bytes, offset: int) -> tuple[int, bool, bytes, int] | None: if offset >= len(packet): return None first = packet[offset] offset += 1 tag = first >> 4 class_tag = bool(first & 0x08) length = first & 0x07 if tag == 0x0F: if offset >= len(packet): return None tag = packet[offset] offset += 1 if length == 5: if offset >= len(packet): return None length = packet[offset] offset += 1 elif length in (6, 7): return None if len(packet) < offset + length: return None value = packet[offset : offset + length] return tag, class_tag, value, offset + length def parse_i_am_apdu(apdu: bytes) -> dict[str, Any] | None: if len(apdu) < 2 or apdu[0] != APDU_UNCONFIRMED_REQUEST or apdu[1] != SERVICE_UNCONFIRMED_I_AM: return None offset = 2 object_id: int | None = None max_apdu: int | None = None segmentation: int | None = None vendor_id: int | None = None for index in range(4): tag = read_application_tag(apdu, offset) if tag is None: return None tag_number, class_tag, value, offset = tag if class_tag: return None if index == 0 and tag_number == 12 and len(value) == 4: raw_object_id = int.from_bytes(value, "big") if raw_object_id >> 22 != OBJECT_TYPE_DEVICE: return None object_id = raw_object_id & OBJECT_ID_INSTANCE_MASK elif index == 1 and tag_number == 2: max_apdu = int.from_bytes(value, "big") elif index == 2 and tag_number == 9: segmentation = int.from_bytes(value, "big") elif index == 3 and tag_number == 2: vendor_id = int.from_bytes(value, "big") else: return None if object_id is None or max_apdu is None or segmentation is None or vendor_id is None: return None return { "bacnet_device_id": object_id, "max_apdu": max_apdu, "segmentation": SEGMENTATION_SUPPORT.get(segmentation, str(segmentation)), "vendor_id": vendor_id, } def parse_i_am_packet(packet: bytes, source: tuple[str, int]) -> dict[str, Any] | None: if len(packet) < 4 or packet[0] != BVLC_TYPE_BACNET_IP: return None length = int.from_bytes(packet[2:4], "big") if length > len(packet): return None function = packet[1] payload = packet[4:length] device_ip, device_port = source if function == BVLC_FORWARDED_NPDU: if len(payload) < 6: return None device_ip, device_port = udp_endpoint_from_bytes(payload[:6]) payload = payload[6:] elif function not in (BVLC_ORIGINAL_UNICAST_NPDU, BVLC_ORIGINAL_BROADCAST_NPDU): return None apdu = npdu_apdu(payload) if apdu is None: return None device = parse_i_am_apdu(apdu) if device is None: return None device["ip"] = device_ip device["port"] = device_port return device def wait_for_bbmd_registration(sock: socket.socket, timeout: float) -> None: deadline = time.monotonic() + timeout while True: remaining = deadline - time.monotonic() if remaining <= 0: raise BacnetCommunicationError("BBMD foreign device registration timeout") sock.settimeout(remaining) try: packet, _source = sock.recvfrom(2048) except socket.timeout as exc: raise BacnetCommunicationError("BBMD foreign device registration timeout") from exc result = parse_bvlc_result(packet) if result is None: continue if result != BVLC_RESULT_SUCCESS: raise BacnetCommunicationError(f"BBMD foreign device registration failed with code {result}") return def collect_bbmd_i_ams( sock: socket.socket, timeout: float, low_limit: int | None, high_limit: int | None, local_device_id: int | None, ) -> list[dict[str, Any]]: devices: dict[tuple[int, str, int], dict[str, Any]] = {} deadline = time.monotonic() + timeout while True: remaining = deadline - time.monotonic() if remaining <= 0: return list(devices.values()) sock.settimeout(remaining) try: packet, source = sock.recvfrom(2048) except socket.timeout: return list(devices.values()) device = parse_i_am_packet(packet, (source[0], source[1])) if device is None: continue device_id = device["bacnet_device_id"] if device_id == 0: continue if local_device_id is not None and device_id == local_device_id: continue if low_limit is not None and high_limit is not None and not low_limit <= device_id <= high_limit: continue devices[(device_id, device["ip"], device["port"])] = device 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 bbmd_whois_devices(request: BacnetBBMDWhoIsRequest) -> list[dict[str, Any]]: local_ip = request.local_ip or get_local_ip(request.bbmd_ip, request.bbmd_port) local_port = request.local_port or get_free_udp_port() bbmd_address = (request.bbmd_ip, request.bbmd_port) logger.info( "BACnet BBMD whois start bbmd=%s:%s local_ip=%s local_port=%s ttl=%s timeout=%s", request.bbmd_ip, request.bbmd_port, local_ip, local_port, request.ttl, request.timeout, ) with closing(create_bbmd_socket(local_ip, local_port)) as sock: sock.sendto(register_foreign_device_packet(request.ttl), bbmd_address) wait_for_bbmd_registration(sock, min(request.timeout, 5.0)) try: sock.sendto(bbmd_whois_packet(request.low_limit, request.high_limit), bbmd_address) return collect_bbmd_i_ams( sock, request.timeout, request.low_limit, request.high_limit, request.local_device_id, ) finally: try: sock.sendto(register_foreign_device_packet(0), bbmd_address) except OSError: logger.debug("BACnet BBMD unregister failed", exc_info=True) 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}) def bbmd_whois(request: BacnetBBMDWhoIsRequest | None = None) -> dict[str, Any]: try: request = request or bbmd_whois_request_from_env() with _BACNET_LOCK: devices = bbmd_whois_devices(request) except BacnetCommunicationError as exc: logger.exception("BACnet BBMD whois communication failed bbmd=%s", bbmd_error_payload(request)) return response_payload(1, str(exc), {"bbmd": bbmd_error_payload(request), "devices": []}) except Exception as exc: logger.exception("BACnet BBMD whois failed bbmd=%s", bbmd_error_payload(request)) return response_payload(1, str(exc), {"bbmd": bbmd_error_payload(request), "devices": []}) return response_payload(0, "success", {"bbmd": bbmd_payload(request), "devices": devices})