From ff989be20266340bc5e85fca595ebac1c1eaf3b1 Mon Sep 17 00:00:00 2001 From: Cory Hawklvelt Date: Fri, 25 Jul 2025 15:38:52 +0930 Subject: [PATCH] PortForwarding work Added workload id to portforward object Updated streamlit to show workload name --- .../api/workload_container_routes.py | 125 +++++++++++------- app/models/models.py | 2 + app/utils/create_workload_container.py | 3 +- docs/curl to add container to pod.md | 31 +++++ docs/curl to launch containers.md | 4 +- streamlit_server/views/containers.py | 3 +- 6 files changed, 119 insertions(+), 49 deletions(-) create mode 100644 docs/curl to add container to pod.md diff --git a/app/controller/api/workload_container_routes.py b/app/controller/api/workload_container_routes.py index c59b5c1..03b4572 100644 --- a/app/controller/api/workload_container_routes.py +++ b/app/controller/api/workload_container_routes.py @@ -1,6 +1,4 @@ #app/controller/api/workload_container_routes.py -import random -import time from flask import json, request, jsonify, abort import requests from app import app, db, logger @@ -8,9 +6,6 @@ from app.models.models import PortForwarding, Workload, WorkloadHost, WorkloadHo from app.models.network import Network, NetworkPort from datetime import datetime from app.controller import api_bp -import uuid -from app.scheduling_filters import ExcludeAllDisabledHosts, MostAvailableCapacity, ExcludeAllOfflineHosts, DockerCapableHosts, ExcludeHostsWithoutNorthSouthIP -from app.services.cloudflare import CloudflareTunnelManager from werkzeug.exceptions import abort from sqlalchemy import or_ from app.utils.standard_responses import api_response @@ -18,40 +13,55 @@ from app.utils.container_deleted import check_deleted_container def validate_payload(payload): """ - Validate the input payload according to the specified rules and return a sanitized version. - + 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. """ - # Ensure the input is a valid JSON object (dict) + import uuid + if not isinstance(payload, dict): raise ValueError("Input must be a valid JSON object (dict).") - - # Check for the top-level 'virtual_data_center' key and validate it as a UUID - if 'virtual_data_center' not in payload: - raise ValueError("Top-level 'virtual_data_center' key is missing.") - try: - uuid.UUID(payload['virtual_data_center'], version=4) - except ValueError: - raise ValueError("'virtual_data_center' must be a valid UUID.") - - # Check for the 'containers' key and ensure it is a list with at least one container + + 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.") - # Initialize the sanitized payload - sanitized_payload = { - 'virtual_data_center': payload['virtual_data_center'], - 'containers': [] - } + sanitized_payload['containers'] = [] # Validate each container in the 'containers' list and build the sanitized version for container in payload['containers']: @@ -100,11 +110,12 @@ def validate_payload(payload): '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 @@ -129,9 +140,8 @@ def validate_payload(payload): # Add the sanitized container to the sanitized payload sanitized_payload['containers'].append(sanitized_container) - - return sanitized_payload + return sanitized_payload @api_bp.route("/workloads/containers", methods=["POST"]) def request_container_workload(): @@ -144,22 +154,32 @@ def request_container_workload(): logger.warning("Validation error: %s", exc) return api_response(success=False, message=str(exc), status=400)[0] - request_vdc = VirtualDataCenter.query.filter_by(id=(validated_data['virtual_data_center'])).first() + # 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 {request_vdc}") + 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=validated_data["virtual_data_center"], + 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( @@ -168,6 +188,7 @@ def request_container_workload(): status=202, )[0] + @api_bp.route('/workloads/containers/status_update/', methods=['PUT']) def update_container_workload(system_container_id): data = request.json @@ -314,7 +335,9 @@ def get_pods(): "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 + "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 ], @@ -369,7 +392,9 @@ def get_pod(pod_id): "internal_port": pf.internal_port, "external_port": pf.external_port, "protocol": pf.protocol, - "ip_address": pf.ip_address + "ip_address": pf.ip_address, + "container_id": pf.workload_id, + "container_name": pf.workload.name, } for pf in pod.port_forwardings ] @@ -378,18 +403,22 @@ def get_pod(pod_id): @api_bp.route('/workloads/containers/', 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: - # Convert workload_id to UUID and ensure it's valid - workload_uuid = (workload_id) + workload_uuid = workload_id # Wrap with uuid.UUID(workload_id) if needed except ValueError: - # Return 404 if it's not a valid UUID return jsonify({"error": "Invalid workload ID"}), 404 - # Query the workload _container = Workload.query.filter( Workload.id == workload_uuid, - or_(Workload.workload_type == "Container", Workload.workload_type == "NSController"), + or_( + Workload.workload_type == "Container", + Workload.workload_type == "NSController" + ), Workload.deleted == False ).first_or_404() @@ -407,15 +436,23 @@ def delete_container_workload(workload_id): ] } } + 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) - websocket_server_response_data = websocket_server_response.json() 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: @@ -424,8 +461,6 @@ def delete_container_workload(workload_id): db.session.add(_container) db.session.commit() - # Dont soft delete here, it's only in delete-pending right now - # _container.soft_delete() logger.info(f"Requested deletion of container {_container.id}") return jsonify({'message': 'Container workload deleted successfully'}), 200 diff --git a/app/models/models.py b/app/models/models.py index fe95bba..51f5772 100644 --- a/app/models/models.py +++ b/app/models/models.py @@ -862,9 +862,11 @@ class PortForwarding(BaseModel): protocol = Column(String(10), default="tcp", nullable=False) ip_address = Column(String(25), nullable=False) dns_record_id = Column(String(36), ForeignKey("cloudflare_dns_records.id"), nullable=True) + workload_id = Column(String(36), ForeignKey("workloads.id"), nullable=True) pod = relationship("ContainerPod", backref="port_forwardings") dns_record = relationship("CloudflareDNSRecord", backref="port_forwardings") + workload=relationship("Workload", backref="port_forwardings") def to_json(self): data = super().to_json() diff --git a/app/utils/create_workload_container.py b/app/utils/create_workload_container.py index 84e3434..cd71c24 100644 --- a/app/utils/create_workload_container.py +++ b/app/utils/create_workload_container.py @@ -142,7 +142,8 @@ def add_container_workload(request_data): internal_port=internal_port, external_port=assigned_external_port, protocol=port_mapping.get("protocol", "tcp"), - ip_address=selected_host.ip_address_northsouth + ip_address=selected_host.ip_address_northsouth, + workload_id=new_container.id, ) db.session.add(new_forward) db.session.flush() # Ensure we have new_forward.id diff --git a/docs/curl to add container to pod.md b/docs/curl to add container to pod.md new file mode 100644 index 0000000..975f455 --- /dev/null +++ b/docs/curl to add container to pod.md @@ -0,0 +1,31 @@ +PODID="" + +curl -X POST http://192.168.64.2:5000/api/workloads/containers \ + -H "Content-Type: application/json" \ + -d '{ + "pod": "$PODID", + "containers": [ + { + "docker_image": "joke_container", + "container_name": "very-important-webserver", + "ports": [ + { + "internal": 8080, + "external": 8765, + "use_dns": false + } + ] + }, + { + "docker_image": "whoami81", + "container_name": "whoami-server", + "ports": [ + { + "internal": 81, + "external": 8081, + "use_dns": false + } + ] + } + ] + }' diff --git a/docs/curl to launch containers.md b/docs/curl to launch containers.md index 19b4312..cf12d23 100644 --- a/docs/curl to launch containers.md +++ b/docs/curl to launch containers.md @@ -1,7 +1,7 @@ -curl -X POST http://localhost:5000/api/workloads/containers \ +curl -X POST http://192.168.64.2:5000/api/workloads/containers \ -H "Content-Type: application/json" \ -d '{ - "vdc": "ce11e1aa-18ca-4128-8d70-c09662f0c039", + "virtual_data_center": "b1d09477-742b-485e-87c3-5dad36dd4f9e", "containers": [ { "docker_image": "joke_container", diff --git a/streamlit_server/views/containers.py b/streamlit_server/views/containers.py index cd2a371..4e8a834 100644 --- a/streamlit_server/views/containers.py +++ b/streamlit_server/views/containers.py @@ -145,10 +145,11 @@ def render_list(): st.markdown("🔁 **Port Forwarding:**") for pf in pod['port_forwardings']: dns_status = "🌐 DNS" if pf.get('dns_record_id') else "🚫 No DNS" - cols = st.columns([1, 3, 2]) + cols = st.columns([1, 2, 2,1]) with cols[0]: st.markdown(f"**{pf['protocol'].upper()}**") with cols[1]: st.text(f"{pf['external_port']} → {pf['internal_port']}") with cols[2]: st.text(dns_status) + with cols[3]: st.text(pf['container_name']) link = pf.get('dns_record_hostname') or f"{pf['ip_address']}:{pf['external_port']}" st.page_link(f"http://{link}", label=f"http://{link}") st.markdown("---")