moving everything into standard api responses

This commit is contained in:
2025-08-04 07:38:40 +09:30
parent a709b59214
commit beae6b181c
12 changed files with 540 additions and 282 deletions
+57 -33
View File
@@ -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")
+4 -4
View File
@@ -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.
+220 -75
View File
@@ -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/<image_id>', 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/<image_id>", 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/<image_id>', 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/<image_id>", 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/<image_id>', 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/<image_id>", 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])
"""
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)},
)
+186 -44
View File
@@ -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/<virtual_data_center_id>', 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/<virtual_data_center_id>', 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/<virtual_data_center_id>", 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/<virtual_data_center_id>", 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/<virtual_data_center_id>', 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/<virtual_data_center_id>", 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)},
)
+52 -94
View File
@@ -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/<uuid:system_container_id>")
@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/<system_container_id>', 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/<uuid:workload_id>")
@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/<workload_id>', 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/<uuid:pod_id>")
@api_bp.output(PodResponse, 200)
@api_bp.doc( summary="Get single pod")
@api_bp.doc(tags=["Pods"])
@api_bp.route('/workloads/pods/<pod_id>', 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/<uuid:workload_id>")
@api_bp.output(SimpleOKResponse, 200)
@api_bp.doc( summary="Mark container for deletion")
@api_bp.doc(tags=["Workloads • Containers"])
@api_bp.route('/workloads/containers/<workload_id>', 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/<uuid:workload_id>/lifecycle/<string:action>"
)
@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/<workload_id>/lifecycle/<action>', 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/<uuid:pod_id>/lifecycle/<string:action>")
@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/<pod_id>/lifecycle/<action>', 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/<uuid:pod_id>")
@api_bp.output(SimpleOKResponse, 200)
@api_bp.doc( summary="Delete pod and its containers")
@api_bp.doc(tags=["Pods"])
@api_bp.route('/workloads/pods/<pod_id>', methods=['DELETE'])
def delete_pod(pod_id):
"""
Delete an entire pod and all its containers via pod-update.
+1
View File
@@ -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:
+5
View File
@@ -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
+10 -27
View File
@@ -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(
+1 -1
View File
@@ -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
+2 -2
View File
@@ -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:
+1 -1
View File
@@ -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 = {
+1 -1
View File
@@ -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