diff --git a/streamlit_server/views/containers.py b/streamlit_server/views/containers.py index a9c45fb..4a3e13d 100644 --- a/streamlit_server/views/containers.py +++ b/streamlit_server/views/containers.py @@ -10,7 +10,7 @@ def render_list(): # Initialize session state flags for pod deletion for flag in ['pod_delete_all_confirmation_shown', 'pod_delete_all_confirmed', - 'pod_delete_confirmation_shown', 'pod_delete_confirmed']: + 'pod_delete_confirmation_shown', 'pod_delete_confirmed','container_delete_confirmation_shown','container_delete_confirmed']: if flag not in st.session_state: st.session_state[flag] = False @@ -192,7 +192,7 @@ def render_detail(): Renders the details of a specific container workload. """ - container = st.session_state.client.get_container(st.session_state.selected_resource['id']) + container = st.session_state.client.get_container(st.session_state.selected_resource['container_id']) if not container: st.error("Container workload not found.") st.session_state.selected_resource = None diff --git a/streamlit_server/views/websocket.py b/streamlit_server/views/websocket.py index 3d582a2..414d94a 100644 --- a/streamlit_server/views/websocket.py +++ b/streamlit_server/views/websocket.py @@ -4,7 +4,7 @@ import socketio import os # Base URL pulled from environment variable -BASE_URL = os.getenv("SOCKET_SERVER_URL", "http://172.17.0.1:6101") +BASE_URL = os.getenv("SOCKET_SERVER_URL", "http://172.17.0.1:6001") API_URL = f"{BASE_URL}/api/connected_clients" def load_clients(): @@ -21,7 +21,7 @@ def render(): st.title("Connected Clients") clients = load_clients() - + if clients: selected_client = st.selectbox("Select a Worker ID", [client["worker_id"] for client in clients]) st.write("## Connected Clients") @@ -38,8 +38,8 @@ def render(): message = st.text_area("Enter your message:") if st.button("Send Message"): if message.strip(): - sio.emit("send_message", {"worker_id": selected_client, "message": message}) - st.success(f"Message sent to {selected_client}") + message_result=sio.emit("send_message", {"worker_id": selected_client, "message": message}) + st.success(f"Message sent to {selected_client} {message_result}") else: st.error("Message cannot be empty!") else: diff --git a/websocket_server/events/messages.py b/websocket_server/events/messages.py index c5de556..5d0f64b 100644 --- a/websocket_server/events/messages.py +++ b/websocket_server/events/messages.py @@ -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}"}) diff --git a/websocket_server/main.py b/websocket_server/main.py index 5807504..3254c7a 100644 --- a/websocket_server/main.py +++ b/websocket_server/main.py @@ -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 )