Files
brewpi/ws/server/ws_server_multi_user.py
T
2020-11-24 21:11:48 +01:00

69 lines
1.7 KiB
Python

import asyncio
from ws.server.ws_server import WsServer
from ws.user import User, UserSet, update
import abc
from ws.connection import IConnection
class IWsServer:
@abc.abstractmethod
def on_recv(self, data, loop):
pass
@abc.abstractmethod
def on_send(self, loop):
pass
class WsServerMultiUser(WsServer):
def __init__(self, loop=None, listener: IConnection = None):
WsServer.__init__(self, loop)
self.listener = listener
self.USERS = UserSet()
self.global_state = {}
self.listener.on_connect(self.loop)
async def notify_state(self, data):
if self.USERS: # asyncio.wait doesn't accept an empty list
await asyncio.wait([user.send(data) for user in self.USERS])
async def notify_users(self):
if self.USERS: # asyncio.wait doesn't accept an empty list
message = {"info": {"type": "users", "count": len(self.USERS)}}
await asyncio.wait([user.send(message) for user in self.USERS])
async def register(self, websocket):
usr = User(websocket)
usr.path_add("info")
self.USERS.add(websocket, usr)
await self.notify_users()
async def unregister(self, websocket):
self.USERS.remove(websocket)
await self.notify_users()
async def handler_recv(self, websocket, path):
while True:
try:
data = await websocket.recv()
usr = self.USERS.get(websocket)
processed = await usr.process(data)
if processed:
await usr.send(self.global_state)
else:
await self.listener.on_recv(data)
except Exception as e:
print(e)
break
async def handler_send(self, websocket, path):
while True:
data = await self.listener.on_send()
try:
await self.notify_state(data)
# Update global state
self.global_state = update(self.global_state, data)
except Exception as e:
print(e)
break