minhvtt commited on
Commit
ee8efa3
·
verified ·
1 Parent(s): 564eafe

Update app/services/alert_hub.py

Browse files
Files changed (1) hide show
  1. app/services/alert_hub.py +51 -48
app/services/alert_hub.py CHANGED
@@ -1,48 +1,51 @@
1
- from collections import defaultdict
2
- from typing import Any
3
-
4
- from fastapi import WebSocket
5
-
6
-
7
- class AlertHub:
8
- def __init__(self) -> None:
9
- self.admin_clients: set[WebSocket] = set()
10
- self.device_clients: dict[str, set[WebSocket]] = defaultdict(set)
11
-
12
- async def connect_admin(self, websocket: WebSocket) -> None:
13
- await websocket.accept()
14
- self.admin_clients.add(websocket)
15
-
16
- async def connect_device(self, device_id: str, websocket: WebSocket) -> None:
17
- await websocket.accept()
18
- self.device_clients[device_id].add(websocket)
19
-
20
- def disconnect_admin(self, websocket: WebSocket) -> None:
21
- self.admin_clients.discard(websocket)
22
-
23
- def disconnect_device(self, device_id: str, websocket: WebSocket) -> None:
24
- if device_id in self.device_clients:
25
- self.device_clients[device_id].discard(websocket)
26
-
27
- async def broadcast_admin(self, payload: dict[str, Any]) -> None:
28
- stale: list[WebSocket] = []
29
- for client in self.admin_clients:
30
- try:
31
- await client.send_json(payload)
32
- except Exception:
33
- stale.append(client)
34
- for client in stale:
35
- self.disconnect_admin(client)
36
-
37
- async def send_device(self, device_id: str, payload: dict[str, Any]) -> None:
38
- stale: list[WebSocket] = []
39
- for client in self.device_clients.get(device_id, set()):
40
- try:
41
- await client.send_json(payload)
42
- except Exception:
43
- stale.append(client)
44
- for client in stale:
45
- self.disconnect_device(device_id, client)
46
-
47
-
48
- alert_hub = AlertHub()
 
 
 
 
1
+ from collections import defaultdict
2
+ from typing import Any
3
+
4
+ from fastapi import WebSocket
5
+
6
+
7
+ class AlertHub:
8
+ def __init__(self) -> None:
9
+ self.admin_clients: set[WebSocket] = set()
10
+ self.device_clients: dict[str, set[WebSocket]] = defaultdict(set)
11
+
12
+ async def connect_admin(self, websocket: WebSocket) -> None:
13
+ await websocket.accept()
14
+ self.admin_clients.add(websocket)
15
+
16
+ async def connect_device(self, device_id: str, websocket: WebSocket) -> None:
17
+ await websocket.accept()
18
+ self.device_clients[device_id].add(websocket)
19
+
20
+ def disconnect_admin(self, websocket: WebSocket) -> None:
21
+ self.admin_clients.discard(websocket)
22
+
23
+ def disconnect_device(self, device_id: str, websocket: WebSocket) -> None:
24
+ if device_id in self.device_clients:
25
+ self.device_clients[device_id].discard(websocket)
26
+
27
+ async def broadcast_admin(self, payload: dict[str, Any]) -> None:
28
+ stale: list[WebSocket] = []
29
+ for client in self.admin_clients:
30
+ try:
31
+ await client.send_json(payload)
32
+ except Exception:
33
+ stale.append(client)
34
+ for client in stale:
35
+ self.disconnect_admin(client)
36
+
37
+ async def send_device(self, device_id: str, payload: dict[str, Any]) -> None:
38
+ stale: list[WebSocket] = []
39
+ for client in self.device_clients.get(device_id, set()):
40
+ try:
41
+ await client.send_json(payload)
42
+ except Exception:
43
+ stale.append(client)
44
+ for client in stale:
45
+ self.disconnect_device(device_id, client)
46
+
47
+ def device_connection_count(self, device_id: str) -> int:
48
+ return len(self.device_clients.get(device_id, set()))
49
+
50
+
51
+ alert_hub = AlertHub()