re-implimented the send_message capability which is useful for troubleshooting
This commit is contained in:
@@ -86,3 +86,32 @@ def register_socketio_handlers(socketio):
|
||||
except Exception as e:
|
||||
logger.exception(f"Error processing join for {worker_id}: {e}")
|
||||
emit("join_reject", {"message": "Internal error during join"}, to=request.sid)
|
||||
|
||||
|
||||
@socketio.on("send_message")
|
||||
def handle_send_message(data):
|
||||
"""
|
||||
Relay a message from a client to a specific worker by worker_id.
|
||||
"""
|
||||
worker_id = data.get("worker_id")
|
||||
message = data.get("message")
|
||||
|
||||
if not worker_id or not message:
|
||||
logger.warning(f"Invalid send_message request: {data}")
|
||||
emit("error", {"message": "Missing worker_id or message"})
|
||||
return
|
||||
|
||||
with worker_lock:
|
||||
target_sid = connected_workers.get(worker_id)
|
||||
|
||||
if not target_sid:
|
||||
logger.error(f"Worker {worker_id} is not connected")
|
||||
emit("error", {"message": f"Worker {worker_id} is not connected"})
|
||||
return
|
||||
|
||||
try:
|
||||
logger.info(f"Relaying message to Worker {worker_id} (SID={target_sid})")
|
||||
emit("message", {"message": message}, to=target_sid)
|
||||
except Exception as e:
|
||||
logger.exception(f"Failed to send message to {worker_id}: {e}")
|
||||
emit("error", {"message": f"Failed to send message to {worker_id}"})
|
||||
|
||||
@@ -51,4 +51,4 @@ def start_background_threads():
|
||||
if __name__ == "__main__":
|
||||
start_background_threads()
|
||||
# TODO - Remove allow_unsafe_werkzeug
|
||||
base.socketio.run(app, host="0.0.0.0", port=6001,allow_unsafe_werkzeug=True )
|
||||
base.socketio.run(app, host="0.0.0.0", port=6001,allow_unsafe_werkzeug=True,debug=True )
|
||||
|
||||
Reference in New Issue
Block a user