diff --git a/app/__init__.py b/app/__init__.py index 62a30f0..e8ef84c 100644 --- a/app/__init__.py +++ b/app/__init__.py @@ -2,7 +2,7 @@ Entry point for the xCloudify Flask application. * Loads environment variables from a `.env` file. -* Configures APIFlask, SQLAlchemy, Celery, and Cloudflare settings. +* Configures Flask, SQLAlchemy, Celery, and Cloudflare settings. * Masks sensitive values in logs for security. """ @@ -12,8 +12,7 @@ import uuid import os import pymysql from dotenv import load_dotenv -from apiflask import APIFlask -from flask import g, request +from flask import Flask, g, jsonify, request from flask_sqlalchemy import SQLAlchemy from flask_migrate import Migrate from celery import Celery, Task @@ -28,43 +27,49 @@ load_dotenv() logger.info("Environment variables loaded from .env") # --------------------------------------------------------------------------- # -# 1. APIFlask application & database # +# 1. Flask application & database # # --------------------------------------------------------------------------- # -app: APIFlask = APIFlask( - __name__, - title="xCloudify IaaS API", - version="0.1.0", - docs_path="/docs", # Swagger-UI - spec_path="/openapi.json", # raw OpenAPI 3.1 -) +app: Flask = Flask(__name__) CORS(app) 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_BROKER_URL=os.getenv("CELERY_BROKER_URL", "redis://172.17.0.1:6379/0"), - CELERY_RESULT_BACKEND=os.getenv("CELERY_RESULT_BACKEND", "redis://172.17.0.1:6379/1"), + + # Celery settings (override via ENV) + CELERY_BROKER_URL=os.getenv("CELERY_BROKER_URL", "redis://172.17.0.1:6379/0"), # Redis Database 0 for celery broker tasks + CELERY_RESULT_BACKEND=os.getenv("CELERY_RESULT_BACKEND", "redis://172.17.0.1:6379/1"), # Redis Database 1 for celery results ) # --------------------------------------------------------------------------- # # 2. Cloudflare configuration # # --------------------------------------------------------------------------- # def _mask(value: str, visible: int = 4) -> str: + """ + Return a partially masked version of a sensitive string. + + Args: + value (str): The string to mask. + visible (int): Number of visible characters to keep at the start. + + Returns: + str: Masked string suitable for logging. + """ if not value: 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)) -app.config["REDIS_URL"] = os.getenv("REDIS_URL", "redis://172.17.0.1:6379/2") -app.config["VNC_SECRET_KEY"] = os.getenv("VNC_SECRET_KEY", "your-very-secret-key") +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"] = 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") logger.info( "Cloudflare configuration set " f"(token={_mask(app.config['CLOUDFLARE_API_TOKEN'])}, " - f"account_id={_mask(app.config['CLOUDFLARE_ACCOUNT_ID'])}, " + f"account_id={_mask(app.config['CLOUDFLARE_ACCOUNT_ID'])})" f"zone_id={_mask(app.config['CLOUDFLARE_ZONE_ID'])})" ) @@ -75,9 +80,15 @@ migrate: Migrate = Migrate(app, db) app.secret_key = "your_secret_key_here" # --------------------------------------------------------------------------- # -# 3. Celery initialisation # +# 2. Celery initialisation # # --------------------------------------------------------------------------- # -def _make_celery(flask_app: APIFlask) -> Celery: +def _make_celery(flask_app: Flask) -> Celery: + """ + Create a Celery instance that shares the Flask application context. + + Every task will inherit `current_app` and can use the database session + without manual `app.app_context()` juggling. + """ celery = Celery( flask_app.import_name, broker=flask_app.config["CELERY_BROKER_URL"], @@ -86,7 +97,9 @@ def _make_celery(flask_app: APIFlask) -> Celery: celery.conf.update(flask_app.config) class ContextTask(Task): + """Celery Task base-class that wraps task execution in app_context.""" abstract = True + def __call__(self, *args, **kwargs): with flask_app.app_context(): return super().__call__(*args, **kwargs) @@ -94,21 +107,21 @@ def _make_celery(flask_app: APIFlask) -> Celery: celery.Task = ContextTask return celery + celery_app: Celery = _make_celery(app) celery_app.autodiscover_tasks(["app.tasks"], force=True) app.extensions["celery"] = celery_app - # --------------------------------------------------------------------------- # -# 4. Blueprints and middleware # +# 3. Blueprints and standard middleware # # --------------------------------------------------------------------------- # -from app.controller import api_bp, ui_bp, celery_bp # noqa: E402 +from app.controller import api_bp, ui_bp, celery_bp # noqa: E402 circular OK app.register_blueprint(api_bp) app.register_blueprint(ui_bp) app.register_blueprint(celery_bp) -# ---------- Universal error envelope ------------------------------------- # -@app.error_processor -def _wrap_all_errors(status_code, message, exc): +# ---------- Error-handlers, request IDs, context processors --------------- # +@app.errorhandler(500) +def handle_internal_server_error(error): logger.error( "Unhandled exception for %s %s (request_id=%s)", request.method, @@ -118,21 +131,31 @@ def _wrap_all_errors(status_code, message, exc): ) return api_response( success=False, - status=status_code, - message=message, - error_type=getattr(exc, "__class__", type("E", (), {})).__name__, + status=500, + message="Internal server error", + error_type="INTERNAL_SERVER_ERROR", + error_details={"path": request.path}, ) -# ---------- Request IDs & context processors ------------------------------ # + +@app.errorhandler(404) +def not_found(e): + return api_response( + success=False, status=404, message=str(e), error_type="RESOURCE_NOT_FOUND" + ) + + @app.before_request def assign_request_id(): g.request_id = str(uuid.uuid4()) + @app.after_request def log_request_id(response): response.headers["X-Request-ID"] = g.request_id return response + @app.context_processor def inject_user_roles(): try: @@ -150,4 +173,5 @@ def inject_user_roles(): logger.error("inject_user_roles - %s", exc) return {} -logger.debug("__init__ complete") + +logger.debug("__init__ complete") \ No newline at end of file diff --git a/app/controller/__init__.py b/app/controller/__init__.py index 0a8a432..8e6b061 100644 --- a/app/controller/__init__.py +++ b/app/controller/__init__.py @@ -1,9 +1,9 @@ -from apiflask import APIBlueprint +from flask import Blueprint # bp = Blueprint('main', __name__) -api_bp = APIBlueprint('api', __name__, url_prefix='/api') -ui_bp = APIBlueprint('ui', __name__, url_prefix='/ui') -celery_bp = APIBlueprint("celery", __name__,url_prefix='/celery') +api_bp = Blueprint('api', __name__, url_prefix='/api') +ui_bp = Blueprint('ui', __name__, url_prefix='/ui') +celery_bp = Blueprint("celery", __name__,url_prefix='/celery') """ Blueprint for the main application routes. diff --git a/app/controller/api/image_routes.py b/app/controller/api/image_routes.py index cb92c49..6ff8d56 100644 --- a/app/controller/api/image_routes.py +++ b/app/controller/api/image_routes.py @@ -1,104 +1,249 @@ -from flask import request, jsonify, abort +""" +API endpoints for managing Images within xCloudify. + +All responses comply with the standard `api_response` envelope to ensure a +consistent contract for consumers. +""" + +from typing import List, Dict, Any + +from flask import request from app import db, logger from app.models.models import Image from app.controller import api_bp +from app.utils.standard_responses import api_response -def validate_image_data(data, is_update=False): + +def validate_image_data(request_payload: Dict[str, Any], is_update: bool = False) -> List[str]: """ - Validates the incoming image data. + Validate the incoming JSON payload for Image creation or update. - Args: - data (dict): The incoming JSON payload. - is_update (bool): Whether this is an update operation (some fields may be optional). + Parameters + ---------- + request_payload : dict + The incoming request data. + is_update : bool, default=False + When True, required-field checks are relaxed for partial updates. - Returns: - dict: The validated data. - - Raises: - HTTPException: If validation fails, aborts with a 400 status code and error message. + Returns + ------- + list[str] + A list of validation error strings (empty if no errors detected). """ - errors = [] + validation_errors: List[str] = [] - # Required fields for creation + # ── Required fields on creation ──────────────────────────────────────────── if not is_update: - if 'location' not in data: - errors.append("Missing required field: location") - if 'size' not in data: - errors.append("Missing required field: size") + if "location" not in request_payload: + validation_errors.append("Missing required field: location") + if "size" not in request_payload: + validation_errors.append("Missing required field: size") - # Validate location (if present) - if 'location' in data and not isinstance(data['location'], str): - errors.append("Field 'location' must be a string") + # ── Type checks ─────────────────────────────────────────────────────────── + if "location" in request_payload and not isinstance(request_payload["location"], str): + validation_errors.append("Field 'location' must be a string") - # Validate size (if present) - if 'size' in data and not isinstance(data['size'], (int, float)): - errors.append("Field 'size' must be a number") + if "size" in request_payload and not isinstance(request_payload["size"], (int, float)): + validation_errors.append("Field 'size' must be a number") - # Validate optional fields - optional_fields = { - 'location_type': str, - 'os_family': str, - 'os_version': str, - 'checksum': str, + optional_field_types = { + "location_type": str, + "os_family": str, + "os_version": str, + "checksum": str, } - for field, field_type in optional_fields.items(): - if field in data and not isinstance(data[field], field_type): - errors.append(f"Field '{field}' must be of type {field_type.__name__}") + for field_name, expected_type in optional_field_types.items(): + if field_name in request_payload and not isinstance(request_payload[field_name], expected_type): + validation_errors.append(f"Field '{field_name}' must be of type {expected_type.__name__}") - if errors: - abort(400, description=", ".join(errors)) + return validation_errors - return data -@api_bp.route('/images', methods=['POST']) +@api_bp.route("/images", methods=["POST"]) def add_image(): - data = request.json - logger.debug(f"Image request {data}") - # Validate incoming data - validated_data = validate_image_data(data) + """ + Create a new Image. - # Create the Image - _image = Image( - name="a", - location=validated_data['location'], - size=validated_data['size'], - location_type=validated_data.get('location_type'), - os_family=validated_data.get('os_family'), - os_version=validated_data.get('os_version'), - checksum=validated_data.get('checksum') - ) - db.session.add(_image) - db.session.commit() - logger.debug(f"Image added to DB with payload {_image.to_json()}") - return jsonify(_image.to_json()), 201 + Returns + ------- + Tuple(Response, int) + Standardised response containing the created Image data. + """ + try: + request_data = request.get_json(force=True) or {} + logger.debug("Received Image creation request: %s", request_data) -@api_bp.route('/images/', methods=['PUT']) + # ── Validate input ──────────────────────────────────────────────────── + validation_errors = validate_image_data(request_data) + if validation_errors: + return api_response( + success=False, + message="Validation error", + status=400, + error_type="VALIDATION_ERROR", + error_details={"errors": validation_errors}, + ) + + # ── Create the Image ───────────────────────────────────────────────── + image_instance = Image( + name=request_data.get("name"), # Name is optional in schema + location=request_data["location"], + size=request_data["size"], + location_type=request_data.get("location_type"), + os_family=request_data.get("os_family"), + os_version=request_data.get("os_version"), + checksum=request_data.get("checksum"), + ) + db.session.add(image_instance) + db.session.commit() + + logger.info("Image '%s' (id=%s) created", image_instance.name, image_instance.id) + + return api_response( + data=image_instance.to_json(), + message="Image created", + status=201, + ) + + except Exception as exc: # pylint: disable=broad-except + logger.exception("Failed to create Image") + return api_response( + success=False, + message="Internal Server Error", + status=500, + error_type=type(exc).__name__, + error_details={"detail": str(exc)}, + ) + + +@api_bp.route("/images/", methods=["PUT"]) def edit_image(image_id): - image = Image.query.get_or_404((image_id)) - data = request.json + """ + Update an existing Image. - # Validate incoming data - validated_data = validate_image_data(data, is_update=True) + Parameters + ---------- + image_id : str + Identifier of the Image to update. - # Update the fields - for key, value in validated_data.items(): - setattr(image, key, value) - db.session.commit() - return jsonify(success=True) + Returns + ------- + Tuple(Response, int) + Standardised response confirming the update. + """ + try: + image_instance = Image.query.get_or_404(image_id) + request_data = request.get_json(force=True) or {} + logger.debug("Received Image update request id=%s payload=%s", image_id, request_data) -@api_bp.route('/images/', methods=['GET']) + # ── Validate input ──────────────────────────────────────────────────── + validation_errors = validate_image_data(request_data, is_update=True) + if validation_errors: + return api_response( + success=False, + message="Validation error", + status=400, + error_type="VALIDATION_ERROR", + error_details={"errors": validation_errors}, + ) + + # ── Apply updates ──────────────────────────────────────────────────── + for field_name, field_value in request_data.items(): + setattr(image_instance, field_name, field_value) + + db.session.commit() + logger.info("Image id=%s updated", image_id) + + return api_response(message="Image updated") + + except Exception as exc: # pylint: disable=broad-except + logger.exception("Failed to update Image %s", image_id) + return api_response( + success=False, + message="Internal Server Error", + status=500, + error_type=type(exc).__name__, + error_details={"detail": str(exc)}, + ) + + +@api_bp.route("/images/", methods=["GET"]) def get_image(image_id): - image = Image.query.get_or_404((image_id)) - return jsonify(image.to_json()) + """ + Retrieve a single Image by its identifier. -@api_bp.route('/images/', methods=['DELETE']) + Returns + ------- + Tuple(Response, int) + Standardised response containing the Image data. + """ + try: + image_instance = Image.query.get_or_404(image_id) + logger.debug("Fetched Image id=%s", image_id) + return api_response(data=image_instance.to_json()) + except Exception as exc: # pylint: disable=broad-except + logger.exception("Failed to retrieve Image %s", image_id) + return api_response( + success=False, + message="Internal Server Error", + status=500, + error_type=type(exc).__name__, + error_details={"detail": str(exc)}, + ) + + +@api_bp.route("/images/", methods=["DELETE"]) def delete_image(image_id): - image = Image.query.get_or_404((image_id)) - image.soft_delete() - db.session.commit() - return jsonify({'message': 'Image deleted successfully'}), 200 + """ + Soft-delete an Image. -@api_bp.route('/images', methods=['GET']) + Returns + ------- + Tuple(Response, int) + Standardised response confirming deletion. + """ + try: + image_instance = Image.query.get_or_404(image_id) + image_instance.soft_delete() + db.session.commit() + + logger.info("Image '%s' (id=%s) soft-deleted", image_instance.name, image_instance.id) + + return api_response(message="Image deleted") + except Exception as exc: # pylint: disable=broad-except + logger.exception("Failed to delete Image %s", image_id) + return api_response( + success=False, + message="Internal Server Error", + status=500, + error_type=type(exc).__name__, + error_details={"detail": str(exc)}, + ) + + +@api_bp.route("/images", methods=["GET"]) def get_images(): - images = Image.query.order_by(Image.name.asc()).all() - return jsonify([image.to_json() for image in images]) \ No newline at end of file + """ + List all Images alphabetically by name. + + Returns + ------- + Tuple(Response, int) + Standardised response containing a list of Images. + """ + try: + image_instances = Image.query.order_by(Image.name.asc()).all() + logger.debug("Fetched %d Images", len(image_instances)) + return api_response( + data=[image.to_json() for image in image_instances], + meta={"total": len(image_instances)}, + ) + except Exception as exc: # pylint: disable=broad-except + logger.exception("Failed to list Images") + return api_response( + success=False, + message="Internal Server Error", + status=500, + error_type=type(exc).__name__, + error_details={"detail": str(exc)}, + ) diff --git a/app/controller/api/virtual_data_center_routes.py b/app/controller/api/virtual_data_center_routes.py index 6bf6201..3fc90ee 100644 --- a/app/controller/api/virtual_data_center_routes.py +++ b/app/controller/api/virtual_data_center_routes.py @@ -1,53 +1,195 @@ -from flask import request, jsonify -from app import app, db, logger +""" +API endpoints for managing Virtual Data Centers within xCloudify. + +All responses use the standard `api_response` envelope to ensure a uniform +contract for client integrations. +""" + +from flask import request +from app import db, logger from app.models.models import VirtualDataCenter from app.controller import api_bp -import uuid +from app.utils.standard_responses import api_response -@api_bp.route('/virtual_data_centers', methods=['POST']) -def add_vdc(): - data = request.json - instance = VirtualDataCenter( - name=data['name'], - description=data.get('description'), - status=data.get('status'), - created_by=data.get('created_by') - ) - db.session.add(instance) - db.session.commit() - logger.debug("vdc added to DB") - return jsonify(instance.to_json()), 201 -@api_bp.route('/virtual_data_centers/', methods=['PUT']) -def edit_vdc(virtual_data_center_id): - vdc = VirtualDataCenter.query.get_or_404(virtual_data_center_id) - data = request.json - - if 'name' in data: - vdc.name = data['name'] - if 'description' in data: - vdc.description = data['description'] - if 'status' in data: - vdc.status = data['status'] - if 'visible' in data: - vdc.visible = data['visible'] - - db.session.commit() - return jsonify(success=True) +@api_bp.route("/virtual_data_centers", methods=["POST"]) +def add_virtual_data_center(): + """ + Create a new Virtual Data Center. -@api_bp.route('/virtual_data_centers/', methods=['GET']) + Returns + ------- + Tuple(Response, int) + Standardised response containing the new Virtual Data Center data. + """ + try: + request_data = request.get_json(force=True) or {} + + virtual_data_center_instance = VirtualDataCenter( + name=request_data["name"], + description=request_data.get("description"), + status=request_data.get("status"), + created_by=request_data.get("created_by"), + ) + + db.session.add(virtual_data_center_instance) + db.session.commit() + + logger.info( + "Virtual Data Center '%s' (id=%s) created", + virtual_data_center_instance.name, + virtual_data_center_instance.id, + ) + + return api_response( + data=virtual_data_center_instance.to_json(), + message="Virtual Data Center created", + status=201, + ) + except Exception as exc: # pylint: disable=broad-except + logger.exception("Failed to create Virtual Data Center") + return api_response( + success=False, + message="Internal Server Error", + status=500, + error_type=type(exc).__name__, + error_details={"detail": str(exc)}, + ) + + +@api_bp.route("/virtual_data_centers/", methods=["PUT"]) +def edit_virtual_data_center(virtual_data_center_id): + """ + Update an existing Virtual Data Center. + + Parameters + ---------- + virtual_data_center_id : str + Identifier of the Virtual Data Center to update. + + Returns + ------- + Tuple(Response, int) + Standardised response confirming the update. + """ + try: + virtual_data_center_instance = VirtualDataCenter.query.get_or_404( + virtual_data_center_id + ) + request_data = request.get_json(force=True) or {} + + # Update allowed fields only + for field in ("name", "description", "status", "visible"): + if field in request_data: + setattr(virtual_data_center_instance, field, request_data[field]) + + db.session.commit() + + logger.info( + "Virtual Data Center '%s' (id=%s) updated", + virtual_data_center_instance.name, + virtual_data_center_instance.id, + ) + + return api_response(message="Virtual Data Center updated") + except Exception as exc: # pylint: disable=broad-except + logger.exception("Failed to update Virtual Data Center %s", virtual_data_center_id) + return api_response( + success=False, + message="Internal Server Error", + status=500, + error_type=type(exc).__name__, + error_details={"detail": str(exc)}, + ) + + +@api_bp.route("/virtual_data_centers/", methods=["GET"]) def get_virtual_data_center(virtual_data_center_id): - vdc = VirtualDataCenter.query.get_or_404((virtual_data_center_id)) - return jsonify(vdc.to_json()) + """ + Retrieve a single Virtual Data Center by its identifier. -@api_bp.route('/virtual_data_centers/', methods=['DELETE']) -def delete_vdc(virtual_data_center_id): - vdc = VirtualDataCenter.query.get_or_404(virtual_data_center_id) - vdc.soft_delete() - db.session.commit() - return jsonify({'message': 'vdc deleted successfully'}), 200 + Returns + ------- + Tuple(Response, int) + Standardised response containing the requested Virtual Data Center data. + """ + try: + virtual_data_center_instance = VirtualDataCenter.query.get_or_404( + virtual_data_center_id + ) + logger.debug("Fetched Virtual Data Center id=%s", virtual_data_center_id) + return api_response(data=virtual_data_center_instance.to_json()) + except Exception as exc: # pylint: disable=broad-except + logger.exception("Failed to retrieve Virtual Data Center %s", virtual_data_center_id) + return api_response( + success=False, + message="Internal Server Error", + status=500, + error_type=type(exc).__name__, + error_details={"detail": str(exc)}, + ) -@api_bp.route('/virtual_data_centers', methods=['GET']) + +@api_bp.route("/virtual_data_centers/", methods=["DELETE"]) +def delete_virtual_data_center(virtual_data_center_id): + """ + Soft-delete a Virtual Data Center. + + Returns + ------- + Tuple(Response, int) + Standardised response confirming deletion. + """ + try: + virtual_data_center_instance = VirtualDataCenter.query.get_or_404( + virtual_data_center_id + ) + virtual_data_center_instance.soft_delete() + db.session.commit() + + logger.info( + "Virtual Data Center '%s' (id=%s) soft-deleted", + virtual_data_center_instance.name, + virtual_data_center_instance.id, + ) + + return api_response(message="Virtual Data Center deleted") + except Exception as exc: # pylint: disable=broad-except + logger.exception("Failed to delete Virtual Data Center %s", virtual_data_center_id) + return api_response( + success=False, + message="Internal Server Error", + status=500, + error_type=type(exc).__name__, + error_details={"detail": str(exc)}, + ) + + +@api_bp.route("/virtual_data_centers", methods=["GET"]) def get_virtual_data_centers(): - virtual_data_centers = VirtualDataCenter.query.order_by(VirtualDataCenter.name.asc()).all() - return jsonify([vdc.to_json() for vdc in virtual_data_centers]) + """ + List all Virtual Data Centers, ordered alphabetically by name. + + Returns + ------- + Tuple(Response, int) + Standardised response containing a list of Virtual Data Centers. + """ + try: + virtual_data_centers = ( + VirtualDataCenter.query.order_by(VirtualDataCenter.name.asc()).all() + ) + logger.debug("Fetched %d Virtual Data Centers", len(virtual_data_centers)) + return api_response( + data=[vdc.to_json() for vdc in virtual_data_centers], + meta={"total": len(virtual_data_centers)}, + ) + except Exception as exc: # pylint: disable=broad-except + logger.exception("Failed to list Virtual Data Centers") + return api_response( + success=False, + message="Internal Server Error", + status=500, + error_type=type(exc).__name__, + error_details={"detail": str(exc)}, + ) diff --git a/app/controller/api/workload_container_routes.py b/app/controller/api/workload_container_routes.py index 036fe5d..1bb3b7e 100644 --- a/app/controller/api/workload_container_routes.py +++ b/app/controller/api/workload_container_routes.py @@ -12,7 +12,6 @@ from app.utils.standard_responses import api_response from app.utils.container_deleted import check_deleted_container from app.models.models import AuditEntry from app.utils.create_workload_container import build_pod_payload -from app.schemas.schemas import WorkloadSchema, WorkloadResponse, SimpleOKResponse, PodResponse, PodListResponse, UpdateStatusSchema, WorkloadRequestSchema, WorkloadListResponse def validate_payload(payload): """ @@ -229,17 +228,7 @@ def validate_payload(payload): return sanitized_payload -@api_bp.post("/workloads/containers") -@api_bp.input(WorkloadRequestSchema, location="json") -@api_bp.doc(summary="Queue new container workload") -@api_bp.output(SimpleOKResponse, 202) -@api_bp.doc( - description=""" -Creates a **WorkloadRequest** and immediately returns 202. -A background Celery task does the heavy lifting. -""", - tags=["Workloads • Containers"] -) +@api_bp.route("/workloads/containers", methods=["POST"]) def request_container_workload(): """ Fast route: create a WorkloadRequest row and queue the heavy job. @@ -284,11 +273,7 @@ def request_container_workload(): status=202, ) -@api_bp.put("/workloads/containers/status_update/") -@api_bp.input(UpdateStatusSchema, location="json") # define tiny schema {new_status, worker_id, timestamp} -@api_bp.output(SimpleOKResponse, 200) -@api_bp.doc( summary="Worker status push") -@api_bp.doc(tags=["Workloads • Containers"]) +@api_bp.route('/workloads/containers/status_update/', methods=['PUT']) def update_container_workload(system_container_id): data = request.json logger.debug(data) @@ -422,38 +407,48 @@ def update_container_workload(system_container_id): error_details={"error": str(e)} ) -@api_bp.get("/workloads/containers/") -@api_bp.output(WorkloadResponse, 200, description="Container details") # ← note api_bp. -@api_bp.doc( - summary="Get a single container workload", # appears in list view - description=""" -Returns metadata about a **Container** or **NSController** workload. - -* **workload_id** – UUID path parameter -* Excludes soft-deleted workloads -* Includes status, resource usage, environment vars, and more -""", - tags=["Workloads • Containers"] -) +@api_bp.route('/workloads/containers/', methods=['GET']) def get_container_workload(workload_id): - workload = Workload.query.filter( - Workload.id == workload_id, - Workload.workload_type.in_(["Container", "NSController"]), - Workload.deleted.is_(False), - ).first_or_404() + try: + # Convert workload_id to UUID and ensure it's valid + uuid.UUID(str(workload_id)) + except ValueError: + # Return 404 if it's not a valid UUID + return api_response( + success=False, + message="Invalid workload ID format", + status=404, + error_type="INVALID_UUID", + error_details={"workload_id": workload_id} + ) - logger.info("Fetched %s (status=%s)", workload.id, workload.status) + try: + # Query the workload + workload = Workload.query.filter( + Workload.id == workload_id, + or_(Workload.workload_type == "Container", Workload.workload_type == "NSController"), + Workload.deleted == False + ).first_or_404() + + logger.info(f"{workload.workload_type} {workload.deleted}") + + # Return the JSON representation of the workload + return api_response( + success=True, + data=workload.to_json(), + message="Container details retrieved successfully" + ) + except Exception as e: + logger.error(f"Error retrieving container {workload_id}: {str(e)}") + return api_response( + success=False, + message="Container not found", + status=404, + error_type="CONTAINER_NOT_FOUND", + error_details={"workload_id": workload_id} + ) - return api_response( - success=True, - data=WorkloadSchema().dump(workload), - message="Container details retrieved successfully", - ) - -@api_bp.get("/workloads/containers") -@api_bp.output(WorkloadListResponse, 200) -@api_bp.doc( summary="List container workloads") -@api_bp.doc(tags=["Workloads • Containers"]) +@api_bp.route('/workloads/containers', methods=['GET']) def get_container_workloads(): try: workloads = Workload.query.filter( @@ -476,10 +471,7 @@ def get_container_workloads(): error_details={"error": str(e)} ) -@api_bp.get("/workloads/pods") -@api_bp.output(PodListResponse, 200) -@api_bp.doc( summary="List pods with active containers") -@api_bp.doc(tags=["Pods"]) +@api_bp.route('/workloads/pods', methods=['GET']) def get_pods(): try: # pods = ContainerPod.query.all() @@ -568,10 +560,7 @@ def get_pods(): error_details={"error": str(e)} ) -@api_bp.get("/workloads/pods/") -@api_bp.output(PodResponse, 200) -@api_bp.doc( summary="Get single pod") -@api_bp.doc(tags=["Pods"]) +@api_bp.route('/workloads/pods/', methods=['GET']) def get_pod(pod_id): try: pod_uuid = (pod_id) @@ -631,10 +620,7 @@ def get_pod(pod_id): error_details={"pod_id": pod_id} ) -@api_bp.delete("/workloads/containers/") -@api_bp.output(SimpleOKResponse, 200) -@api_bp.doc( summary="Mark container for deletion") -@api_bp.doc(tags=["Workloads • Containers"]) +@api_bp.route('/workloads/containers/', methods=['DELETE']) def delete_container_workload(workload_id): """ Mark a single container for deletion via pod-update. @@ -685,34 +671,15 @@ def delete_container_workload(workload_id): error_type="CONTAINER_DELETE_ERROR", error_details={"error": str(e)} ) -@api_bp.post( - "/workloads/containers//lifecycle/" -) -@api_bp.output(SimpleOKResponse, 202, - description="Task queued – see X-Request-ID header for trace") -@api_bp.doc( - # Document the two path parameters explicitly - - tags=["Workloads • Containers"], -) + +@api_bp.route('/workloads/containers//lifecycle/', methods=['POST']) def container_lifecycle_action(workload_id, action): """ - Run a lifecycle action on a container - Queue a **Celery** task that performs one of three lifecycle operations - on an existing Container / NSController workload. - - The call is non-blocking – you receive **HTTP 202** immediately and can - follow the progress via the event bus or by polling `/workloads/containers/{id}`. - - **Accepted actions** - - | verb | effect on container | - |---------|---------------------| - | start | boot a stopped or never-started container | - | stop | gracefully stop a running container | - | restart | stop + start in a single transaction | + Perform lifecycle operations (start, stop, restart) on a container. + + This is a non-blocking API that queues a Celery task. """ - valid_actions = {"start", "stop", "restart"} + valid_actions = ["start", "stop", "restart"] if action not in valid_actions: return api_response( success=False, @@ -746,13 +713,7 @@ def container_lifecycle_action(workload_id, action): status=500 ) -@api_bp.post("/workloads/pods//lifecycle/") -@api_bp.output(SimpleOKResponse, 202) -@api_bp.doc( summary="Lifecycle op on pod") -@api_bp.doc( - description="Applies to **all** containers in the pod", - tags=["Pods"] -) +@api_bp.route('/workloads/pods//lifecycle/', methods=['POST']) def pod_lifecycle_action(pod_id, action): """ Perform lifecycle operations (start, stop, restart) on all containers in a pod. @@ -789,10 +750,7 @@ def pod_lifecycle_action(pod_id, action): status=500 ) -@api_bp.delete("/workloads/pods/") -@api_bp.output(SimpleOKResponse, 200) -@api_bp.doc( summary="Delete pod and its containers") -@api_bp.doc(tags=["Pods"]) +@api_bp.route('/workloads/pods/', methods=['DELETE']) def delete_pod(pod_id): """ Delete an entire pod and all its containers via pod-update. diff --git a/app/models/models.py b/app/models/models.py index 2b55149..3d888a3 100644 --- a/app/models/models.py +++ b/app/models/models.py @@ -551,6 +551,7 @@ class Workload(BaseModel): "pending-stop": ["deleted","running","stopped"], "pending-start": ["deleted","running","stopped","launch_failed"], "pending-restart": ["deleted","running","stopped","launch_failed"], + "host_failed": ["*"], } if old_status not in valid_transitions: diff --git a/app/tasks/host_monitoring.py b/app/tasks/host_monitoring.py index d669a70..3c0e380 100644 --- a/app/tasks/host_monitoring.py +++ b/app/tasks/host_monitoring.py @@ -54,6 +54,11 @@ def handle_host_downtime(self, host_id: str) -> None: for workload in workloads: # Update workload status to host_failed + + if workload.status not in ["running"]: + logger.debug(f"Skipping setting workload {workload.id} to filed, becuase it's currrently {workload.status}") + continue + workload.set_status("host_failed") # Collect pod IDs and VM IDs diff --git a/app/utils/create_workload_container.py b/app/utils/create_workload_container.py index 10074e3..20ba1f0 100644 --- a/app/utils/create_workload_container.py +++ b/app/utils/create_workload_container.py @@ -129,31 +129,6 @@ def _create_new_pod_flow(request_row: WorkloadRequest, payload: Dict): SortByPlacementPriority(), ] - - - # # Hypothetical setup - # customer_id = workload.owner_id - # affinity_label = workload.labels.get("Affinity:web-servers") - # anti_affinity_label = workload.labels.get("AntiAffinity:web-servers") - # all_workloads = get_all_workloads_for_customer(customer_id) - - # filters = [] - - # if affinity_label: - # filters.append(AffinityFilter("Affinity:web-servers", customer_id, all_workloads)) - - # if anti_affinity_label: - # filters.append(AntiAffinityFilter("AntiAffinity:web-servers", customer_id, all_workloads)) - - # filters += [ - # ExcludeHostsWithoutNorthSouthIP(), - # DockerCapableHosts(), - # ExcludeAllOfflineHosts(), - # ExcludeAllDisabledHosts(), - # HasSufficientResources(), - # SortByPlacementPriority(), - # ] - requirements = { "cpu": {"min": 4}, "ram": {"min": 8192}, @@ -173,8 +148,16 @@ def _create_new_pod_flow(request_row: WorkloadRequest, payload: Dict): if not eligible_hosts: logger.warning(f"No eligible hosts available for region {vdc.region.id}") - host = eligible_hosts[0] - logger.info("Selected host %s for new pod", host.id) + if eligible_hosts: + host = eligible_hosts[0] + logger.info("Selected host %s for new pod", host.id) + else: + logger.warning(f"No eligable hosts found, setting state to waiting_host") + host=None + # TODO - This is a big chunk of work, but we need to fork the _build_pod_objects_and_enqueue + # to enable the workload to be created in the DB but not scheduled to a task or a host + # the workload should sit in waiting_host till a new ost joins the region and we try to schedule it there + return # 2. Build DB objects & enqueue task _build_pod_objects_and_enqueue( diff --git a/app/utils/standard_responses.py b/app/utils/standard_responses.py index 985c70f..5778996 100644 --- a/app/utils/standard_responses.py +++ b/app/utils/standard_responses.py @@ -50,5 +50,5 @@ def api_response(*, data=None, success=True, # ------------------------------------------------------------------ # DEBUG log: one-liner with request_id and compact JSON payload. # ------------------------------------------------------------------ - logger.debug("api_response (request_id=%s): %s", g.request_id, payload) + # logger.debug("api_response (request_id=%s): %s", g.request_id, payload) return jsonify(payload), status diff --git a/websocket_server/events/logs.py b/websocket_server/events/logs.py index b12e54e..5b4f55a 100644 --- a/websocket_server/events/logs.py +++ b/websocket_server/events/logs.py @@ -36,7 +36,7 @@ def register_socketio_handlers(socketio): }, to=request.sid) return - container_info = container_response.json() + container_info = container_response.json()['data'] worker_id = container_info["workload_host_id"] if worker_id not in connected_workers: @@ -110,7 +110,7 @@ def register_socketio_handlers(socketio): 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.raise_for_status() - container_info = container_response.json() + container_info = container_response.json()['data'] worker_id = container_info["workload_host_id"] if worker_id not in connected_workers: diff --git a/websocket_server/events/terminal.py b/websocket_server/events/terminal.py index 50774ed..5933c37 100644 --- a/websocket_server/events/terminal.py +++ b/websocket_server/events/terminal.py @@ -28,7 +28,7 @@ def register_socketio_handlers(socketio): # Lookup container details response = requests.get(f"{api_server_url}/workloads/containers/{container_id}") response.raise_for_status() - info = response.json() + info = response.json()['data'] worker_id = info["workload_host_id"] context = { diff --git a/worker/docker-compose.yml b/worker/docker-compose.yml index 981b35b..a86257f 100644 --- a/worker/docker-compose.yml +++ b/worker/docker-compose.yml @@ -7,7 +7,7 @@ services: #user: "1000:120" restart: unless-stopped environment: - - WEBSOCKET_SERVER_URL=http://192.168.64.2:6001 + - WEBSOCKET_SERVER_URL=https://ws.xcloudify.tech - REGION_ID=1 - REGION_ENROLLMENT_KEY=a - API_BASE_URL=http://192.168.64.2:5000