Feat: Implemented X-Internal-Xcloudify-API
This commit is contained in:
@@ -60,6 +60,9 @@ CHECK_PERMISSION_URL=http://172.17.0.1:5000/api/check_permission
|
||||
|
||||
JWT_SECRET_KEY=change-me-to-a-strong-random-secret
|
||||
|
||||
# Shared key for all API calls, sent as the X-Internal-Xcloudify-API header.
|
||||
XCLOUDIFY_API_KEY=change-me-generate-a-real-key
|
||||
|
||||
# Base domain for user-facing resources (certificates, etc.)
|
||||
BASE_DOMAIN=xcloudify.tech
|
||||
# Tunnel domain used for container/workload hostnames via Cloudflare tunnels
|
||||
|
||||
+22
-1
@@ -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 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
|
||||
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, XCLOUDIFY_API_KEY
|
||||
|
||||
# --------------------------------------------------------------------------- #
|
||||
# 1. Flask application & database #
|
||||
@@ -111,6 +111,9 @@ 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
|
||||
app.config["XCLOUDIFY_API_KEY"] = XCLOUDIFY_API_KEY
|
||||
if not [k for k in XCLOUDIFY_API_KEY.split(",") if k.strip()]:
|
||||
logger.error("XCLOUDIFY_API_KEY is empty — every authenticated endpoint will return 401")
|
||||
|
||||
logger.info(
|
||||
"Cloudflare configuration set "
|
||||
@@ -198,6 +201,24 @@ def assign_request_id():
|
||||
g.request_id = str(uuid.uuid4())
|
||||
|
||||
|
||||
_PUBLIC_PATHS = {"/api/healthz"}
|
||||
|
||||
|
||||
@app.before_request
|
||||
def enforce_authentication():
|
||||
from app.utils.auth_utils import authenticate_request
|
||||
authenticate_request()
|
||||
|
||||
path = request.path
|
||||
if not path.startswith("/api"):
|
||||
return
|
||||
if path in _PUBLIC_PATHS:
|
||||
return
|
||||
if getattr(g, "is_service", False):
|
||||
return
|
||||
return api_response(success=False, status=401, message="Authentication required", error_type="UNAUTHORIZED")
|
||||
|
||||
|
||||
@app.after_request
|
||||
def log_request_id(response):
|
||||
response.headers["X-Request-ID"] = g.request_id
|
||||
|
||||
+22
-2
@@ -1,7 +1,10 @@
|
||||
import hmac
|
||||
from typing import Optional
|
||||
from flask import Request, current_app
|
||||
from flask import Request, current_app, g, request as flask_request
|
||||
import jwt
|
||||
|
||||
API_KEY_HEADER = "X-Internal-Xcloudify-API"
|
||||
|
||||
|
||||
def get_request_user_id(request: Request) -> Optional[str]:
|
||||
"""
|
||||
@@ -25,4 +28,21 @@ def get_request_user_id(request: Request) -> Optional[str]:
|
||||
return str(sub) if sub is not None else None
|
||||
except Exception:
|
||||
# On any JWT decode/validation error, do not block the request; just omit user_id
|
||||
return None
|
||||
return None
|
||||
|
||||
|
||||
def _accepted_keys() -> list:
|
||||
raw = current_app.config.get("XCLOUDIFY_API_KEY") or ""
|
||||
return [k.strip() for k in raw.split(",") if k.strip()]
|
||||
|
||||
|
||||
def authenticate_request():
|
||||
g.is_service = False
|
||||
presented = flask_request.headers.get(API_KEY_HEADER)
|
||||
if not presented:
|
||||
return
|
||||
presented_bytes = presented.encode("utf-8")
|
||||
for key in _accepted_keys():
|
||||
if hmac.compare_digest(presented_bytes, key.encode("utf-8")):
|
||||
g.is_service = True
|
||||
return
|
||||
|
||||
@@ -53,7 +53,7 @@ class ApiClient:
|
||||
def __init__(self, base_url: str, api_key: str) -> None:
|
||||
self.base_url = base_url.rstrip("/")
|
||||
self.session = requests.Session()
|
||||
self.session.headers.update({"X-API-KEY": api_key, "Content-Type": "application/json"})
|
||||
self.session.headers.update({"X-Internal-Xcloudify-API": api_key, "Content-Type": "application/json"})
|
||||
|
||||
# Generic helpers -------------------------------------------------------- #
|
||||
def _url(self, path: str) -> str:
|
||||
|
||||
@@ -0,0 +1,39 @@
|
||||
"""Stamp the internal service key on every outbound API request.
|
||||
|
||||
Installing this once per process adds `X-Internal-Xcloudify-API` to any
|
||||
request whose URL targets the API base, so the backend's global auth gate treats
|
||||
them as a trusted service. Scoped to the API base URL so the key never leaks to
|
||||
other hosts
|
||||
"""
|
||||
from urllib.parse import urlparse
|
||||
|
||||
import requests
|
||||
|
||||
API_KEY_HEADER = "X-Internal-Xcloudify-API"
|
||||
|
||||
|
||||
def _origin(url: str) -> str:
|
||||
"""Return scheme://host:port for a URL, or "" when it isn't absolute."""
|
||||
parsed = urlparse(url or "")
|
||||
if not parsed.scheme or not parsed.netloc:
|
||||
return ""
|
||||
return f"{parsed.scheme.lower()}://{parsed.netloc.lower()}"
|
||||
|
||||
|
||||
def install_internal_auth(api_base_url: str, api_key: str) -> None:
|
||||
if not api_key:
|
||||
return
|
||||
base = _origin(api_base_url)
|
||||
original = requests.sessions.Session.request
|
||||
if getattr(original, "_internal_auth_installed", False):
|
||||
return
|
||||
|
||||
def request(self, method, url, **kwargs):
|
||||
if base and isinstance(url, str) and _origin(url) == base:
|
||||
headers = dict(kwargs.get("headers") or {})
|
||||
headers.setdefault(API_KEY_HEADER, api_key)
|
||||
kwargs["headers"] = headers
|
||||
return original(self, method, url, **kwargs)
|
||||
|
||||
request._internal_auth_installed = True
|
||||
requests.sessions.Session.request = request
|
||||
@@ -1,4 +1,5 @@
|
||||
flask
|
||||
flask_cors
|
||||
flasgger
|
||||
colorlog
|
||||
python-socketio
|
||||
|
||||
@@ -57,6 +57,7 @@ CHECK_PERMISSION_URL = _env("CHECK_PERMISSION_URL", "http://172.17.0.1:5000/api/
|
||||
APP_ENV = _env("APP_ENV", "")
|
||||
FLASK_ENV = _env("FLASK_ENV", "development")
|
||||
JWT_SECRET_KEY = _env("JWT_SECRET_KEY", "your-very-secret-key")
|
||||
XCLOUDIFY_API_KEY = _env("XCLOUDIFY_API_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)
|
||||
|
||||
@@ -1,10 +1,13 @@
|
||||
import streamlit as st
|
||||
# from streamlit_server.api.client import IaaSClient
|
||||
from api_client.client import IaaSClient
|
||||
from internal_http import install_internal_auth
|
||||
from runtime_urls import API_BASE_URL
|
||||
|
||||
def init_session_state():
|
||||
# Initialize all required session state keys before use
|
||||
if 'client' not in st.session_state:
|
||||
install_internal_auth(API_BASE_URL, st.session_state.get('api_key', ''))
|
||||
st.session_state.client = IaaSClient(st.session_state.get('api_key', ''))
|
||||
if 'api_key' not in st.session_state:
|
||||
st.session_state.api_key = ""
|
||||
|
||||
+4
-1
@@ -8,7 +8,10 @@ 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, VNC_PROXY_WS_PATH, 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, API_BASE_URL, XCLOUDIFY_API_KEY
|
||||
from internal_http import install_internal_auth
|
||||
|
||||
install_internal_auth(API_BASE_URL, XCLOUDIFY_API_KEY)
|
||||
|
||||
logging.basicConfig(level=logging.DEBUG, format="%(asctime)s [%(levelname)s] %(message)s")
|
||||
logger = logging.getLogger("middleware")
|
||||
|
||||
@@ -8,7 +8,10 @@ import logging
|
||||
from websocket_server.config import setup_logging
|
||||
from websocket_server.worker_manager import worker_dispatch_flag_check, all_worker_watchdog
|
||||
from websocket_server.events import base
|
||||
from runtime_urls import WEBSOCKET_SERVER_PATH
|
||||
from runtime_urls import WEBSOCKET_SERVER_PATH, API_BASE_URL, XCLOUDIFY_API_KEY
|
||||
from internal_http import install_internal_auth
|
||||
|
||||
install_internal_auth(API_BASE_URL, XCLOUDIFY_API_KEY)
|
||||
# Setup logging
|
||||
logger = setup_logging()
|
||||
logger.info("Starting WebSocket server...")
|
||||
|
||||
@@ -7,6 +7,8 @@
|
||||
# API Server base URL - where worker sends registration/status updates
|
||||
API_BASE_URL=http://api-server:5000/api
|
||||
|
||||
XCLOUDIFY_API_KEY=change-me-generate-a-real-key
|
||||
|
||||
# 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
|
||||
|
||||
@@ -0,0 +1,37 @@
|
||||
"""Stamp the internal service key on every outbound API request .
|
||||
|
||||
The worker runs from its own directory, so it carries a local copy of the
|
||||
injector.
|
||||
"""
|
||||
from urllib.parse import urlparse
|
||||
|
||||
import requests
|
||||
|
||||
API_KEY_HEADER = "X-Internal-Xcloudify-API"
|
||||
|
||||
|
||||
def _origin(url: str) -> str:
|
||||
"""Return scheme://host:port for a URL, or "" when it isn't absolute."""
|
||||
parsed = urlparse(url or "")
|
||||
if not parsed.scheme or not parsed.netloc:
|
||||
return ""
|
||||
return f"{parsed.scheme.lower()}://{parsed.netloc.lower()}"
|
||||
|
||||
|
||||
def install_internal_auth(api_base_url: str, api_key: str) -> None:
|
||||
if not api_key:
|
||||
return
|
||||
base = _origin(api_base_url)
|
||||
original = requests.sessions.Session.request
|
||||
if getattr(original, "_internal_auth_installed", False):
|
||||
return
|
||||
|
||||
def request(self, method, url, **kwargs):
|
||||
if base and isinstance(url, str) and _origin(url) == base:
|
||||
headers = dict(kwargs.get("headers") or {})
|
||||
headers.setdefault(API_KEY_HEADER, api_key)
|
||||
kwargs["headers"] = headers
|
||||
return original(self, method, url, **kwargs)
|
||||
|
||||
request._internal_auth_installed = True
|
||||
requests.sessions.Session.request = request
|
||||
@@ -13,6 +13,8 @@ from datetime import datetime
|
||||
# from dotenv import load_dotenv
|
||||
from logger import logger
|
||||
from settings import settings # Import the global settings instance
|
||||
from runtime_urls import API_BASE_URL, XCLOUDIFY_API_KEY
|
||||
from internal_http import install_internal_auth
|
||||
from dockerMonitor import DockerMonitor
|
||||
from workerClient import WorkerClient
|
||||
from worker_tasks.metadata_worker_proxy import run_proxy
|
||||
@@ -312,6 +314,14 @@ def main():
|
||||
# Update settings from CLI arguments
|
||||
settings.update_from_args(args)
|
||||
|
||||
# Stamp the internal service key on all outbound API calls (auth for the API gate)
|
||||
if not XCLOUDIFY_API_KEY:
|
||||
logger.error(
|
||||
"XCLOUDIFY_API_KEY is not set — the API will reject every request from this "
|
||||
"worker with 401. Set it in the worker .env to match the API server's key."
|
||||
)
|
||||
install_internal_auth(settings.get_value("API_BASE_URL", API_BASE_URL), XCLOUDIFY_API_KEY)
|
||||
|
||||
logger.info("Starting Worker CLI with monitoring capabilities")
|
||||
logger.info(f"Settings loaded. Current settings: {settings.get_all_settings()}")
|
||||
|
||||
|
||||
@@ -36,6 +36,7 @@ def get_worker_env(name: str, default: str | None = None) -> str | None:
|
||||
|
||||
|
||||
API_BASE_URL = _env("API_BASE_URL", "http://localhost:5000/api")
|
||||
XCLOUDIFY_API_KEY = _env("XCLOUDIFY_API_KEY", "")
|
||||
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")
|
||||
|
||||
@@ -14,7 +14,14 @@ sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__))))
|
||||
|
||||
from logger import logger
|
||||
from settings import settings
|
||||
from runtime_urls import API_BASE_URL, METADATA_SOCKET_PATH
|
||||
from runtime_urls import API_BASE_URL, METADATA_SOCKET_PATH, XCLOUDIFY_API_KEY
|
||||
from internal_http import API_KEY_HEADER
|
||||
|
||||
if not XCLOUDIFY_API_KEY:
|
||||
logger.error(
|
||||
"XCLOUDIFY_API_KEY is not set — the API will reject every cloud-init metadata "
|
||||
"request with 401 and VMs will boot without SSH keys or user-data"
|
||||
)
|
||||
|
||||
app = Flask(__name__)
|
||||
_STRIP_RESPONSE_HEADERS = {
|
||||
@@ -56,6 +63,12 @@ def _forward_request(path: str, headers=None):
|
||||
|
||||
forward_headers = dict(request.headers)
|
||||
forward_headers.pop('Host', None)
|
||||
|
||||
for name in [k for k in forward_headers if k.lower() == API_KEY_HEADER.lower()]:
|
||||
forward_headers.pop(name)
|
||||
if XCLOUDIFY_API_KEY:
|
||||
forward_headers[API_KEY_HEADER] = XCLOUDIFY_API_KEY
|
||||
|
||||
if headers:
|
||||
forward_headers.update(headers)
|
||||
|
||||
@@ -99,6 +112,12 @@ def _forward_request(path: str, headers=None):
|
||||
logger.error(f"Error forwarding request to {target_url}: {e}")
|
||||
return jsonify({"error": "Internal server error"}), 500
|
||||
|
||||
if response.status_code == 401:
|
||||
logger.error(
|
||||
f"API rejected metadata request path={path} with 401 — the XCLOUDIFY_API_KEY on this "
|
||||
f"worker is missing or does not match the one configured on the API server ({target_url})"
|
||||
)
|
||||
|
||||
return Response(body, status=response.status_code, headers=_sanitize_response_headers(response.headers))
|
||||
|
||||
|
||||
@@ -118,11 +137,35 @@ def get_userdata():
|
||||
return _forward_request("/v1/user-data")
|
||||
|
||||
|
||||
# Only these paths are proxied to the API. This proxy attaches the internal service
|
||||
# key, so a catch-all would let any VM that can reach 169.254.169.254 use it as an
|
||||
# authenticated tunnel into the control plane. Unknown paths are answered here with a
|
||||
# 404 — the same status the API returns for them — so cloud-init behaviour is unchanged.
|
||||
_ALLOWED_PATHS = frozenset({
|
||||
"/v1.json",
|
||||
"/v1/user-data",
|
||||
"/openstack",
|
||||
"/openstack/latest",
|
||||
"/openstack/latest/meta_data.json",
|
||||
"/openstack/latest/user_data",
|
||||
"/openstack/latest/network_data.json",
|
||||
"/openstack/latest/vendor_data.json",
|
||||
"/openstack/latest/vendor_data2.json",
|
||||
})
|
||||
|
||||
|
||||
@app.get("/<path:path>")
|
||||
def proxy_any(path):
|
||||
"""Forward any remaining metadata path to the API server."""
|
||||
logger.info(f"Metadata request path={path} from {request.remote_addr or 'unix-socket'}")
|
||||
return _forward_request(f"/{path}")
|
||||
"""Forward a known metadata path to the API server; refuse anything else."""
|
||||
normalized = f"/{path}".rstrip("/") or "/"
|
||||
if normalized not in _ALLOWED_PATHS:
|
||||
logger.warning(
|
||||
f"Blocked non-metadata path={normalized} from {request.remote_addr or 'unix-socket'}"
|
||||
)
|
||||
return jsonify({"error": "Not Found"}), 404
|
||||
|
||||
logger.info(f"Metadata request path={normalized} from {request.remote_addr or 'unix-socket'}")
|
||||
return _forward_request(normalized)
|
||||
|
||||
|
||||
@app.route("/")
|
||||
|
||||
Reference in New Issue
Block a user