VNC proxy working in the worker!!
Confirmed multi client access too Streamlit button added
This commit is contained in:
@@ -36,7 +36,10 @@ payload = {
|
||||
token_bytes = json.dumps(payload).encode("utf-8")
|
||||
token_b64 = base64.urlsafe_b64encode(token_bytes).decode("utf-8").rstrip("=")
|
||||
|
||||
print(payload)
|
||||
# print(payload)
|
||||
print(token_b64)
|
||||
|
||||
#URL=http://127.0.0.1:8080/vnc.html?autoconnect=true&host=127.0.0.1&port=6002&path=vnc?token={token_b64}
|
||||
|
||||
# Exmaple output
|
||||
#URL=http://127.0.0.1:8080/vnc.html?autoconnect=true&host=127.0.0.1&port=6002&path=vnc?token=eyJ2aXJ0dWFsX21hY2hpbmVfaWQiOiAiNDAxMmE2YTgtMGIyNC00ZDZlLTlkYmItNzY3NzYyZTI4ZTBhIiwgInVzZXJfdG9rZW4iOiAiZXlKaGJHY2lPaUpJVXpJMU5pSXNJblI1Y0NJNklrcFhWQ0o5LmV5SnpkV0lpT2lKMWMyVnlYMmxrWHpFeU16UWlMQ0p1WVcxbElqb2lRV3hwWTJVaUxDSnliMnhsSWpvaVlXUnRhVzRpTENKbGVIQWlPakUzTlRrd016Z3lPVEFzSW1saGRDSTZNVGMwT0RJek9ESTVNSDAuNUdpcTRKa196RDJoRmhCcGhZWlQ5ZF9WMlJWNEFzd3JTTldHX05nQVEzNCJ9
|
||||
@@ -1,3 +1,5 @@
|
||||
#vnc_worker.py
|
||||
#Example code that needs to be integrated into the production code
|
||||
import asyncio
|
||||
import logging
|
||||
import socketio
|
||||
@@ -2,6 +2,37 @@ import json
|
||||
import streamlit as st
|
||||
from streamlit_server.utils.helpers import format_timestamp
|
||||
|
||||
import jwt
|
||||
import base64
|
||||
import json
|
||||
from datetime import datetime, timedelta
|
||||
|
||||
def generate_vnc_url(virtual_machine_id: str, secret_key: str) -> str:
|
||||
"""
|
||||
Generates a signed and base64-encoded VNC token URL for a given VM ID.
|
||||
"""
|
||||
user_payload = {
|
||||
"sub": "user_id_1234",
|
||||
"name": "Console User",
|
||||
"role": "admin",
|
||||
"exp": datetime.utcnow() + timedelta(hours=3000),
|
||||
"iat": datetime.utcnow()
|
||||
}
|
||||
|
||||
user_token = jwt.encode(user_payload, secret_key, algorithm="HS256")
|
||||
if isinstance(user_token, bytes):
|
||||
user_token = user_token.decode("utf-8")
|
||||
|
||||
outer_payload = {
|
||||
"virtual_machine_id": virtual_machine_id,
|
||||
"user_token": user_token
|
||||
}
|
||||
|
||||
token_bytes = json.dumps(outer_payload).encode("utf-8")
|
||||
token_b64 = base64.urlsafe_b64encode(token_bytes).decode("utf-8").rstrip("=")
|
||||
|
||||
return f"http://127.0.0.1:8080/vnc.html?autoconnect=true&host=127.0.0.1&port=6002&path=vnc?token={token_b64}"
|
||||
|
||||
|
||||
def render_list():
|
||||
"""
|
||||
@@ -212,6 +243,13 @@ def render_list():
|
||||
st.session_state[confirm_key] = False
|
||||
|
||||
with st.popover("Actions", use_container_width=True):
|
||||
if VirtualMachine['status'] == "running":
|
||||
vnc_url = generate_vnc_url(VirtualMachine['id'], "your-very-secret-key") # replace with secure source
|
||||
st.page_link(vnc_url,label="VNC")
|
||||
# if st.button("Open Console", key=f"console_{VirtualMachine['id']}", use_container_width=True):
|
||||
# js = f"window.open('{vnc_url}', '_blank')"
|
||||
# st.markdown(f"<script>{js}</script>", unsafe_allow_html=True)
|
||||
|
||||
# View Details Button
|
||||
if st.button("View Details", key=f"popover_view_{VirtualMachine['id']}", use_container_width=True):
|
||||
st.session_state.selected_resource = VirtualMachine
|
||||
|
||||
+4
-2
@@ -17,9 +17,9 @@ logger = logging.getLogger("middleware")
|
||||
|
||||
# Configure socketio and engineio loggers to use WARNING level to reduce verbosity
|
||||
sio_logger = logging.getLogger("socketio.client")
|
||||
sio_logger.setLevel(logging.WARNING)
|
||||
sio_logger.setLevel(logging.INFO)
|
||||
engineio_logger = logging.getLogger("engineio.client")
|
||||
engineio_logger.setLevel(logging.WARNING)
|
||||
engineio_logger.setLevel(logging.INFO)
|
||||
|
||||
# sio = socketio.AsyncServer(async_mode="aiohttp", cors_allowed_origins="*")
|
||||
app = web.Application()
|
||||
@@ -42,6 +42,7 @@ async def disconnect(sid):
|
||||
|
||||
@sio.on("vnc_frame_from_worker")
|
||||
async def vnc_frame_from_worker(payload):
|
||||
logger.debug(f"Got a frame from worker")
|
||||
vnc_request_id = payload.get('vnc_request_id')
|
||||
|
||||
if not vnc_request_id:
|
||||
@@ -140,6 +141,7 @@ async def vnc_handler(request):
|
||||
"vnc_request_id": req_id,
|
||||
"data": msg.data
|
||||
})
|
||||
logger.debug(f"Got a frame from novnc")
|
||||
elif msg.type == web.WSMsgType.ERROR:
|
||||
logger.error(f"WebSocket error: {ws.exception()}")
|
||||
break
|
||||
|
||||
+5
-2
@@ -1114,6 +1114,7 @@ def vnc_frame_from_worker(payload):
|
||||
|
||||
@socketio.on("vnc_frame_from_novnc")
|
||||
def vnc_frame_from_novnc(payload):
|
||||
logger.debug(f"Got a frame from novnc")
|
||||
#Send the VNC frame back to the VNC Proxy that matches the requestID
|
||||
# logger.debug(f"Scored a VNC frame from novnc {data}")
|
||||
emit("vnc_frame_from_novnc",payload,broadcast=True)
|
||||
@@ -1126,7 +1127,7 @@ def start_vnc_stream(data):
|
||||
virtual_machine_id = data.get("virtual_machine_id")
|
||||
user_token = data.get("user_token")
|
||||
vnc_proxy_sid = request.sid # Who sent the request
|
||||
|
||||
logger.debug(f"User token: {user_token}")
|
||||
if not vnc_request_id or not virtual_machine_id or not user_token:
|
||||
logger.error("Missing required data in VNC stream request")
|
||||
return
|
||||
@@ -1161,7 +1162,9 @@ def start_vnc_stream(data):
|
||||
emit("start_vnc_stream_on_worker", {
|
||||
"vnc_request_id": vnc_request_id,
|
||||
"virtual_machine_id": virtual_machine_id,
|
||||
"user_token": user_token
|
||||
"workload_host_id": workload_host_id,
|
||||
"vnc_port": 5900,
|
||||
"vnc_ip_address": "127.0.0.1"
|
||||
}, to=socket_id)
|
||||
|
||||
logger.info(f"Dispatched VNC stream request {vnc_request_id} to worker {workload_host_id}, from proxy {vnc_proxy_sid}")
|
||||
|
||||
@@ -0,0 +1,44 @@
|
||||
import socket
|
||||
|
||||
class VNCSession:
|
||||
def __init__(self, logger, request_id, vnc_host, vnc_port, on_output, cancel_flag):
|
||||
self.logger = logger
|
||||
self.request_id = request_id
|
||||
self.vnc_host = vnc_host
|
||||
self.vnc_port = int(vnc_port)
|
||||
self.on_output = on_output
|
||||
self.cancel_flag = cancel_flag
|
||||
self.sock = None
|
||||
|
||||
def start(self):
|
||||
try:
|
||||
self.logger.info(f"[VNC] Connecting to {self.vnc_host}:{self.vnc_port}")
|
||||
self.sock = socket.create_connection((self.vnc_host, self.vnc_port))
|
||||
self.logger.info(f"[VNC] Connected to {self.vnc_host}:{self.vnc_port}")
|
||||
|
||||
while not self.cancel_flag():
|
||||
data = self.sock.recv(4096)
|
||||
if not data:
|
||||
self.logger.info(f"[VNC] No more data from {self.request_id}")
|
||||
break
|
||||
self.on_output(data)
|
||||
|
||||
except Exception as e:
|
||||
self.logger.warning(f"[VNC] Error in session {self.request_id}: {e}")
|
||||
finally:
|
||||
self.close()
|
||||
|
||||
def write(self, data):
|
||||
try:
|
||||
self.sock.sendall(data)
|
||||
except Exception as e:
|
||||
self.logger.warning(f"[VNC] Write error on session {self.request_id}: {e}")
|
||||
|
||||
def close(self):
|
||||
if self.sock:
|
||||
try:
|
||||
self.sock.shutdown(socket.SHUT_RDWR)
|
||||
self.sock.close()
|
||||
self.logger.info(f"[VNC] Closed socket for {self.request_id}")
|
||||
except Exception:
|
||||
pass
|
||||
+41
-34
@@ -1,5 +1,7 @@
|
||||
#Worker.py
|
||||
from logger import logger
|
||||
|
||||
from worker.vncsession import VNCSession
|
||||
from worker_tasks.container import ContainerTask
|
||||
from worker_tasks.file_presence import FilePresenceTask
|
||||
from worker_tasks.ping import PingTask
|
||||
@@ -47,7 +49,7 @@ class WorkerClient:
|
||||
self.sio.on("join_accept", self.on_join_accept)
|
||||
self.sio.on("join_reject", self.on_join_reject)
|
||||
self.sio.on("message", self.on_message)
|
||||
self.sio.on("*", self.debug_all_events)
|
||||
# self.sio.on("*", self.debug_all_events)
|
||||
|
||||
|
||||
self.sio.on("container-log",self.logs_retrieve)
|
||||
@@ -65,9 +67,9 @@ class WorkerClient:
|
||||
|
||||
|
||||
self.vnc_sessions = {}
|
||||
self.sio.on("start_vnc_stream", self.start_vnc_stream)
|
||||
self.sio.on("stop_vnc_stream", self.stop_vnc_stream)
|
||||
self.sio.on("vnc_input", self.vnc_input)
|
||||
self.sio.on("start_vnc_stream_on_worker", self.start_vnc_stream)
|
||||
self.sio.on("stop_vnc_stream_on_worker", self.stop_vnc_stream)
|
||||
self.sio.on("vnc_frame_from_novnc", self.vnc_frame_from_novnc) # Add this line
|
||||
|
||||
|
||||
async def on_connect(self):
|
||||
@@ -198,9 +200,9 @@ class WorkerClient:
|
||||
await self.sio.emit("docker_event", payload)
|
||||
logger.debug(f"Sent {source} docker event: {event_type}\n{payload}")
|
||||
|
||||
async def debug_all_events(self, event, data):
|
||||
"""Debug all incoming data."""
|
||||
logger.debug(f"Event: {event} | Data: {json.dumps(data, indent=2)}")
|
||||
# async def debug_all_events(self, event, data):
|
||||
# """Debug all incoming data."""
|
||||
# logger.debug(f"Event: {event} | Data: {json.dumps(data, indent=2)}")
|
||||
|
||||
async def process_event_queue(self):
|
||||
"""Process events from the queue and send them to the server."""
|
||||
@@ -365,39 +367,44 @@ class WorkerClient:
|
||||
logger.info(f"Closed terminal session for {request_id}")
|
||||
|
||||
|
||||
|
||||
|
||||
async def start_vnc_stream(self, data):
|
||||
request_id = data["request_id"]
|
||||
port = data.get("vnc_port", 5900)
|
||||
logger.debug(f"start vnc {data}")
|
||||
request_id = data["vnc_request_id"]
|
||||
vnc_port = int(data["vnc_port"])
|
||||
vnc_host = "127.0.0.1"
|
||||
|
||||
logger.info(f"[VNC] Starting session {request_id} on port {vnc_port}")
|
||||
|
||||
if request_id in self.vnc_sessions:
|
||||
logger.warning(f"[{request_id}] VNC stream already running")
|
||||
logger.warning(f"[VNC] Session already running for {request_id}")
|
||||
return
|
||||
|
||||
from worker_tasks.libvirt import LibvirtVirtualMachineTask
|
||||
LibvirtVirtualMachineTask.start_vnc_stream_thread(
|
||||
request_id=request_id,
|
||||
port=port,
|
||||
emit_func=self.sio.emit,
|
||||
session_registry=self.vnc_sessions,
|
||||
logger=logger
|
||||
)
|
||||
logger.info(f"[{request_id}] VNC stream started on port {port}")
|
||||
def cancel_flag():
|
||||
return request_id not in self.vnc_sessions
|
||||
|
||||
def on_output(vnc_data):
|
||||
asyncio.run(self.sio.emit("vnc_frame_from_worker", {
|
||||
"vnc_request_id": request_id,
|
||||
"data": vnc_data
|
||||
}))
|
||||
|
||||
session = VNCSession(logger, request_id, vnc_host, vnc_port, on_output, cancel_flag)
|
||||
self.vnc_sessions[request_id] = session
|
||||
|
||||
thread = threading.Thread(target=session.start, daemon=True)
|
||||
thread.start()
|
||||
logger.info(f"[VNC] Session thread started for {request_id}")
|
||||
|
||||
async def vnc_frame_from_novnc(self, data):
|
||||
request_id = data["vnc_request_id"]
|
||||
session = self.vnc_sessions.get(request_id)
|
||||
if session:
|
||||
session.write(data["data"])
|
||||
|
||||
|
||||
async def stop_vnc_stream(self, data):
|
||||
request_id = data["request_id"]
|
||||
request_id = data["vnc_request_id"]
|
||||
session = self.vnc_sessions.pop(request_id, None)
|
||||
if session:
|
||||
session["stop_event"].set()
|
||||
for task in session["tasks"]:
|
||||
task.cancel()
|
||||
self.logger.info(f"[{request_id}] VNC stream stopped")
|
||||
|
||||
async def vnc_input(self, data):
|
||||
for req_id, sess in self.vnc_sessions.items():
|
||||
try:
|
||||
sess["writer"].write(data)
|
||||
await sess["writer"].drain()
|
||||
except Exception:
|
||||
pass
|
||||
session.close()
|
||||
logger.info(f"[VNC] Stopped session for {request_id}")
|
||||
|
||||
+1
-46
@@ -1,3 +1,4 @@
|
||||
#libvirt.py
|
||||
import re
|
||||
import libvirt
|
||||
import difflib
|
||||
@@ -8,8 +9,6 @@ from xml.dom import minidom
|
||||
|
||||
from worker_tasks.volumes import VolumeProcessor
|
||||
from worker.settings import settings # Import global settings
|
||||
import asyncio
|
||||
import socket
|
||||
|
||||
class LibvirtVirtualMachineTask:
|
||||
def __init__(self, params, logger):
|
||||
@@ -453,47 +452,3 @@ class LibvirtVirtualMachineTask:
|
||||
|
||||
# Return True if there are differences
|
||||
return len(diff) > 0
|
||||
|
||||
|
||||
@staticmethod
|
||||
def start_vnc_stream_thread(request_id, port, emit_func, session_registry, logger):
|
||||
"""Start a VNC stream thread that reads from a local VNC server and emits data via Socket.IO."""
|
||||
stop_event = asyncio.Event()
|
||||
logger.debug("Starting ")
|
||||
async def vnc_to_middleware(reader):
|
||||
try:
|
||||
while not stop_event.is_set():
|
||||
buf = await reader.read(4096)
|
||||
if not buf:
|
||||
break
|
||||
await emit_func("vnc_output", {
|
||||
"request_id": request_id,
|
||||
"data": buf
|
||||
})
|
||||
logger.debug("Sent frame ")
|
||||
|
||||
except Exception as e:
|
||||
logger.warning(f"[{request_id}] Read loop failed: {e}")
|
||||
finally:
|
||||
stop_event.set()
|
||||
|
||||
async def wait_for_stop(writer):
|
||||
await stop_event.wait()
|
||||
writer.close()
|
||||
|
||||
async def start_stream():
|
||||
try:
|
||||
reader, writer = await asyncio.open_connection("127.0.0.1", port)
|
||||
logger.info(f"[{request_id}] Connected to VNC server on port {port}")
|
||||
session_registry[request_id] = {
|
||||
"stop_event": stop_event,
|
||||
"writer": writer,
|
||||
"tasks": [
|
||||
asyncio.create_task(vnc_to_middleware(reader)),
|
||||
asyncio.create_task(wait_for_stop(writer)),
|
||||
]
|
||||
}
|
||||
except Exception as e:
|
||||
logger.error(f"[{request_id}] Failed to connect to VNC port {port}: {e}")
|
||||
|
||||
asyncio.create_task(start_stream())
|
||||
Reference in New Issue
Block a user