Fix: remove hardcoded vnc vm port & remove consts

This commit is contained in:
2026-05-25 19:07:25 +05:45
parent 0e59fee995
commit b7dbcaf315
10 changed files with 87 additions and 46 deletions
+7
View File
@@ -45,6 +45,8 @@ VNC_BASE_URL=https://vnc-console.xcloudify.tech/vnc.html
VNC_PROXY_HOST=vnc-proxy.xcloudify.tech
VNC_PROXY_PORT=443
VNC_PROXY_WS_PATH=/vnc
# Ports used by the local VNC proxy process (vnc_proxy.py)
VNC_PROXY_HTTP_PORT=6002
VNC_PORT_LOCAL=5900
@@ -58,5 +60,10 @@ CHECK_PERMISSION_URL=http://172.17.0.1:5000/api/check_permission
JWT_SECRET_KEY=change-me-to-a-strong-random-secret
# Base domain for user-facing resources (certificates, etc.)
BASE_DOMAIN=xcloudify.tech
# Tunnel domain used for container/workload hostnames via Cloudflare tunnels
TUNNEL_DOMAIN=hawkvelt.tech
# Default OVS bridge name used when creating networks (optional)
XCF_DEFAULT_OVS_BRIDGE=br-xcloudify
+5 -1
View File
@@ -20,7 +20,7 @@ from logger import logger
from app.utils.standard_responses import api_response
from flask_cors import CORS
from runtime_urls import APP_ENV, CLOUDFLARE_ACCOUNT_ID, CLOUDFLARE_API_TOKEN, CLOUDFLARE_ZONE_ID, DATABASE_URL, JWT_SECRET_KEY, PING_HEARTBEAT_TIMEOUT_SECONDS, REDIS_BROKER_URL, REDIS_RESULT_BACKEND_URL, REDIS_URL, SCHEDULER_MAX_ALLOCATION_ATTEMPTS, SCHEDULER_RETRY_BACKOFF_FACTOR, SCHEDULER_RETRY_BASE_DELAY_SECONDS, SCHEDULER_RETRY_JITTER_SECONDS, SCHEDULER_RETRY_MAX_DELAY_SECONDS, VNC_BASE_URL, VNC_PROXY_HOST, VNC_PROXY_PORT, VNC_SECRET_KEY, WEBSOCKET_SERVER_URL
from runtime_urls import API_BASE_URL, APP_ENV, BASE_DOMAIN, CLOUDFLARE_ACCOUNT_ID, CLOUDFLARE_API_TOKEN, CLOUDFLARE_ZONE_ID, DATABASE_URL, JWT_SECRET_KEY, PING_HEARTBEAT_TIMEOUT_SECONDS, REDIS_BROKER_URL, REDIS_RESULT_BACKEND_URL, REDIS_URL, SCHEDULER_MAX_ALLOCATION_ATTEMPTS, SCHEDULER_RETRY_BACKOFF_FACTOR, SCHEDULER_RETRY_BASE_DELAY_SECONDS, SCHEDULER_RETRY_JITTER_SECONDS, SCHEDULER_RETRY_MAX_DELAY_SECONDS, TUNNEL_DOMAIN, VNC_BASE_URL, VNC_PROXY_HOST, VNC_PROXY_PORT, VNC_PROXY_WS_PATH, VNC_SECRET_KEY, WEBSOCKET_SERVER_URL
# --------------------------------------------------------------------------- #
# 1. Flask application & database #
@@ -102,9 +102,13 @@ app.config["CLOUDFLARE_ZONE_ID"] = CLOUDFLARE_ZONE_ID
app.config["PING_HEARTBEAT_TIMEOUT_SECONDS"] = PING_HEARTBEAT_TIMEOUT_SECONDS # How long before a Websocket PING\PONG is classed as a failure and triggers a worker offline event
app.config["REDIS_URL"] = REDIS_URL # Redis Database 2 for operational tasks like websocket server and PING logging
app.config["VNC_SECRET_KEY"] = VNC_SECRET_KEY
app.config["API_BASE_URL"] = API_BASE_URL
app.config["BASE_DOMAIN"] = BASE_DOMAIN
app.config["TUNNEL_DOMAIN"] = TUNNEL_DOMAIN
app.config["VNC_BASE_URL"] = VNC_BASE_URL
app.config["VNC_PROXY_HOST"] = VNC_PROXY_HOST
app.config["VNC_PROXY_PORT"] = VNC_PROXY_PORT
app.config["VNC_PROXY_WS_PATH"] = VNC_PROXY_WS_PATH
# JWT secret used for decoding Authorization bearer tokens for audit attribution
app.config["JWT_SECRET_KEY"] = JWT_SECRET_KEY
+4 -5
View File
@@ -4,7 +4,7 @@ Certificate API Routes
This module provides API routes for managing Certificate Authorities and Certificates.
"""
from flask import request
from flask import request, current_app
from app import db, logger
from app.models.certificate_models import CertificateAuthority, Certificate, CertificateRevocationList
from app.models.models import Project, AuditEntry
@@ -88,11 +88,10 @@ def create_ca_for_project(project_id):
)
# Create a new CA
domain_name = f"{project_id}.xcloudify.tech"
domain_name = f"{project_id}.{current_app.config['BASE_DOMAIN']}"
common_name = f"rootca.{domain_name}"
new_ca_id=uuid.uuid4()
# TODO - Change the URL to be a BASE_URL Param sourced from ENV
crl_url = f"https://api.xcloudify.tech/api/certificates/crl/{new_ca_id}"
new_ca_id = uuid.uuid4()
crl_url = f"{current_app.config['API_BASE_URL']}/certificates/crl/{new_ca_id}"
ca_data = generate_self_signed_ca(
common_name=common_name,
+3 -5
View File
@@ -1,4 +1,4 @@
from flask import request
from flask import request, current_app
from werkzeug.exceptions import BadRequest
from datetime import datetime
from enum import Enum
@@ -9,9 +9,6 @@ from app.controller import api_bp
from app.utils.standard_responses import api_response
# Constants
SECRET_KEY = "your-very-secret-key"
# Enum for supported resource types
class ResourceType(Enum):
PROJECT = "project"
@@ -22,7 +19,8 @@ class ResourceType(Enum):
# Token decoding utility
def decode_token(token):
try:
payload = jwt.decode(token, SECRET_KEY, algorithms=["HS256"])
secret_key = current_app.config["VNC_SECRET_KEY"]
payload = jwt.decode(token, secret_key, algorithms=["HS256"])
return {'user_id': payload['sub']}
except jwt.ExpiredSignatureError:
raise BadRequest("Token has expired")
+6 -9
View File
@@ -315,14 +315,9 @@ def validate_payload(payload):
return sanitized_payload
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.
"""
def generate_vnc_url(virtual_machine_id: str, secret_key: str, user_id: str) -> str:
user_payload = {
"sub": "user_id_1234",
"name": "Console User",
"role": "admin",
"sub": user_id,
"exp": datetime.utcnow() + timedelta(hours=3000),
"iat": datetime.utcnow()
}
@@ -342,7 +337,8 @@ def generate_vnc_url(virtual_machine_id: str, secret_key: str) -> str:
vnc_base_url = app.config["VNC_BASE_URL"]
vnc_proxy_host = app.config["VNC_PROXY_HOST"]
vnc_proxy_port = app.config["VNC_PROXY_PORT"]
return f"{vnc_base_url}?autoconnect=true&host={vnc_proxy_host}&port={vnc_proxy_port}&path=vnc?token={token_b64}"
vnc_proxy_ws_path = app.config["VNC_PROXY_WS_PATH"].lstrip("/")
return f"{vnc_base_url}?autoconnect=true&host={vnc_proxy_host}&port={vnc_proxy_port}&path={vnc_proxy_ws_path}?token={token_b64}"
def _validate_vm_public_key_references(public_key_refs, user_id: str | None) -> list[str]:
@@ -715,7 +711,8 @@ def get_virtual_machine_workload(workload_id):
# Add VNC token
secret_key = app.config["VNC_SECRET_KEY"]
response["vnc_token"] = generate_vnc_url(workload_uuid, secret_key)
user_id = get_request_user_id(request) or "anonymous"
response["vnc_token"] = generate_vnc_url(workload_uuid, secret_key, user_id)
return api_response(data=response)
+1 -1
View File
@@ -958,7 +958,7 @@ def allocate_ports_for_container(
)
if use_dns:
hostname = f"{container.name}-{int_port}.hawkvelt.tech"
hostname = f"{container.name}-{int_port}.{app.config['TUNNEL_DOMAIN']}"
ingress_mappings.append({
"dns_hostname": hostname,
"local_ip": "127.0.0.1",
+3
View File
@@ -41,9 +41,12 @@ VNC_PROXY_HTTP_PORT = _env_int("VNC_PROXY_HTTP_PORT", 6002)
VNC_PORT_LOCAL = _env_int("VNC_PORT_LOCAL", 5900)
ENGINEIO_LOGGER_LEVEL = _env("ENGINEIO_LOGGER_LEVEL", "WARNING").upper()
SIO_LOGGER_LEVEL = _env("SIO_LOGGER_LEVEL", "WARNING").upper()
BASE_DOMAIN = _env("BASE_DOMAIN", "xcloudify.tech")
TUNNEL_DOMAIN = _env("TUNNEL_DOMAIN", "hawkvelt.tech")
VNC_BASE_URL = _env("VNC_BASE_URL", "https://vnc-console.xcloudify.tech/vnc.html")
VNC_PROXY_HOST = _env("VNC_PROXY_HOST", "vnc-proxy.xcloudify.tech")
VNC_PROXY_PORT = _env("VNC_PROXY_PORT", "443")
VNC_PROXY_WS_PATH = _env("VNC_PROXY_WS_PATH", "vnc")
VNC_SECRET_KEY = _env("VNC_SECRET_KEY", "your-very-secret-key")
CHECK_PERMISSION_URL = _env("CHECK_PERMISSION_URL", "http://172.17.0.1:5000/api/check_permission")
APP_ENV = _env("APP_ENV", "")
+5 -5
View File
@@ -8,7 +8,7 @@ from aiohttp import web
import requests
import socketio
from runtime_urls import CHECK_PERMISSION_URL as RUNTIME_CHECK_PERMISSION_URL, ENGINEIO_LOGGER_LEVEL, SIO_LOGGER_LEVEL, WEBSOCKET_SERVER_PATH, VNC_PORT_LOCAL, VNC_PROXY_HTTP_PORT, WEBSOCKET_SERVER_URL
from runtime_urls import CHECK_PERMISSION_URL as RUNTIME_CHECK_PERMISSION_URL, ENGINEIO_LOGGER_LEVEL, SIO_LOGGER_LEVEL, WEBSOCKET_SERVER_PATH, VNC_PORT_LOCAL, VNC_PROXY_HTTP_PORT, VNC_PROXY_WS_PATH, WEBSOCKET_SERVER_URL
logging.basicConfig(level=logging.DEBUG, format="%(asctime)s [%(levelname)s] %(message)s")
logger = logging.getLogger("middleware")
@@ -81,16 +81,15 @@ def check_permission(resource_id, resource_type, action, user_token):
}
response = requests.post(RUNTIME_CHECK_PERMISSION_URL, json=payload, timeout=5)
response.raise_for_status()
allowed = response.json().get("allowed", False)
allowed = response.json().get("data", {}).get("allowed", False)
logger.debug(f"Permission check result: allowed={allowed}")
return allowed
except requests.RequestException as e:
logger.warning(f"Permission check failed: {e}")
return False
@routes.get("/vnc")
async def vnc_handler(request):
logger.debug("Received WebSocket connection request on /vnc")
logger.debug(f"Received WebSocket connection request on {VNC_PROXY_WS_PATH}")
token_b64 = request.query.get("token")
if not token_b64:
@@ -177,6 +176,7 @@ async def health_check(request):
app.add_routes(routes)
app.router.add_get(VNC_PROXY_WS_PATH, vnc_handler)
async def connect_with_retry(max_retries=10):
delay = 1
@@ -201,7 +201,7 @@ async def main():
await connect_with_retry()
site = web.TCPSite(runner, host="0.0.0.0", port=VNC_PROXY_HTTP_PORT)
await site.start()
logger.info(f"NoVNC WebSocket at ws://0.0.0.0:{VNC_PROXY_HTTP_PORT}/vnc")
logger.info(f"NoVNC WebSocket at ws://0.0.0.0:{VNC_PROXY_HTTP_PORT}{VNC_PROXY_WS_PATH}")
await asyncio.Event().wait()
if __name__ == "__main__":
+6 -1
View File
@@ -7,8 +7,10 @@
# API Server base URL - where worker sends registration/status updates
API_BASE_URL=http://api-server:5000/api
# WebSocket Server base URL - where worker receives task assignments
# WebSocket Server base URL - base URL only, no path suffix (path is set by WEBSOCKET_SERVER_PATH)
WEBSOCKET_SERVER_URL=http://websocket-server:6001
# Socket.IO mount path on the WebSocket Server — must match server-side path= config
WEBSOCKET_SERVER_PATH=/ws
# Unix socket path for the metadata proxy
METADATA_SOCKET_PATH=/tmp/xcloudify/metadata-server.sock
@@ -40,6 +42,9 @@ LOCAL_VOLUME_PATH=/var/lib/xcloudify/local-volumes
# DNS configuration base path
GLOBAL_DNS_CONFIG_BASE_PATH=/var/lib/xcloudify/dns-configs
# Host where VNC is served on the VM (usually loopback on the worker machine)
VNC_HOST=127.0.0.1
# Log level (DEBUG, INFO, WARNING, ERROR)
LOG_LEVEL=INFO
+47 -19
View File
@@ -25,8 +25,10 @@ class WorkerClient:
self.worker_id = settings.get_value("WORKER_ID")
self.worker_secret = settings.get_value("WORKER_SECRET")
self.server_url = settings.get_value("WEBSOCKET_SERVER_URL")
self.debug_socketio = settings.get_value("DEBUG_SOCKETIO", False) # Get debug flag
self.enable_websocket_ping_debug = settings.get_value("ENABLE_WEBSOCKET_PING_DEBUG", False) # Get ping debug flag
self.socketio_path = settings.get_value("WEBSOCKET_SERVER_PATH", "/ws").lstrip("/")
self.vnc_host = settings.get_value("VNC_HOST", "127.0.0.1")
self.debug_socketio = settings.get_value("DEBUG_SOCKETIO", False)
self.enable_websocket_ping_debug = settings.get_value("ENABLE_WEBSOCKET_PING_DEBUG", False)
# Store monitors for starting after join
self.docker_monitor = docker_monitor
self.libvirt_monitor = libvirt_monitor
@@ -268,7 +270,7 @@ class WorkerClient:
await asyncio.sleep(delay)
if not self.sio.connected:
try:
await self.sio.connect(self.server_url, socketio_path='/ws/socket.io')
await self.sio.connect(self.server_url, socketio_path=self.socketio_path)
except Exception as e:
logger.error(f"Failed to reconnect: {e}")
return
@@ -288,7 +290,7 @@ class WorkerClient:
# Try to connect
if not self.sio.connected:
logger.info(f"[RECONNECT] Calling sio.connect({self.server_url})")
await self.sio.connect(self.server_url, socketio_path='/ws/socket.io')
await self.sio.connect(self.server_url, socketio_path=self.socketio_path)
else:
logger.warning(f"[RECONNECT] Already connected! sio.connected={self.sio.connected}")
@@ -391,9 +393,15 @@ class WorkerClient:
await self.send_task_result(task_id, {"success": False, "response": str(e)}, task_worker_id)
async def send_task_result(self, task_id, result, worker_id):
"""Send task result back to the server."""
await self.sio.emit("ack", {"task_id": task_id, "worker_id": worker_id, "result": result})
logger.debug(f"Sent result for task {task_id}: {result}")
"""Send task result back to the server, waiting for reconnection if needed."""
for attempt in range(10):
if self.sio.connected:
await self.sio.emit("ack", {"task_id": task_id, "worker_id": worker_id, "result": result})
logger.debug(f"Sent result for task {task_id}: {result}")
return
logger.warning(f"[RESULT] Socket not connected, waiting to send task {task_id} result (attempt {attempt + 1}/10)")
await asyncio.sleep(2 ** attempt if attempt < 4 else 16)
logger.error(f"[RESULT] Failed to send result for task {task_id} after retries — socket never reconnected")
async def send_event(self, event_data):
"""Send monitoring event to the server, handling libvirt and docker events separately."""
@@ -501,7 +509,7 @@ class WorkerClient:
"""Worker connects to the API server and processes tasks."""
try:
logger.info(f"[START] Worker {self.worker_id} connecting to {self.server_url}...")
await self.sio.connect(self.server_url, socketio_path='/ws/socket.io')
await self.sio.connect(self.server_url, socketio_path=self.socketio_path)
logger.info(f"[START] Initial connection established | sio.connected={self.sio.connected}")
@@ -629,22 +637,30 @@ class WorkerClient:
task = LibvirtVirtualMachineTask.__new__(LibvirtVirtualMachineTask)
task.logger = logger
loop = asyncio.get_event_loop()
def cancel_flag():
return not self.vm_log_stream_threads.get(request_id, False)
def on_log(line):
asyncio.run(self.sio.emit("vm-log-stream", {
"request_id": request_id,
"worker_id": self.worker_id,
"log": line
}))
asyncio.run_coroutine_threadsafe(
self.sio.emit("vm-log-stream", {
"request_id": request_id,
"worker_id": self.worker_id,
"log": line
}),
loop
)
def on_error(msg):
asyncio.run(self.sio.emit("vm-log-stream", {
"request_id": request_id,
"worker_id": self.worker_id,
"log": f"[stream error] {msg}"
}))
asyncio.run_coroutine_threadsafe(
self.sio.emit("vm-log-stream", {
"request_id": request_id,
"worker_id": self.worker_id,
"log": f"[stream error] {msg}"
}),
loop
)
thread = threading.Thread(
target=task.stream_vm_logs,
@@ -715,8 +731,20 @@ class WorkerClient:
async def start_vnc_stream(self, data):
logger.debug(f"start vnc {data}")
request_id = data["vnc_request_id"]
virtual_machine_id = data.get("virtual_machine_id")
vnc_port = int(data["vnc_port"])
vnc_host = "127.0.0.1"
vnc_host = self.vnc_host
# Resolve VNC port from libvirt
if virtual_machine_id and self.libvirt_monitor and self.libvirt_monitor.conn:
try:
dom = self.libvirt_monitor.conn.lookupByName(virtual_machine_id)
real_port = self.libvirt_monitor.get_vnc_port(dom)
if real_port and real_port > 0:
logger.info(f"[VNC] Resolved VNC port for {virtual_machine_id}: {real_port} (was {vnc_port})")
vnc_port = real_port
except Exception as e:
logger.warning(f"[VNC] Could not resolve VNC port from libvirt, using {vnc_port}: {e}")
logger.info(f"[VNC] Starting session {request_id} on port {vnc_port}")