Files
3cloud-backend/app/controller/api/workload_container_routes.py
T
coryHawkvelt ff989be202 PortForwarding work
Added workload id to portforward object
Updated streamlit to show workload name
2025-07-25 15:38:52 +09:30

549 lines
22 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
#app/controller/api/workload_container_routes.py
from flask import json, request, jsonify, abort
import requests
from app import app, db, logger
from app.models.models import PortForwarding, Workload, WorkloadHost, WorkloadHostPortMapping, VirtualDataCenter, Label, VolumeWorkloadMapping, ContainerPod, ContainerPodContainer, CloudflareDNSRecord, CloudflareTunnel, WorkloadRequest
from app.models.network import Network, NetworkPort
from datetime import datetime
from app.controller import api_bp
from werkzeug.exceptions import abort
from sqlalchemy import or_
from app.utils.standard_responses import api_response
from app.utils.container_deleted import check_deleted_container
def validate_payload(payload):
"""
Validate the input payload for container deployment.
Exactly one of 'virtual_data_center' or 'pod' must be provided.
Args:
payload (dict): The input JSON object to validate.
Returns:
dict: A sanitized version of the input payload containing only the expected parameters.
Raises:
ValueError: If the payload fails any validation rule.
"""
import uuid
if not isinstance(payload, dict):
raise ValueError("Input must be a valid JSON object (dict).")
sanitized_payload = {}
# Enforce exactly one of 'virtual_data_center' or 'pod'
has_vdc = 'virtual_data_center' in payload
has_pod = 'pod' in payload
if not (has_vdc or has_pod):
raise ValueError("Either 'virtual_data_center' or 'pod' must be provided.")
if has_vdc and has_pod:
raise ValueError("Only one of 'virtual_data_center' or 'pod' can be provided — not both.")
if has_vdc:
try:
uuid.UUID(payload['virtual_data_center'], version=4)
sanitized_payload['virtual_data_center'] = payload['virtual_data_center']
except ValueError:
raise ValueError("'virtual_data_center' must be a valid UUID.")
if has_pod:
try:
uuid.UUID(payload['pod'], version=4)
sanitized_payload['pod'] = payload['pod']
except ValueError:
raise ValueError("'pod' must be a valid UUID.")
# Validate container list
if 'containers' not in payload:
raise ValueError("'containers' key is missing.")
if not isinstance(payload['containers'], list) or len(payload['containers']) == 0:
raise ValueError("'containers' must be a list with at least one container.")
sanitized_payload['containers'] = []
# Validate each container in the 'containers' list and build the sanitized version
for container in payload['containers']:
if not isinstance(container, dict):
raise ValueError("Each container must be a dictionary.")
# Check for required keys: 'docker_image' and 'container_name'
if 'docker_image' not in container:
raise ValueError("Container is missing 'docker_image'.")
if 'container_name' not in container:
raise ValueError("Container is missing 'container_name'.")
# Initialize the sanitized container
sanitized_container = {
'docker_image': container['docker_image'],
'container_name': container['container_name']
}
# Add 'cpu_shares' if present and valid
if 'cpu_shares' in container:
if not isinstance(container['cpu_shares'], (int, float)) or container['cpu_shares'] <= 0:
raise ValueError("'cpu_shares' must be a positive number.")
sanitized_container['cpu_shares'] = container['cpu_shares']
# Add 'mem_limit' if present and valid
if 'mem_limit' in container:
if not isinstance(container['mem_limit'], (int, float)) or container['mem_limit'] <= 0:
raise ValueError("'mem_limit' must be a positive number.")
sanitized_container['mem_limit'] = container['mem_limit']
# Add 'ports' if present and valid
if 'ports' in container:
if not isinstance(container['ports'], list):
raise ValueError("'ports' must be a list.")
sanitized_ports = []
for port_mapping in container['ports']:
if not isinstance(port_mapping, dict):
raise ValueError("Each port mapping must be a dictionary with 'internal' and 'external' keys.")
if 'internal' not in port_mapping or 'external' not in port_mapping:
raise ValueError("Each port mapping must have 'internal' and 'external' keys.")
if not isinstance(port_mapping['internal'], int) or not isinstance(port_mapping['external'], int):
raise ValueError("'internal' and 'external' ports must be integers.")
# Handle use_dns in port mapping
port_data = {
'internal': port_mapping['internal'],
'external': port_mapping['external']
}
if 'use_dns' in port_mapping:
port_data['use_dns'] = bool(port_mapping['use_dns'])
sanitized_ports.append(port_data)
if sanitized_ports:
sanitized_container['ports'] = sanitized_ports
# Add 'networks' if present and valid
if 'networks' in container:
if isinstance(container['networks'], str):
if container['networks'].strip():
sanitized_container['networks'] = container['networks'].strip()
elif isinstance(container['networks'], list):
sanitized_networks = []
for net in container['networks']:
if not isinstance(net, str):
raise ValueError("Each network name must be a string.")
if net.strip():
sanitized_networks.append(net.strip())
if sanitized_networks:
sanitized_container['networks'] = sanitized_networks
elif container['networks'] is None:
pass # allow it to be missing or null
else:
raise ValueError("'networks' must be a string or a list of strings.")
# Add the sanitized container to the sanitized payload
sanitized_payload['containers'].append(sanitized_container)
return sanitized_payload
@api_bp.route("/workloads/containers", methods=["POST"])
def request_container_workload():
"""
Fast route: create a WorkloadRequest row and queue the heavy job.
"""
try:
validated_data = validate_payload(request.get_json(force=True))
except ValueError as exc:
logger.warning("Validation error: %s", exc)
return api_response(success=False, message=str(exc), status=400)[0]
# Resolve Virtual Data Center ID from either pod or direct input
if 'pod' in validated_data:
container_pod = ContainerPod.query.filter_by(id=validated_data['pod']).first()
if not container_pod:
logger.warning(f"Referenced pod ID not found: {validated_data['pod']}")
abort(404, description="Container Pod not found")
vdc_id = container_pod.vdc_id
else:
vdc_id = validated_data['virtual_data_center']
request_vdc = VirtualDataCenter.query.filter_by(id=vdc_id).first()
if request_vdc is None:
logger.warning(f"Requested a Virtual Data Center that does not exist: {vdc_id}")
abort(404, description="Virtual Data Center not found")
req = WorkloadRequest(
vdc_id=vdc_id,
payload=validated_data,
status="pending-scheduling",
workload_type="container"
)
db.session.add(req)
db.session.commit()
logger.info("Created workload request %s – enqueuing Celery task", req.id)
from app.tasks.process_workload_request import process_workload
process_workload.delay(str(req.id))
return api_response(
data={"request_id": str(req.id)},
message="Workload queued for scheduling",
status=202,
)[0]
@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)
# TODO - need to add error to autid table if there is one
# Validate new_status
valid_statuses = ["running", "deleted", "stopped","dead", "launch_failed"]
new_status=data.get("new_status")
if new_status not in valid_statuses:
error_message = f"Invalid status {new_status} Must be one of: {', '.join(valid_statuses)}"
logger.error(error_message)
return jsonify({"success": False, "message": error_message}), 400
# Validate worker_id as UUIDs
try:
worker_id = (data.get("worker_id"))
except (ValueError, TypeError) as e:
error_message = "Invalid UUID format for worker_id."
logger.error(f"{error_message} Error: {str(e)}")
return jsonify({"success": False, "message": error_message}), 400
# Validate system_container_id as UUID
try:
system_container_id = (system_container_id)
except (ValueError, TypeError) as e:
error_message = "Invalid UUID format for system_container_id."
logger.error(f"{error_message} Error: {str(e)}")
return jsonify({"success": False, "message": error_message}), 400
# Validate timestamp as ISO format
try:
datetime.fromisoformat(data.get("timestamp"))
except (ValueError, TypeError) as e:
error_message = "Invalid ISO timestamp format."
logger.error(f"{error_message} Error: {str(e)}")
return jsonify({"success": False, "message": error_message}), 400
# Fetch the container from the database
try:
container = Workload.query.filter(
Workload.id == system_container_id,
or_(Workload.workload_type == "Container", Workload.workload_type == "NSController"),
Workload.deleted == False
).first_or_404()
except Exception as e:
error_message = f"Failed to fetch container with ID {system_container_id}."
logger.error(f"{error_message} Error: {str(e)}")
return jsonify({"success": False, "message": error_message}), 404
# Check if the container is actually on the reported worker
if container.workload_host_id != worker_id:
error_message = f"Workload host in DB for container ID {system_container_id} is {container.workload_host_id}, but was reported from {worker_id}. Ignoring."
logger.error(error_message)
return jsonify({"success": False, "message": error_message}), 400
# Update the status of the container
try:
new_status = data.get('new_status')
logger.info(f"Updating Container ID: {system_container_id} to new status {new_status}")
container.set_status(new_status)
if new_status == "deleted":
container.soft_delete()
# If container status is "deleted", check if all containers in the same pod are deleted
check_deleted_container(container)
db.session.commit()
return jsonify({"success": True, "message": "Status updated successfully."}), 200
except Exception as e:
error_message = f"Failed to update status for Container ID: {system_container_id}."
logger.error(f"{error_message} Error: {str(e)}")
db.session.rollback()
return jsonify({"success": False, "message": error_message}), 500
@api_bp.route('/workloads/containers/<workload_id>', methods=['GET'])
def get_container_workload(workload_id):
try:
# Convert workload_id to UUID and ensure it's valid
workload_uuid = (workload_id)
except ValueError:
# Return 404 if it's not a valid UUID
return jsonify({"error": "Invalid workload ID"}), 404
# Query the workload
workload = Workload.query.filter(
Workload.id == workload_uuid,
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 jsonify(workload.to_json())
@api_bp.route('/workloads/containers', methods=['GET'])
def get_container_workloads():
workloads = Workload.query.filter(or_(Workload.workload_type == "Container", Workload.workload_type == "NSController"), Workload.deleted==False).all()
return jsonify([workload.to_json() for workload in workloads])
@api_bp.route('/workloads/pods', methods=['GET'])
def get_pods():
# pods = ContainerPod.query.all()
active_pods = ContainerPod.query_with_only_active_containers()
response = []
for item in active_pods:
pod=item['pod']
containers=item['containers']
# Get associated tunnel (if exists)
tunnel = CloudflareTunnel.query.filter_by(
nscontroller_workload_id=pod.nscontroller_workload_id,
deleted=False
).first()
# Get DNS records for this tunnel
dns_records = []
if tunnel:
dns_records = CloudflareDNSRecord.query.filter_by(
tunnel_id=tunnel.id,
deleted=False
).all()
pod_data = {
"pod_id": str(pod.id),
"pod_name": pod.name,
"workload_host_id": str(pod.workload_host_id),
"vdc_id": str(pod.vdc_id),
"region_id": str(pod.vdc.region_id) if pod.vdc else None,
"nscontroller_workload_id": str(pod.nscontroller_workload_id),
"containers": [
{
"container_id": str(_container.id),
"container_name": _container.name,
"status": _container.status,
"launch_params": json.loads(_container.launch_params or "{}")
}
for _container in containers
],
"port_forwardings": [
{
"internal_port": pf.internal_port,
"external_port": pf.external_port,
"protocol": pf.protocol,
"ip_address": pf.ip_address,
"dns_record_id": pf.dns_record_id,
"dns_record_hostname": pf.dns_record.hostname if pf.dns_record else None,
"container_id": pf.workload_id,
"container_name": pf.workload.name if pf.workload else None,
}
for pf in pod.port_forwardings
],
"cloudflare_tunnel": {
"tunnel_id": tunnel.tunnel_id if tunnel else None,
"account_id": tunnel.account_id if tunnel else None,
"associated_hostname": tunnel.associated_hostname if tunnel else None,
"dns_records": [
{
"hostname": record.hostname,
"record_type": record.record_type,
"content": record.content,
"proxied": record.proxied,
"ttl": record.ttl
}
for record in dns_records
]
} if tunnel else None
}
response.append(pod_data)
return jsonify(response), 200
@api_bp.route('/workloads/pods/<pod_id>', methods=['GET'])
def get_pod(pod_id):
try:
pod_uuid = (pod_id)
except ValueError:
return jsonify({"error": "Invalid pod ID"}), 404
pod = ContainerPod.query.filter_by(id=pod_uuid).first_or_404()
response = {
"pod_id": str(pod.id),
"pod_name": pod.name,
"workload_host_id": str(pod.workload_host_id),
"vdc_id": str(pod.vdc_id),
"region_id": str(pod.vdc.region_id) if pod.vdc else None,
"nscontroller_workload_id": str(pod.nscontroller_workload_id),
"containers": [
{
"container_id": str(mapping.container.id),
"container_name": mapping.container.name,
"status": mapping.container.status,
"launch_params": json.loads(mapping.container.launch_params or "{}")
}
for mapping in pod.container_mappings
],
"port_forwardings": [
{
"internal_port": pf.internal_port,
"external_port": pf.external_port,
"protocol": pf.protocol,
"ip_address": pf.ip_address,
"container_id": pf.workload_id,
"container_name": pf.workload.name,
}
for pf in pod.port_forwardings
]
}
return jsonify(response), 200
@api_bp.route('/workloads/containers/<workload_id>', methods=['DELETE'])
def delete_container_workload(workload_id):
"""
Handle deletion request for a container workload.
Sends the deletion request to the websocket server and marks status accordingly.
Also soft-deletes any associated PortForwarding records.
"""
try:
workload_uuid = workload_id # Wrap with uuid.UUID(workload_id) if needed
except ValueError:
return jsonify({"error": "Invalid workload ID"}), 404
_container = Workload.query.filter(
Workload.id == workload_uuid,
or_(
Workload.workload_type == "Container",
Workload.workload_type == "NSController"
),
Workload.deleted == False
).first_or_404()
# TODO - Dont send a container deletes - send a pod update with desired state=deleted for the container
#Send the request off to the websocket server to have this container deleted
payload = {
"worker_id": _container.workload_host_id,
"task_type": "container-delete",
"job_details": {
"container": [
{
"container_id": _container.id,
"desired_state": "deleted"
}
]
}
}
logger.debug(f"Sending payload to websocket server {payload}")
headers = {"Content-Type": "application/json"}
websocket_server_response = requests.post(app.config['WEBSOCKET_SERVER_URL'], data=json.dumps(payload), headers=headers)
if websocket_server_response.status_code == 201:
logger.info(f"Task created successfully for container {_container.id}")
_container.set_status("pending-deleted")
# Soft delete any associated port forwards
forwards = PortForwarding.query.filter_by(container_workload_id=_container.id, deleted=False).all()
if forwards:
logger.debug(f"Soft-deleting {len(forwards)} port forward(s) for container {_container.id}")
for pf in forwards:
pf.soft_delete()
db.session.add(pf)
db.session.add(_container)
db.session.commit()
else:
logger.error(f"Deletion request failed for container {_container.id}")
_container.set_status("failed-deleting")
db.session.add(_container)
db.session.commit()
logger.info(f"Requested deletion of container {_container.id}")
return jsonify({'message': 'Container workload deleted successfully'}), 200
@api_bp.route('/workloads/pods/<pod_id>', methods=['DELETE'])
def delete_pod(pod_id):
"""
Delete a pod and all its containers:
1. Mark all containers as pending-deleted
2. Mark the pod as pending-deleted
3. Send tasks via websocket to update each container's desired state to "deleted"
4. Delete any associated Cloudflare tunnels if the pod has an NSController
"""
try:
pod_uuid = (pod_id)
except ValueError:
return jsonify({"error": "Invalid pod ID"}), 404
pod = ContainerPod.query.filter_by(id=pod_uuid).first_or_404()
# Cleanup host port mappings
for mapping in pod.container_mappings:
port_mappings = WorkloadHostPortMapping.query.filter_by(container_workload_id=mapping.container.id).all()
for pm in port_mappings:
db.session.delete(pm)
# Set desired state on the pod
pod.status = "pending-deleted"
db.session.add(pod)
# Mark all containers inside as "pending-deleted" (for tracking)
# and collect them for the websocket request
container_delete_requests = []
for mapping in pod.container_mappings:
container = mapping.container
container.set_status("pending-deleted")
db.session.add(container)
# Add this container to our websocket deletion request
container_delete_requests.append({
"container_id": str(container.id),
"desired_state": "deleted"
})
# Also make sure to include the NSController container for deletion
nscontroller = None
if pod.nscontroller_workload_id:
nscontroller = Workload.query.get(pod.nscontroller_workload_id)
if nscontroller and not nscontroller.deleted:
nscontroller.set_status("pending-deleted")
db.session.add(nscontroller)
db.session.commit()
# Send websocket task to delete all containers
if container_delete_requests:
payload = {
"worker_id": str(pod.workload_host_id),
"task_type": "pod-update",
"job_details": {
"pod_id": pod_id,
"nscontroller": {
"container_id": str(nscontroller.id) if nscontroller else None,
"desired_state": "deleted"
},
"containers": container_delete_requests
}
}
logger.info(f"Sending deletion payload for pod {pod_id} to websocket server: {payload}")
headers = {"Content-Type": "application/json"}
try:
websocket_server_response = requests.post(app.config['WEBSOCKET_SERVER_URL'], data=json.dumps(payload), headers=headers)
if websocket_server_response.status_code == 201:
logger.info(f"Deletion tasks created successfully for pod {pod.id} with {len(container_delete_requests)} containers")
else:
logger.error(f"Deletion request failed for pod {pod.id}. Response: {websocket_server_response.status_code} - {websocket_server_response.text}")
# Continue with the deletion process anyway, as the status updates may come through other channels
except Exception as e:
logger.error(f"Exception when sending deletion request for pod {pod.id}: {str(e)}")
# Continue with the process despite the error
logger.info(f"Marked pod {pod.id} and its containers as pending-deleted")
return jsonify({'message': 'Pod and containers set to pending-deleted'}), 200