diff --git a/.env.example b/.env.example index fcbb1e6..86a23ea 100644 --- a/.env.example +++ b/.env.example @@ -1,71 +1,44 @@ - -# ----------------------------------------------------------------------------- -# App -# ----------------------------------------------------------------------------- - # Environment: "development" (default) enables /apidocs + /apispec.json. # Set to "production" to disable Swagger UI entirely. -APP_ENV=development - -# ----------------------------------------------------------------------------- -# Database -# ----------------------------------------------------------------------------- +FLASK_ENV=development # MySQL connection string used by Flask-SQLAlchemy and the websocket server. -SQLALCHEMY_DATABASE_URI=mysql://root:password@172.17.0.1:3306/theapi - -# ----------------------------------------------------------------------------- -# Redis -# ----------------------------------------------------------------------------- +DATABASE_URL=mysql://theapi_user:theapi_password@172.17.0.1:3306/theapi # DB 0 — Celery broker -BROKER_URL=redis://172.17.0.1:6379/0 +REDIS_BROKER_URL=redis://172.17.0.1:6379/0 # DB 1 — Celery result backend -RESULT_BACKEND=redis://172.17.0.1:6379/1 +REDIS_RESULT_BACKEND_URL=redis://172.17.0.1:6379/1 # DB 2 — Operational tasks (websocket server, PING logging) REDIS_URL=redis://172.17.0.1:6379/2 -# ----------------------------------------------------------------------------- -# Celery task scheduler — retry / backoff tuning (all optional) -# ----------------------------------------------------------------------------- - SCHEDULER_MAX_ALLOCATION_ATTEMPTS=8 SCHEDULER_RETRY_BASE_DELAY_SECONDS=60 SCHEDULER_RETRY_BACKOFF_FACTOR=2.0 SCHEDULER_RETRY_JITTER_SECONDS=30 SCHEDULER_RETRY_MAX_DELAY_SECONDS=1800 -# ----------------------------------------------------------------------------- -# Cloudflare -# ----------------------------------------------------------------------------- - CLOUDFLARE_API_TOKEN=your-cloudflare-api-token CLOUDFLARE_ACCOUNT_ID=your-cloudflare-account-id CLOUDFLARE_ZONE_ID=your-cloudflare-zone-id -# ----------------------------------------------------------------------------- -# WebSocket server -# ----------------------------------------------------------------------------- +API_BASE_URL=http://172.17.0.1:5000/api -# URL the API uses to dispatch tasks to the websocket server -WEBSOCKET_SERVER_URL=http://172.17.0.1:6001/api/create_task +# WebSocket Server base URL +WEBSOCKET_SERVER_URL=http://172.17.0.1:6001 -# URL used by the Streamlit dashboard and VNC proxy to reach the socket server -SOCKET_SERVER_URL=http://172.17.0.1:6001 +# WebSocket task creation endpoint (full path, used by api-server to dispatch tasks) +WEBSOCKET_TASK_URL=http://172.17.0.1:6001/api/create_task -# Seconds without a PING/PONG before a worker is marked offline +# Ping heartbeat configuration PING_HEARTBEAT_TIMEOUT_SECONDS=30 # Debug flags for the websocket server (set to "true" to enable) ENABLE_WEBSOCKET_PING_DEBUG=false ENABLE_TASK_ASSIGNMENT_DEBUG=false -# ----------------------------------------------------------------------------- -# VNC console -# ----------------------------------------------------------------------------- - VNC_SECRET_KEY=change-me-to-a-strong-random-secret VNC_BASE_URL=https://vnc-console.xcloudify.tech/vnc.html VNC_PROXY_HOST=vnc-proxy.xcloudify.tech @@ -80,17 +53,9 @@ ENGINEIO_LOGGER_LEVEL=WARNING SIO_LOGGER_LEVEL=WARNING # URL the VNC proxy calls to verify session permissions -CHECK_PERMISSION_URL=http://localhost:5000/api/check_permission - -# ----------------------------------------------------------------------------- -# Auth / JWT -# ----------------------------------------------------------------------------- +CHECK_PERMISSION_URL=http://172.17.0.1:5000/api/check_permission JWT_SECRET_KEY=change-me-to-a-strong-random-secret -# ----------------------------------------------------------------------------- -# Networking -# ----------------------------------------------------------------------------- - # Default OVS bridge name used when creating networks (optional) XCF_DEFAULT_OVS_BRIDGE=br-xcloudify \ No newline at end of file diff --git a/.gitignore b/.gitignore index c92c866..e66d10f 100644 --- a/.gitignore +++ b/.gitignore @@ -46,6 +46,7 @@ vault_pass.txt *.vault group_vars/*/vault.yml host_vars/*/vault.yml +.env # Test artifacts test_results/ diff --git a/api_client/__init__.py b/api_client/__init__.py index e91553a..088b12d 100644 --- a/api_client/__init__.py +++ b/api_client/__init__.py @@ -1,3 +1,5 @@ from .client import IaaSClient, XCloudifyAPIError +__version__ = "0.0.0" + __all__ = ['IaaSClient', 'XCloudifyAPIError'] \ No newline at end of file diff --git a/api_client/client.py b/api_client/client.py index 25d27e5..32ce105 100644 --- a/api_client/client.py +++ b/api_client/client.py @@ -1,6 +1,8 @@ # api_client/client.py import requests import time + +from runtime_urls import API_BASE_URL class XCloudifyAPIError(Exception): """ @@ -16,16 +18,17 @@ class XCloudifyAPIError(Exception): self.details = details or {} class IaaSClient: - def __init__(self, api_key=None, access_token=None, refresh_token=None, base_url="http://172.17.0.1:5000/api", verify_ssl=True): + def __init__(self, api_key=None, access_token=None, refresh_token=None, base_url=API_BASE_URL, verify_ssl=True): """ Initialize the IaaSClient. :param api_key: API key for X-API-KEY auth (legacy, use access_token for Bearer) :param access_token: Bearer access token (preferred) :param refresh_token: Refresh token for auto-refresh on 401 - :param base_url: Base API URL (default: http://172.17.0.1:5000/api) + :param base_url: Base API URL (default: shared runtime API URL) :param verify_ssl: Verify SSL certificates (default: True; set False for self-signed) """ + self.api_key = api_key self.base_url = base_url.rstrip('/') self.verify_ssl = verify_ssl self.refresh_token = refresh_token diff --git a/app/__init__.py b/app/__init__.py index b38a641..4d00c70 100644 --- a/app/__init__.py +++ b/app/__init__.py @@ -11,7 +11,6 @@ print("-----In init----------") import uuid import os import pymysql -from dotenv import load_dotenv from flask import Flask, g, jsonify, request from flask_sqlalchemy import SQLAlchemy from flask_migrate import Migrate @@ -21,11 +20,7 @@ from logger import logger from app.utils.standard_responses import api_response from flask_cors import CORS -# --------------------------------------------------------------------------- # -# 0. Environment setup # -# --------------------------------------------------------------------------- # -load_dotenv() -logger.info("Environment variables loaded from .env") +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_TASK_URL # --------------------------------------------------------------------------- # # 1. Flask application & database # @@ -34,7 +29,7 @@ app: Flask = Flask(__name__) CORS(app) # Set APP_ENV=development in prod .env to disable /apidocs and /apispec.json. -_env = os.getenv("APP_ENV", "").lower() +_env = APP_ENV.lower() _swagger_enabled = _env == "development" if _swagger_enabled: @@ -76,20 +71,12 @@ if _swagger_enabled: else: logger.info("Swagger UI disabled (APP_ENV=%s)", _env) -app.config.update( - SQLALCHEMY_DATABASE_URI="mysql://root:password@172.17.0.1:3306/theapi", - WEBSOCKET_SERVER_URL="http://172.17.0.1:6001/api/create_task", - - # Celery settings (override via ENV) - use new style keys to avoid mix error - broker_url=os.getenv("BROKER_URL", os.getenv("CELERY_BROKER_URL", "redis://172.17.0.1:6379/0")), # Redis DB 0 for Celery broker - result_backend=os.getenv("RESULT_BACKEND", os.getenv("CELERY_RESULT_BACKEND", "redis://172.17.0.1:6379/1")), # Redis DB 1 for Celery results -) -# Scheduler retry/backoff defaults (tunable via ENV) -app.config["SCHEDULER_MAX_ALLOCATION_ATTEMPTS"] = int(os.getenv("SCHEDULER_MAX_ALLOCATION_ATTEMPTS", 8)) -app.config["SCHEDULER_RETRY_BASE_DELAY_SECONDS"] = int(os.getenv("SCHEDULER_RETRY_BASE_DELAY_SECONDS", 60)) -app.config["SCHEDULER_RETRY_BACKOFF_FACTOR"] = float(os.getenv("SCHEDULER_RETRY_BACKOFF_FACTOR", 2.0)) -app.config["SCHEDULER_RETRY_JITTER_SECONDS"] = int(os.getenv("SCHEDULER_RETRY_JITTER_SECONDS", 30)) -app.config["SCHEDULER_RETRY_MAX_DELAY_SECONDS"] = int(os.getenv("SCHEDULER_RETRY_MAX_DELAY_SECONDS", 1800)) +app.config.update(SQLALCHEMY_DATABASE_URI=DATABASE_URL, WEBSOCKET_SERVER_URL=WEBSOCKET_TASK_URL, broker_url=REDIS_BROKER_URL, result_backend=REDIS_RESULT_BACKEND_URL) +app.config["SCHEDULER_MAX_ALLOCATION_ATTEMPTS"] = SCHEDULER_MAX_ALLOCATION_ATTEMPTS +app.config["SCHEDULER_RETRY_BASE_DELAY_SECONDS"] = SCHEDULER_RETRY_BASE_DELAY_SECONDS +app.config["SCHEDULER_RETRY_BACKOFF_FACTOR"] = SCHEDULER_RETRY_BACKOFF_FACTOR +app.config["SCHEDULER_RETRY_JITTER_SECONDS"] = SCHEDULER_RETRY_JITTER_SECONDS +app.config["SCHEDULER_RETRY_MAX_DELAY_SECONDS"] = SCHEDULER_RETRY_MAX_DELAY_SECONDS # --------------------------------------------------------------------------- # # 2. Cloudflare configuration # @@ -109,17 +96,17 @@ def _mask(value: str, visible: int = 4) -> str: return "" return f"{value[:visible]}{'*' * (len(value) - visible)}" -app.config["CLOUDFLARE_API_TOKEN"] = os.getenv("CLOUDFLARE_API_TOKEN") -app.config["CLOUDFLARE_ACCOUNT_ID"] = os.getenv("CLOUDFLARE_ACCOUNT_ID") -app.config["CLOUDFLARE_ZONE_ID"] = os.getenv("CLOUDFLARE_ZONE_ID") -app.config["PING_HEARTBEAT_TIMEOUT_SECONDS"] = int(os.getenv("PING_HEARTBEAT_TIMEOUT_SECONDS",30)) # How long before a Websocket PING\PONG is classed as a failure and triggers a worker offline event -app.config["REDIS_URL"] = os.getenv("REDIS_URL","redis://172.17.0.1:6379/2") # Redis Database 2 for operational tasks like websocket server and PING logging -app.config["VNC_SECRET_KEY"] = os.getenv("VNC_SECRET_KEY","your-very-secret-key") -app.config["VNC_BASE_URL"] = os.getenv("VNC_BASE_URL", "https://vnc-console.xcloudify.tech/vnc.html") -app.config["VNC_PROXY_HOST"] = os.getenv("VNC_PROXY_HOST", "vnc-proxy.xcloudify.tech") -app.config["VNC_PROXY_PORT"] = os.getenv("VNC_PROXY_PORT", "443") +app.config["CLOUDFLARE_API_TOKEN"] = CLOUDFLARE_API_TOKEN +app.config["CLOUDFLARE_ACCOUNT_ID"] = CLOUDFLARE_ACCOUNT_ID +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["VNC_BASE_URL"] = VNC_BASE_URL +app.config["VNC_PROXY_HOST"] = VNC_PROXY_HOST +app.config["VNC_PROXY_PORT"] = VNC_PROXY_PORT # JWT secret used for decoding Authorization bearer tokens for audit attribution -app.config["JWT_SECRET_KEY"] = os.getenv("JWT_SECRET_KEY", "your-very-secret-key") +app.config["JWT_SECRET_KEY"] = JWT_SECRET_KEY logger.info( "Cloudflare configuration set " diff --git a/app/controller/api/workload_vm_routes.py b/app/controller/api/workload_vm_routes.py index cf3d922..9107610 100644 --- a/app/controller/api/workload_vm_routes.py +++ b/app/controller/api/workload_vm_routes.py @@ -16,7 +16,9 @@ from datetime import datetime, timedelta import jwt import base64 from app.utils.auth_utils import get_request_user_id -websocket_server_url = "http://172.17.0.1:6001/api/create_task" +from runtime_urls import WEBSOCKET_TASK_URL + +websocket_server_url = app.config.get("WEBSOCKET_SERVER_URL", WEBSOCKET_TASK_URL) def validate_port_entry(port_entry): """ diff --git a/app/controller/ui_routes.py b/app/controller/ui_routes.py index c143be7..27c2126 100644 --- a/app/controller/ui_routes.py +++ b/app/controller/ui_routes.py @@ -1,9 +1,9 @@ from flask import request, jsonify, render_template, g, url_for -import requests from app.models.models import Universe, Project, VirtualDataCenter, Region, Workload from app.controller import ui_bp from app.models.network import * from sqlalchemy import or_ +from runtime_urls import API_BASE_URL # UI Routes with prefix @@ -56,16 +56,6 @@ def container_logs_specific(container_id): @ui_bp.route("/vnc.html") def console(): """Render the noVNC client page for a specific VM""" - # Fetch the WebSocket URL from the API - API_SERVER_URL = "http://localhost:5000" # API to get VNC session token - - # response = requests.post(f"{API_SERVER_URL}/get-console-url", json={"username": "test_user", "worker_id": vm_id}) - - # if response.status_code != 200: - # return "Error: Unable to get VNC session", 500 - - # console_data = response.json() - # console_url = console_data.get("console_url") # Get VNC connection parameters from request args or use defaults websocket_host = '127.0.0.1' websocket_port = 6080 # Standard noVNC port diff --git a/app/utils/create_workload_virtual_machine.py b/app/utils/create_workload_virtual_machine.py index cf0ca07..a3208d2 100644 --- a/app/utils/create_workload_virtual_machine.py +++ b/app/utils/create_workload_virtual_machine.py @@ -15,8 +15,9 @@ from app.utils.sdn_helpers import send_sdn_updates_for_networks from app.utils.dns_helpers import send_dns_updates_for_vdc from app.utils.standard_responses import api_response from app.compute.routes.gpu import schedule_gpu_for_vm +from runtime_urls import WEBSOCKET_TASK_URL -websocket_server_url = "http://172.17.0.1:6001/api/create_task" +websocket_server_url = app.config.get("WEBSOCKET_SERVER_URL", WEBSOCKET_TASK_URL) def add_VirtualMachine_workload(request_data): diff --git a/app/utils/sdn_helpers.py b/app/utils/sdn_helpers.py index bef07ad..f10754d 100644 --- a/app/utils/sdn_helpers.py +++ b/app/utils/sdn_helpers.py @@ -5,6 +5,7 @@ from app.models.network import NetworkPort from flask import json import requests from app.utils.standard_responses import api_response +from runtime_urls import WEBSOCKET_TASK_URL from typing import List, Optional from app import db, logger @@ -42,7 +43,7 @@ def send_sdn_updates_for_networks(network_ids: List[str], return_task_for_host: logger.info("Hosts involved in network update: %s", list(host_ids)) return_task_id = None - websocket_server_url = "http://172.17.0.1:6001/api/create_task" + websocket_server_url = app.config.get("WEBSOCKET_SERVER_URL", WEBSOCKET_TASK_URL) headers = {"Content-Type": "application/json"} for host_id in host_ids: diff --git a/docker-compose.yml b/docker-compose.yml index 90a1cc0..3210ecf 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -45,6 +45,7 @@ services: context: ./app image: xcloudify-flask-base container_name: db-migrate + env_file: .env volumes: - .:/app working_dir: /app @@ -59,6 +60,10 @@ services: context: ./app image: xcloudify-flask-base container_name: api-server + env_file: .env + environment: + FLASK_APP: api_server.py + FLASK_ENV: ${FLASK_ENV:-development} volumes: - .:/app working_dir: /app @@ -66,9 +71,6 @@ services: sh -c "flask run --host=0.0.0.0 --port=5000 --debug" ports: - "5000:5000" - environment: - FLASK_APP: api_server.py - FLASK_ENV: development depends_on: mariadb: condition: service_healthy @@ -90,6 +92,7 @@ services: command: python3 -m websocket_server.main ports: - "6001:6001" + env_file: .env volumes: - .:/app working_dir: /app @@ -144,6 +147,7 @@ services: context: app/ container_name: celery-worker command: celery -A app.celery_app.celery worker --loglevel=info + env_file: .env volumes: - .:/app working_dir: /app @@ -158,6 +162,7 @@ services: context: app/ container_name: celery-beat command: celery -A app.celery_app.celery beat --loglevel=info + env_file: .env volumes: - .:/app working_dir: /app diff --git a/run-master.sh b/run-master.sh new file mode 100755 index 0000000..374a6d7 --- /dev/null +++ b/run-master.sh @@ -0,0 +1,59 @@ +#!/usr/bin/env bash + +set -e + + +show_help() { + echo "Usage: $0 [OPTIONS]" + echo "" + echo "Options:" + echo " --build Rebuild containers without cache and start services" + echo " --clean Remove containers, volumes, images, and rebuild everything" + echo " --help Show this help message" +} + +# Detect docker compose command +if command -v docker-compose >/dev/null 2>&1; then + DOCKER_COMPOSE="docker-compose" +elif docker compose version >/dev/null 2>&1; then + DOCKER_COMPOSE="docker compose" +else + echo "Error: Neither docker-compose nor docker compose is installed." + exit 1 +fi + +case "$1" in + --help) + show_help + ;; + + --build) + echo "Building containers without cache..." + $DOCKER_COMPOSE build --no-cache + echo "Starting services..." + $DOCKER_COMPOSE up -d + ;; + + --clean) + echo "Cleaning containers, volumes, images, and orphans..." + $DOCKER_COMPOSE down -v --rmi all --remove-orphans + + echo "Rebuilding containers without cache..." + $DOCKER_COMPOSE build --no-cache + + echo "Starting services..." + $DOCKER_COMPOSE up -d + ;; + + "") + echo "Starting services..." + $DOCKER_COMPOSE up -d + ;; + + *) + echo "Unknown option: $1" + echo "" + show_help + exit 1 + ;; +esac \ No newline at end of file diff --git a/runtime_urls.py b/runtime_urls.py new file mode 100644 index 0000000..5f05219 --- /dev/null +++ b/runtime_urls.py @@ -0,0 +1,63 @@ +"""Centralized runtime URL and configuration values for the main app. + +This module loads the root .env file once and exposes resolved values so the +rest of the app can import config instead of reading environment variables at +each call site. +""" + +import os + +from dotenv import load_dotenv + +load_dotenv() + + +def _env(name: str, default: str) -> str: + return os.getenv(name, default) + + +def _env_bool(name: str, default: bool = False) -> bool: + value = os.getenv(name) + return default if value is None else value.lower() in ("true", "1", "t", "yes", "y") + + +def _env_int(name: str, default: int) -> int: + value = os.getenv(name) + return default if value is None else int(value) + + +def _env_float(name: str, default: float) -> float: + value = os.getenv(name) + return default if value is None else float(value) + +DATABASE_URL = _env("DATABASE_URL", "mysql://theapi_user:theapi_password@172.17.0.1:3306/theapi") +REDIS_URL = _env("REDIS_URL", "redis://172.17.0.1:6379/2") +REDIS_BROKER_URL = _env("BROKER_URL", _env("CELERY_BROKER_URL", "redis://172.17.0.1:6379/0")) +REDIS_RESULT_BACKEND_URL = _env("RESULT_BACKEND", _env("CELERY_RESULT_BACKEND", "redis://172.17.0.1:6379/1")) +API_BASE_URL = _env("API_BASE_URL", "http://172.17.0.1:5000/api") +WEBSOCKET_SERVER_URL = _env("WEBSOCKET_SERVER_URL", "http://172.17.0.1:6001") +WEBSOCKET_TASK_URL = _env("WEBSOCKET_TASK_URL", "http://172.17.0.1:6001/api/create_task") +WEBSOCKET_SERVER_PATH = _env("WEBSOCKET_SERVER_PATH", "/ws") +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() +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_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", "") +FLASK_ENV = _env("FLASK_ENV", "development") +JWT_SECRET_KEY = _env("JWT_SECRET_KEY", "your-very-secret-key") +SCHEDULER_MAX_ALLOCATION_ATTEMPTS = _env_int("SCHEDULER_MAX_ALLOCATION_ATTEMPTS", 8) +SCHEDULER_RETRY_BASE_DELAY_SECONDS = _env_int("SCHEDULER_RETRY_BASE_DELAY_SECONDS", 60) +SCHEDULER_RETRY_BACKOFF_FACTOR = _env_float("SCHEDULER_RETRY_BACKOFF_FACTOR", 2.0) +SCHEDULER_RETRY_JITTER_SECONDS = _env_int("SCHEDULER_RETRY_JITTER_SECONDS", 30) +SCHEDULER_RETRY_MAX_DELAY_SECONDS = _env_int("SCHEDULER_RETRY_MAX_DELAY_SECONDS", 1800) +CLOUDFLARE_API_TOKEN = _env("CLOUDFLARE_API_TOKEN", "") +CLOUDFLARE_ACCOUNT_ID = _env("CLOUDFLARE_ACCOUNT_ID", "") +CLOUDFLARE_ZONE_ID = _env("CLOUDFLARE_ZONE_ID", "") +PING_HEARTBEAT_TIMEOUT_SECONDS = _env_int("PING_HEARTBEAT_TIMEOUT_SECONDS", 30) +ENABLE_TASK_ASSIGNMENT_DEBUG = _env_bool("ENABLE_TASK_ASSIGNMENT_DEBUG", False) +ENABLE_WEBSOCKET_PING_DEBUG = _env_bool("ENABLE_WEBSOCKET_PING_DEBUG", False) \ No newline at end of file diff --git a/sandbox/task_creator.py b/sandbox/task_creator.py index 9de05cd..f1c8844 100644 --- a/sandbox/task_creator.py +++ b/sandbox/task_creator.py @@ -22,7 +22,9 @@ redis_client = redis.StrictRedis(host="localhost", port=6379, decode_responses=T # SQLAlchemy setup Base = declarative_base() -DATABASE_URL = "mysql://root:password@172.17.0.1:3306/defaultdb" +from runtime_urls import DATABASE_URL as RUNTIME_DATABASE_URL + +DATABASE_URL = RUNTIME_DATABASE_URL.replace("theapi", "defaultdb") engine = create_engine(DATABASE_URL) Session = sessionmaker(bind=engine) diff --git a/streamlit_server/views/websocket.py b/streamlit_server/views/websocket.py index 414d94a..a761f14 100644 --- a/streamlit_server/views/websocket.py +++ b/streamlit_server/views/websocket.py @@ -1,11 +1,10 @@ import streamlit as st import requests import socketio -import os -# Base URL pulled from environment variable -BASE_URL = os.getenv("SOCKET_SERVER_URL", "http://172.17.0.1:6001") -API_URL = f"{BASE_URL}/api/connected_clients" +from runtime_urls import WEBSOCKET_SERVER_URL + +API_URL = f"{WEBSOCKET_SERVER_URL}/api/connected_clients" def load_clients(): response = requests.get(API_URL) @@ -16,7 +15,7 @@ def load_clients(): def render(): # Initialize SocketIO client sio = socketio.Client() - sio.connect(BASE_URL) + sio.connect(WEBSOCKET_SERVER_URL) st.title("Connected Clients") diff --git a/tests/client_test.py b/tests/client_test.py index c9c8984..de9f309 100644 --- a/tests/client_test.py +++ b/tests/client_test.py @@ -13,12 +13,12 @@ def mock_client(): def test_create_container_validation(mock_client): """Test create_container with valid data passes validation.""" data = { - "virtual_data_center": "123e4567-e89b-12d3-a456-426614174000", + "vdc_id": "123e4567-e89b-12d3-a456-426614174000", "containers": [{ "docker_image": "nginx:latest", "container_name": "test-nginx", "cpu": 1.0, - "mem_limit": 128, + "memory_mb": 128, "ports": [{"internal": 80, "external": 8080, "use_dns": False}], "storage": [{"volume_id": "vol-123", "mount_point": "/data", "read_only": False}], "networks": ["default"], @@ -36,8 +36,8 @@ def test_create_container_validation(mock_client): def test_create_container_invalid_data(mock_client): """Test create_container raises on invalid data.""" - invalid_data = {"virtual_data_center": "invalid-uuid", "containers": []} - with pytest.raises(XCloudifyAPIError, match="Validation failed"): + invalid_data = {"containers": []} + with pytest.raises(XCloudifyAPIError, match="Invalid data"): mock_client.create_container(invalid_data) def test_patch_container(mock_client): diff --git a/vnc_proxy.py b/vnc_proxy.py index e4856ce..43ca652 100644 --- a/vnc_proxy.py +++ b/vnc_proxy.py @@ -7,17 +7,8 @@ import uuid from aiohttp import web import requests import socketio -import os -# Config -# Config from environment -HTTP_PORT = int(os.getenv("VNC_PROXY_HTTP_PORT", 6002)) -VNC_PORT_LOCAL = int(os.getenv("VNC_PORT_LOCAL", 5900)) -SOCKET_SERVER_URL = os.getenv("SOCKET_SERVER_URL", "http://172.17.0.1:6001") -SOCKET_SERVER_PATH = os.getenv("WEBSOCKET_SERVER_PATH", "/ws") -CHECK_PERMISSION_URL = os.getenv("CHECK_PERMISSION_URL", "http://localhost:5000/api/check_permission") -ENGINEIO_LOGGER_LEVEL = os.getenv("ENGINEIO_LOGGER_LEVEL", "WARNING").upper() -SIO_LOGGER_LEVEL = os.getenv("SIO_LOGGER_LEVEL", "WARNING").upper() +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 logging.basicConfig(level=logging.DEBUG, format="%(asctime)s [%(levelname)s] %(message)s") logger = logging.getLogger("middleware") @@ -88,7 +79,7 @@ def check_permission(resource_id, resource_type, action, user_token): "resource_type": resource_type, "action": action } - response = requests.post(CHECK_PERMISSION_URL, json=payload, timeout=5) + response = requests.post(RUNTIME_CHECK_PERMISSION_URL, json=payload, timeout=5) response.raise_for_status() allowed = response.json().get("allowed", False) logger.debug(f"Permission check result: allowed={allowed}") @@ -189,11 +180,11 @@ app.add_routes(routes) async def connect_with_retry(max_retries=10): delay = 1 - socketio_path = SOCKET_SERVER_PATH.lstrip("/") or "socket.io" + socketio_path = WEBSOCKET_SERVER_PATH.lstrip("/") or "socket.io" for attempt in range(1, max_retries + 1): try: logger.info(f"Connecting to WebSocket server (attempt {attempt})...") - await sio.connect(SOCKET_SERVER_URL, socketio_path=socketio_path) + await sio.connect(WEBSOCKET_SERVER_URL, socketio_path=socketio_path) logger.info("WebSocket connection established") return except Exception as e: @@ -208,9 +199,9 @@ async def main(): runner = web.AppRunner(app) await runner.setup() await connect_with_retry() - site = web.TCPSite(runner, host="0.0.0.0", port=HTTP_PORT) + 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:{HTTP_PORT}/vnc") + logger.info(f"NoVNC WebSocket at ws://0.0.0.0:{VNC_PROXY_HTTP_PORT}/vnc") await asyncio.Event().wait() if __name__ == "__main__": diff --git a/websocket_server/config.py b/websocket_server/config.py index 678498f..fb00d7c 100644 --- a/websocket_server/config.py +++ b/websocket_server/config.py @@ -9,8 +9,7 @@ from sqlalchemy import create_engine from sqlalchemy.orm import sessionmaker, declarative_base import pymysql -# Redis configuration -REDIS_URL = "redis://172.17.0.1:6379/2" +from runtime_urls import API_BASE_URL, DATABASE_URL, ENABLE_TASK_ASSIGNMENT_DEBUG, ENABLE_WEBSOCKET_PING_DEBUG, REDIS_URL # Configurable intervals @@ -21,12 +20,6 @@ PING_EXPIRY_SECONDS = 5 # How long between PING and PONG before we ASSIGN_INTERVAL_SECONDS = 30 LIVENESS_CHECK_INTERVAL_SECONDS = 5 # How often to cross check connected_workers with ping results from Redis -ENABLE_WEBSOCKET_PING_DEBUG = os.getenv('ENABLE_WEBSOCKET_PING_DEBUG', 'false').lower() == 'true' - -ENABLE_TASK_ASSIGNMENT_DEBUG = os.getenv('ENABLE_TASK_ASSIGNMENT_DEBUG', 'false').lower() == 'true' -# Database configuration -DATABASE_URL = "mysql://root:password@172.17.0.1:3306/theapi" - # Connection pool for Redis redis_connection_pool = redis.ConnectionPool.from_url( REDIS_URL, diff --git a/websocket_server/events/docker.py b/websocket_server/events/docker.py index 2d6225f..8486890 100644 --- a/websocket_server/events/docker.py +++ b/websocket_server/events/docker.py @@ -6,13 +6,12 @@ import logging import requests import json +from websocket_server.config import API_BASE_URL from websocket_server.events import base logger = logging.getLogger("websocket_server") socketio = base.socketio -api_server_url = "http://172.17.0.1:5000/api" - def register_socketio_handlers(socketio): @@ -91,7 +90,7 @@ def register_socketio_handlers(socketio): headers = {"Content-Type": "application/json"} websocket_server_response = requests.put( - f"{api_server_url}/{action_URL}", + f"{API_BASE_URL}/{action_URL}", data=json.dumps(payload), headers=headers, ) diff --git a/websocket_server/events/libvirt.py b/websocket_server/events/libvirt.py index 3b9daed..e216885 100644 --- a/websocket_server/events/libvirt.py +++ b/websocket_server/events/libvirt.py @@ -6,13 +6,12 @@ import logging import requests import json +from websocket_server.config import API_BASE_URL from websocket_server.events import base logger = logging.getLogger("websocket_server") socketio = base.socketio -api_server_url = "http://172.17.0.1:5000/api" - def register_socketio_handlers(socketio): @socketio.on("libvirt_event") @@ -60,7 +59,7 @@ def register_socketio_handlers(socketio): } headers = {"Content-Type": "application/json"} - websocket_server_response = requests.put(f"{api_server_url}/{action_URL}", data=json.dumps(payload), headers=headers) + websocket_server_response = requests.put(f"{API_BASE_URL}/{action_URL}", data=json.dumps(payload), headers=headers) logger.info("Status update sent to API server") if event_type=="libvirt_stopped": @@ -80,7 +79,7 @@ def register_socketio_handlers(socketio): } headers = {"Content-Type": "application/json"} - websocket_server_response = requests.put(f"{api_server_url}/{action_URL}", data=json.dumps(payload), headers=headers) + websocket_server_response = requests.put(f"{API_BASE_URL}/{action_URL}", data=json.dumps(payload), headers=headers) logger.info("Status update sent to API server") @@ -101,7 +100,7 @@ def register_socketio_handlers(socketio): } headers = {"Content-Type": "application/json"} - websocket_server_response = requests.put(f"{api_server_url}/{action_URL}", data=json.dumps(payload), headers=headers) + websocket_server_response = requests.put(f"{API_BASE_URL}/{action_URL}", data=json.dumps(payload), headers=headers) logger.info("Status update sent to API server") except Exception as e: diff --git a/websocket_server/events/logs.py b/websocket_server/events/logs.py index 2cca68f..3f9b170 100644 --- a/websocket_server/events/logs.py +++ b/websocket_server/events/logs.py @@ -6,15 +6,13 @@ import logging import json import requests -from websocket_server.config import get_redis_client +from websocket_server.config import get_redis_client, API_BASE_URL from websocket_server.events import base from websocket_server.shared_state import connected_sids, connected_sids_lock, connected_workers logger = logging.getLogger("websocket_server") socketio = base.socketio - -api_server_url = "http://172.17.0.1:5000/api" def register_socketio_handlers(socketio): @socketio.on("user_request_container_logs") @@ -27,7 +25,7 @@ def register_socketio_handlers(socketio): logger.info(f"[{request_id}] User {user_id} requested logs for container {container_id}") try: - container_response = requests.get(f"{api_server_url}/workloads/containers/{container_id}") + container_response = requests.get(f"{API_BASE_URL}/workloads/containers/{container_id}") if container_response.status_code != 200: emit("user_log_response", { "success": False, @@ -108,7 +106,7 @@ def register_socketio_handlers(socketio): try: logger.info(f"[{request_id}] Starting log stream for container {container_id}") - container_response = requests.get(f"{api_server_url}/workloads/containers/{container_id}") + container_response = requests.get(f"{API_BASE_URL}/workloads/containers/{container_id}") container_response.raise_for_status() container_info = container_response.json()['data'] @@ -205,7 +203,7 @@ def register_socketio_handlers(socketio): logger.info(f"[{request_id}] User '{user_id}' requested VM log stream for VM '{vm_id}'") # Look up which worker hosts this VM - vm_response = requests.get(f"{api_server_url}/workloads/virtual_machines/{vm_id}") + vm_response = requests.get(f"{API_BASE_URL}/workloads/virtual_machines/{vm_id}") if vm_response.status_code != 200: emit("user_vm_log_stream_error", { "success": False, diff --git a/websocket_server/events/terminal.py b/websocket_server/events/terminal.py index 5933c37..4b035fe 100644 --- a/websocket_server/events/terminal.py +++ b/websocket_server/events/terminal.py @@ -6,14 +6,13 @@ import logging import requests from flask import request -from websocket_server.config import get_redis_client +from websocket_server.config import get_redis_client, API_BASE_URL from websocket_server.events import base from websocket_server.shared_state import connected_workers logger = logging.getLogger("websocket_server") socketio = base.socketio -api_server_url = "http://172.17.0.1:5000/api" def register_socketio_handlers(socketio): @socketio.on("start_terminal_session") @@ -26,7 +25,7 @@ def register_socketio_handlers(socketio): try: # Lookup container details - response = requests.get(f"{api_server_url}/workloads/containers/{container_id}") + response = requests.get(f"{API_BASE_URL}/workloads/containers/{container_id}") response.raise_for_status() info = response.json()['data'] worker_id = info["workload_host_id"] diff --git a/websocket_server/events/vnc.py b/websocket_server/events/vnc.py index aab78c5..942fa8a 100644 --- a/websocket_server/events/vnc.py +++ b/websocket_server/events/vnc.py @@ -5,6 +5,7 @@ import logging import requests from flask import request +from websocket_server.config import API_BASE_URL from websocket_server.events import base from websocket_server.shared_state import connected_workers @@ -12,8 +13,6 @@ from websocket_server.shared_state import connected_workers logger = logging.getLogger("websocket_server") socketio = base.socketio -api_server_url = "http://172.17.0.1:5000/api" - # In-memory request map: vnc_request_id -> {worker_id, vnc_proxy_sid, vm_id} vnc_request_worker_map = {} def register_socketio_handlers(socketio): @@ -59,7 +58,7 @@ def register_socketio_handlers(socketio): return try: - vm_url = f"{api_server_url}/workloads/virtual_machines/{virtual_machine_id}" + vm_url = f"{API_BASE_URL}/workloads/virtual_machines/{virtual_machine_id}" response = requests.get(vm_url, timeout=5) response.raise_for_status() vm_info = response.json() diff --git a/websocket_server/worker_manager.py b/websocket_server/worker_manager.py index 910059e..8425c19 100644 --- a/websocket_server/worker_manager.py +++ b/websocket_server/worker_manager.py @@ -7,7 +7,7 @@ import threading import logging import json from websocket_server.shared_state import connected_workers, worker_lock, ping_tracker -from websocket_server.config import get_redis_client,PING_INTERVAL_SECONDS, ASSIGN_INTERVAL_SECONDS, LIVENESS_CHECK_INTERVAL_SECONDS, PING_EXPIRY_SECONDS, ENABLE_WEBSOCKET_PING_DEBUG, ENABLE_TASK_ASSIGNMENT_DEBUG +from websocket_server.config import get_redis_client, PING_INTERVAL_SECONDS, ASSIGN_INTERVAL_SECONDS, LIVENESS_CHECK_INTERVAL_SECONDS, PING_EXPIRY_SECONDS, ENABLE_WEBSOCKET_PING_DEBUG, ENABLE_TASK_ASSIGNMENT_DEBUG, API_BASE_URL from websocket_server.task_assigner import assign_task_to_worker from websocket_server.events import base from websocket_server.models import Task @@ -17,9 +17,6 @@ from websocket_server.config import engine logger = logging.getLogger("websocket_server") -# API server endpoint -api_server_url = "http://172.17.0.1:5000/api" - # Tracks threads for each worker worker_threads = {} @@ -42,7 +39,7 @@ def notify_worker_online(worker_id): headers = {"Content-Type": "application/json"} try: - response = requests.put(f"{api_server_url}/workload_hosts/{worker_id}", json=payload, headers=headers) + response = requests.put(f"{API_BASE_URL}/workload_hosts/{worker_id}", json=payload, headers=headers) logger.info(f"[{worker_id}] API server acknowledged online state") # Start host reconciliation process @@ -61,7 +58,7 @@ def notify_worker_disconnect(worker_id): headers = {"Content-Type": "application/json"} try: - response = requests.put(f"{api_server_url}/workload_hosts/{worker_id}", json=payload, headers=headers) + response = requests.put(f"{API_BASE_URL}/workload_hosts/{worker_id}", json=payload, headers=headers) logger.debug(response.text) logger.info(f"[{worker_id}] API server acknowledged offline state") except Exception as e: @@ -105,7 +102,7 @@ def update_host_status(worker_id, status): headers = {"Content-Type": "application/json"} try: - response = requests.put(f"{api_server_url}/workload_hosts/{worker_id}", json=payload, headers=headers) + response = requests.put(f"{API_BASE_URL}/workload_hosts/{worker_id}", json=payload, headers=headers) if response.status_code == 200: logger.info(f"[{worker_id}] Host status updated to: {status}") else: @@ -121,7 +118,7 @@ def get_expected_containers_for_host(worker_id): logger.debug(f"[{worker_id}] Fetching expected containers for host") try: - response = requests.get(f"{api_server_url}/workload_hosts/{worker_id}/container_workloads") + response = requests.get(f"{API_BASE_URL}/workload_hosts/{worker_id}/container_workloads") if response.status_code == 200: data = response.json() if data.get('success') and 'data' in data: diff --git a/worker/.env.example b/worker/.env.example new file mode 100644 index 0000000..60b2750 --- /dev/null +++ b/worker/.env.example @@ -0,0 +1,59 @@ +# Worker-Specific Environment Configuration +# This file is used when running the worker docker-compose independently +# If running worker with the main docker-compose, use the root .env file instead + +# API Server & WebSocket endpoints + +# 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_URL=http://websocket-server:6001 + +# Unix socket path for the metadata proxy +METADATA_SOCKET_PATH=/tmp/xcloudify/metadata-server.sock + +# Unique identifier for this worker instance +# WORKER_ID= + +# Secret key for worker authentication +# WORKER_SECRET= + +# Region identifier where this worker operates +# REGION_ID= + +# Enrollment key to authorize worker joining a region +# REGION_ENROLLMENT_KEY= + +# Disable libvirt-based virtualization +NO_LIBVIRT=false + +# Disable Docker support +NO_DOCKER=false + +# Shared filesystem mount point (for NFS/shared storage) +SHAREDFS_ROOT=/mnt/shared + +# Local volume storage path +LOCAL_VOLUME_PATH=/var/lib/xcloudify/local-volumes + +# DNS configuration base path +GLOBAL_DNS_CONFIG_BASE_PATH=/var/lib/xcloudify/dns-configs + +# Log level (DEBUG, INFO, WARNING, ERROR) +LOG_LEVEL=INFO + +# Enable WebSocket ping debug output +ENABLE_WEBSOCKET_PING_DEBUG=false + +# Enable task assignment debug output +ENABLE_TASK_ASSIGNMENT_DEBUG=false + +# Seconds without a PING/PONG before worker is marked offline +PING_HEARTBEAT_TIMEOUT_SECONDS=30 + +# Default restart policy for containers created by this worker +RESTART_POLICY=always + +# Default volume path for worker-created volumes +DEFAULT_VOLUME_PATH=/tmp diff --git a/worker/docker-compose.yml b/worker/docker-compose.yml index d736c3d..6e347ab 100644 --- a/worker/docker-compose.yml +++ b/worker/docker-compose.yml @@ -6,16 +6,11 @@ services: container_name: worker #user: "1000:120" restart: unless-stopped - # environment: - # - WEBSOCKET_SERVER_URL=http://10.121.20.36:6001 - # - REGION_ID=3bcaf6fd-7638-45fe-9701-d416630b08ff - # - REGION_ENROLLMENT_KEY=283cb3a7-ed19-4152-a212-549991a63959 - # - API_BASE_URL=http://10.121.20.36:5000 - # - NO_LIBVIRT=True - # - SHAREDFS_ROOT=/mnt/nfs/cloud/ + env_file: .env volumes: - .:/app - /run/libvirt:/run/libvirt + - /var/lib/xcloudify:/var/lib/xcloudify - /var/run/docker.sock:/var/run/docker.sock - /var/run/openvswitch:/var/run/openvswitch - /run/openvswitch:/run/openvswitch diff --git a/worker/dockerMonitor.py b/worker/dockerMonitor.py index 149f8fe..4b51d49 100644 --- a/worker/dockerMonitor.py +++ b/worker/dockerMonitor.py @@ -1,15 +1,15 @@ import docker from logger import logger import time -import os from settings import settings +from runtime_urls import DOCKER_HOST class DockerMonitor: def __init__(self, event_queue): self.event_queue = event_queue # Explicitly set the Docker socket path - self.docker_socket = os.getenv('DOCKER_HOST', 'unix:///var/run/docker.sock') + self.docker_socket = DOCKER_HOST try: self.docker_client = docker.DockerClient(base_url=self.docker_socket) # Verify connection diff --git a/worker/main.py b/worker/main.py index 3e18e10..f635f40 100644 --- a/worker/main.py +++ b/worker/main.py @@ -33,7 +33,7 @@ def enroll_worker(): logger.error(f"Enrollment failed: Missing required setting: {e}") exit(1) - API_URL = f"{api_base_url}/api/workload_hosts/enroll" + API_URL = f"{api_base_url}/workload_hosts/enroll" # Construct JSON payload diff --git a/worker/run-worker.sh b/worker/run-worker.sh new file mode 100755 index 0000000..7e4117a --- /dev/null +++ b/worker/run-worker.sh @@ -0,0 +1,64 @@ +#!/usr/bin/env bash + +set -e + +# Ensure worker_settings.json exists +if [ ! -f worker_settings.json ]; then + touch worker_settings.json + echo "Created worker_settings.json" +fi + +show_help() { + echo "Usage: $0 [OPTIONS]" + echo "" + echo "Options:" + echo " --build Rebuild containers without cache and start services" + echo " --clean Remove containers, volumes, images, and rebuild everything" + echo " --help Show this help message" +} + +# Detect docker compose command +if command -v docker-compose >/dev/null 2>&1; then + DOCKER_COMPOSE="docker-compose" +elif docker compose version >/dev/null 2>&1; then + DOCKER_COMPOSE="docker compose" +else + echo "Error: Neither docker-compose nor docker compose is installed." + exit 1 +fi + +case "$1" in + --help) + show_help + ;; + + --build) + echo "Building containers without cache..." + $DOCKER_COMPOSE build --no-cache + echo "Starting services..." + $DOCKER_COMPOSE up -d + ;; + + --clean) + echo "Cleaning containers, volumes, images, and orphans..." + $DOCKER_COMPOSE down -v --rmi all --remove-orphans + + echo "Rebuilding containers without cache..." + $DOCKER_COMPOSE build --no-cache + + echo "Starting services..." + $DOCKER_COMPOSE up -d + ;; + + "") + echo "Starting services..." + $DOCKER_COMPOSE up -d + ;; + + *) + echo "Unknown option: $1" + echo "" + show_help + exit 1 + ;; +esac \ No newline at end of file diff --git a/worker/runtime_urls.py b/worker/runtime_urls.py new file mode 100644 index 0000000..bc59c2c --- /dev/null +++ b/worker/runtime_urls.py @@ -0,0 +1,48 @@ +"""Centralized runtime URL and configuration values for the worker. + +This module loads the worker .env file once and exposes resolved values so the +worker can import config instead of reading environment variables directly at +each call site. +""" + +import os + +from dotenv import load_dotenv + +load_dotenv() + + +def _env(name: str, default: str) -> str: + return os.getenv(name, default) + + +def _env_bool(name: str, default: bool = False) -> bool: + value = os.getenv(name) + return default if value is None else value.lower() in ("true", "1", "t", "yes", "y") + + +def _env_int(name: str, default: int) -> int: + value = os.getenv(name) + return default if value is None else int(value) + + +def load_worker_env(dotenv_path: str = ".env") -> None: + load_dotenv(dotenv_path=dotenv_path, override=False) + + +def get_worker_env(name: str, default: str | None = None) -> str | None: + value = os.getenv(name) + return default if value is None else value + + +API_BASE_URL = _env("API_BASE_URL", "http://localhost:5000/api") +WEBSOCKET_SERVER_URL = _env("WEBSOCKET_SERVER_URL", "http://localhost:6001") +DOCKER_HOST = _env("DOCKER_HOST", "unix:///var/run/docker.sock") +METADATA_SOCKET_PATH = _env("METADATA_SOCKET_PATH", "/tmp/xcloudify/metadata-server.sock") +DEFAULT_VOLUME_PATH = "/tmp" +GLOBAL_DNS_CONFIG_BASE_PATH = "/var/lib/xcloudify/dns-configs" +LOCAL_VOLUME_PATH = "/var/lib/xcloudify/local-volumes" +SHAREDFS_ROOT = "/mnt/shared" +RESTART_POLICY = "always" +LOG_LEVEL = "INFO" +ENABLE_WEBSOCKET_PING_DEBUG = _env_bool("ENABLE_WEBSOCKET_PING_DEBUG", False) \ No newline at end of file diff --git a/worker/settings.py b/worker/settings.py index a730c38..fa4bcbe 100644 --- a/worker/settings.py +++ b/worker/settings.py @@ -1,10 +1,11 @@ import os import json import logging -from dotenv import load_dotenv from threading import Lock from logger import logger +from runtime_urls import API_BASE_URL, DEFAULT_VOLUME_PATH, ENABLE_WEBSOCKET_PING_DEBUG, GLOBAL_DNS_CONFIG_BASE_PATH, LOCAL_VOLUME_PATH, LOG_LEVEL, RESTART_POLICY, SHAREDFS_ROOT, WEBSOCKET_SERVER_URL, get_worker_env, load_worker_env + class SettingsManager: """Manages application settings, loading from .env and persisting to a JSON file.""" @@ -30,16 +31,18 @@ class SettingsManager: # Default values for optional settings DEFAULT_SETTINGS = { - "DEFAULT_VOLUME_PATH": "/tmp", + "API_BASE_URL": API_BASE_URL, + "DEFAULT_VOLUME_PATH": DEFAULT_VOLUME_PATH, "DEBUG_SOCKETIO": False, - "ENABLE_WEBSOCKET_PING_DEBUG": False, - "GLOBAL_DNS_CONFIG_BASE_PATH": "/var/lib/xcloudify/dns-configs", - "LOG_LEVEL": "INFO", - "LOCAL_VOLUME_PATH": "/var/lib/xcloudify/local-volumes", + "ENABLE_WEBSOCKET_PING_DEBUG": ENABLE_WEBSOCKET_PING_DEBUG, + "GLOBAL_DNS_CONFIG_BASE_PATH": GLOBAL_DNS_CONFIG_BASE_PATH, + "LOG_LEVEL": LOG_LEVEL, + "LOCAL_VOLUME_PATH": LOCAL_VOLUME_PATH, "NO_DOCKER": False, "NO_LIBVIRT": False, - "SHAREDFS_ROOT": "/mnt/shared", - "RESTART_POLICY": "always", # Default restart policy for containers + "WEBSOCKET_SERVER_URL": WEBSOCKET_SERVER_URL, + "SHAREDFS_ROOT": SHAREDFS_ROOT, + "RESTART_POLICY": RESTART_POLICY, # Default restart policy for containers # Add other settings with defaults here } @@ -95,8 +98,8 @@ class SettingsManager: """Loads settings from the .env file and environment variables.""" settings_from_env = {} try: - # Load the .env file first - load_dotenv(dotenv_path=self._env_file, override=False) # override=False: env vars take precedence + # Load the configured .env file for non-URL worker settings. + load_worker_env(self._env_file) logger.info(f"Checked for environment variables from {self._env_file}") # Check all known keys (required, optional defaults, cli overridable) @@ -104,7 +107,7 @@ class SettingsManager: loaded_count = 0 for key in all_known_keys: - value = os.getenv(key) + value = get_worker_env(key) if value is not None: # Basic type conversion for known defaults if key in self.DEFAULT_SETTINGS: diff --git a/worker/worker_tasks/libvirt.py b/worker/worker_tasks/libvirt.py index dd953f1..0fb8b70 100644 --- a/worker/worker_tasks/libvirt.py +++ b/worker/worker_tasks/libvirt.py @@ -279,7 +279,7 @@ class LibvirtVirtualMachineTask: return try: resp = requests.post( - f"{api_base_url}/api/gpus/release", + f"{api_base_url}/gpus/release", json={"workload_id": workload_id}, headers={"Authorization": f"Bearer {api_token}", "Content-Type": "application/json"}, timeout=10, diff --git a/worker/worker_tasks/metadata_worker_proxy.py b/worker/worker_tasks/metadata_worker_proxy.py index 4435405..62330bd 100644 --- a/worker/worker_tasks/metadata_worker_proxy.py +++ b/worker/worker_tasks/metadata_worker_proxy.py @@ -12,16 +12,14 @@ sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__)))) from logger import logger from settings import settings - -# Metadata socket path -METADATA_SOCKET_PATH = "/tmp/xcloudify/metadata-server.sock" +from runtime_urls import API_BASE_URL, METADATA_SOCKET_PATH app = Flask(__name__) def _forward_request(path: str, headers=None): """Forward request to the API server.""" - api_url = settings.get_value("API_BASE_URL", "http://172.17.0.1:5000") + api_url = settings.get_value("API_BASE_URL", API_BASE_URL) target_url = f"{api_url}{path}" forward_headers = dict(request.headers) @@ -78,7 +76,7 @@ def root(): def healthz(): """Health check endpoint.""" try: - api_url = settings.get_value("API_BASE_URL", "http://172.17.0.1:5000") + api_url = settings.get_value("API_BASE_URL", API_BASE_URL) response = requests.get(f"{api_url}/api/metadata/healthz", timeout=5) return jsonify({ "status": "ok" if response.status_code == 200 else "degraded", @@ -117,7 +115,7 @@ def run_proxy(): except Exception as e: logger.warning(f"Could not remove stale socket: {e}") - api_url = settings.get_value("API_BASE_URL", "http://172.17.0.1:5000") + api_url = settings.get_value("API_BASE_URL", API_BASE_URL) logger.info(f"Starting metadata proxy on {METADATA_SOCKET_PATH}") logger.info(f"Forwarding to API server at {api_url}") diff --git a/worker/worker_tasks/ovs_bridge_scanner.py b/worker/worker_tasks/ovs_bridge_scanner.py index b08098d..0d80e70 100644 --- a/worker/worker_tasks/ovs_bridge_scanner.py +++ b/worker/worker_tasks/ovs_bridge_scanner.py @@ -18,7 +18,7 @@ class OVSBridgeScannerTask: raise TypeError("Logger must be an instance of the logging.Logger class.") self.logger = logger - self.host_id = settings.get_value("WORKER_ID") or os.getenv("WORKER_ID") + self.host_id = settings.get_value("WORKER_ID") self.api_base_url = settings.get_value("API_BASE_URL") self.api_auth_token = settings.get_value("WORKER_SECRET") # Use worker secret as token @@ -86,7 +86,7 @@ class OVSBridgeScannerTask: self.logger.info("No bridges to report; skipping API call.") return True # No error if empty - url = f"{self.api_base_url}/api/workload_hosts/{self.host_id}/ovs_bridges" + url = f"{self.api_base_url}/workload_hosts/{self.host_id}/ovs_bridges" payload = {"bridges": bridges} headers = { "Content-Type": "application/json", diff --git a/worker/worker_tasks/pci_device_scanner.py b/worker/worker_tasks/pci_device_scanner.py index f0c6359..23a477d 100644 --- a/worker/worker_tasks/pci_device_scanner.py +++ b/worker/worker_tasks/pci_device_scanner.py @@ -70,7 +70,7 @@ class PCIDeviceScannerTask: """ try: # First, get the host info to find the region_id - url = f"{self.api_base_url}/api/workload_hosts/{self.host_id}" + url = f"{self.api_base_url}/workload_hosts/{self.host_id}" headers = {"Authorization": f"Bearer {self.api_auth_token}"} response = requests.get(url, headers=headers, timeout=10) @@ -86,7 +86,7 @@ class PCIDeviceScannerTask: return False # Now get the region config - url = f"{self.api_base_url}/api/regions/{self.region_id}" + url = f"{self.api_base_url}/regions/{self.region_id}" response = requests.get(url, headers=headers, timeout=10) if response.status_code != 200: