Lu Xianghui 1 месяц назад
Родитель
Сommit
33629d7a56
1 измененных файлов с 103 добавлено и 50 удалено
  1. 103 50
      app/services/bacnet_service.py

+ 103 - 50
app/services/bacnet_service.py

@@ -7,6 +7,7 @@ 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.*"
@@ -14,8 +15,13 @@ BAC0_REGISTERED_WARNING = r"object type .* for vendor identifier 842 already reg
 warnings.filterwarnings("ignore", message=BAC0_REGISTERED_WARNING)
 
 import BAC0
+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
 
@@ -123,6 +129,49 @@ def get_free_udp_port() -> int:
         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
+    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
+        yield
+    finally:
+        bacpypes_app.NormalLinkLayer_ipv4 = original_normal_link_layer
+        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"))
@@ -165,6 +214,8 @@ def ignore_cancelled_bacnet_broadcast_endpoint(loop: asyncio.AbstractEventLoop,
     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)
 
 
@@ -223,15 +274,16 @@ async def read_points_async(request: BacnetPointReadRequest) -> list[dict[str, A
 
     with warnings.catch_warnings():
         warnings.filterwarnings("ignore", message=BAC0_REGISTERED_WARNING)
-        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),
-                    }
-                )
+        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
 
 
@@ -250,47 +302,48 @@ async def search_points_async(request: BacnetPointSearchRequest) -> list[dict[st
 
     with warnings.catch_warnings():
         warnings.filterwarnings("ignore", message=BAC0_REGISTERED_WARNING)
-        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"),
-                    }
-                )
+        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