45 lines
1.1 KiB
Python
45 lines
1.1 KiB
Python
import queue
|
|
import time
|
|
import asyncio
|
|
import websockets
|
|
import json
|
|
|
|
|
|
async def message_consumer(websocket, path):
|
|
subscription = wxbot_instance.subscribe()
|
|
while not websocket.closed:
|
|
try:
|
|
message = subscription.queue.get_nowait()
|
|
subscription.lastRead = time.time()
|
|
if message == None:
|
|
await ws.close(code=1000)
|
|
break
|
|
await websocket.send(json.dumps(message))
|
|
except queue.Empty:
|
|
await asyncio.sleep(1)
|
|
continue
|
|
except WebSocketDisconnected:
|
|
wxbot_instance.unsubscribe(subscription)
|
|
subscription = None
|
|
return
|
|
except:
|
|
wxbot_instance.unsubscribe(subscription)
|
|
subscription = None
|
|
await ws.close(code=1001)
|
|
if subscription != None:
|
|
wxbot_instance.unsubscribe()
|
|
|
|
|
|
wxbot_instance = {"current": None}
|
|
|
|
|
|
class MessageTransfer:
|
|
def __init__(self, msg_queue):
|
|
message_queue = msg_queue
|
|
self.msg_queue = msg_queue
|
|
|
|
def consume_message():
|
|
msg = msg_queue.get()
|
|
for session in sessions:
|
|
pass
|