- ws/message.py, ws/server/ws_server_multi_user.py: drop the loop= kwarg and replace asyncio.wait(coroutines) with asyncio.gather(), both of which were removed/forbidden in Python 3.10+. - brewpi/requirements.txt: pin websockets==10.4, the latest version that still supports this code's serve()/handler API. - config.json.sim: move Model to the top level and add the missing gain field so it matches the schema brewpi.py and config.json.templ expect. - brewpi/brewpi.py: add -c/-m/-s/-d CLI args for selecting config/model files, simulation mode, and debug. - .gitignore: exclude __pycache__, logs, the local config.json, and editor/venv directories. Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
146 lines
3.3 KiB
Python
146 lines
3.3 KiB
Python
import asyncio
|
|
from ws.connection import IConnection
|
|
import json
|
|
import queue
|
|
|
|
class MsgIo:
|
|
def __init__(self, key, send, send_from_thread):
|
|
self.key = key
|
|
self.sender = send
|
|
self.send_from_thread_cb = send_from_thread
|
|
self.receiver = None
|
|
|
|
def set_recv_handler(self, handler):
|
|
self.receiver = handler
|
|
|
|
async def on_recv(self, data):
|
|
print("MessageHandler {}".format(data))
|
|
if self.receiver is not None:
|
|
await self.receiver(data[self.key])
|
|
|
|
async def send(self, data):
|
|
return await self.sender({self.key: data})
|
|
|
|
def send_from_thread(self, data):
|
|
return self.send_from_thread_cb({self.key: data})
|
|
|
|
def can_key(self, key):
|
|
return self.key == key
|
|
|
|
def get_key(self):
|
|
return self.key
|
|
|
|
|
|
class MessageDispatcher(IConnection):
|
|
def __init__(self, auto_subscribe=False, loop=None):
|
|
self.msg_handlers: MsgIo = []
|
|
self.state = None
|
|
self.state = asyncio.Queue()
|
|
self.auto_subscribe = auto_subscribe
|
|
self.loop = loop
|
|
|
|
def msgio_get(self, key):
|
|
obj = MsgIo(key, self.send, self.send_from_thread)
|
|
if key not in self.msg_handlers:
|
|
self.msg_handlers.append(obj)
|
|
return obj
|
|
return None
|
|
|
|
async def on_connect(self):
|
|
if self.auto_subscribe:
|
|
for handler in self.msg_handlers:
|
|
key = handler.get_key()
|
|
print("Would subscribe {}".format(key))
|
|
await self.send({'+': key})
|
|
|
|
async def on_disconnect(self):
|
|
pass
|
|
|
|
async def on_recv(self, data):
|
|
d = json.loads(data)
|
|
for key_req in d.keys():
|
|
for h in self.msg_handlers:
|
|
if h.can_key(key_req):
|
|
await h.on_recv(d)
|
|
|
|
async def on_send(self):
|
|
return await self.state.get()
|
|
|
|
async def send(self, data):
|
|
await self.state.put(data)
|
|
|
|
def send_from_thread(self, data):
|
|
self.loop.create_task(self.send(data))
|
|
|
|
|
|
class MsgIoSync:
|
|
def __init__(self, key, send):
|
|
self.key = key
|
|
self.sender = send
|
|
self.receiver = None
|
|
|
|
def set_recv_handler(self, handler):
|
|
self.receiver = handler
|
|
|
|
def on_recv(self, data):
|
|
print("MessageHandler {}".format(data))
|
|
if self.receiver is not None:
|
|
self.receiver(data[self.key])
|
|
|
|
def send(self, data):
|
|
return self.sender({self.key: data})
|
|
|
|
def can_key(self, key):
|
|
return self.key == key
|
|
|
|
def get_key(self):
|
|
return self.key
|
|
|
|
|
|
class MessageDispatcherSync(IConnection):
|
|
def __init__(self, auto_subscribe=False):
|
|
self.msg_handlers: MsgIoSync = []
|
|
self.state = None
|
|
self.state = queue.Queue()
|
|
self.auto_subscribe = auto_subscribe
|
|
|
|
def msgio_get(self, key):
|
|
obj = MsgIoSync(key, self.send)
|
|
if key not in self.msg_handlers:
|
|
self.msg_handlers.append(obj)
|
|
return obj
|
|
return None
|
|
|
|
async def on_connect(self):
|
|
if self.auto_subscribe:
|
|
asyncio.get_event_loop().run_in_executor(None, self.subscribe)
|
|
|
|
def subscribe(self):
|
|
for handler in self.msg_handlers:
|
|
key = handler.get_key()
|
|
print("Would subscribe {}".format(key))
|
|
self.send({'+': key})
|
|
|
|
async def on_disconnect(self):
|
|
pass
|
|
|
|
async def on_recv(self, data):
|
|
d = json.loads(data)
|
|
for key_req in d.keys():
|
|
for h in self.msg_handlers:
|
|
if h.can_key(key_req):
|
|
await asyncio.get_event_loop().run_in_executor(None, h.on_recv, d)
|
|
|
|
async def on_send(self):
|
|
result = None
|
|
try:
|
|
result = await asyncio.get_event_loop().run_in_executor(None, self.state.get, True, 1.0)
|
|
print("Send subscribe {}".format(result))
|
|
except Exception as e:
|
|
pass
|
|
return result
|
|
|
|
def send(self, data):
|
|
self.state.put(data)
|
|
|