From 904b2e011d40d450fb9baad2eb14c0faa96d090d Mon Sep 17 00:00:00 2001 From: Cory Hawklvelt Date: Sat, 26 Jul 2025 10:31:59 +0930 Subject: [PATCH] Adding a new container to an existing pod doesnt break everything! --- app/Dockerfile | 29 +++ app/__init__.py | 3 + app/celery_app.py | 6 +- .../api/workload_container_routes.py | 176 +++++++----------- app/controller/api/workload_host_routes.py | 2 - app/models/models.py | 19 +- app/requirements.txt | 5 +- app/tasks/reconcile_online_workers.py | 76 +++++--- app/templates/container_logs.html | 2 +- app/utils/container_deleted.py | 28 ++- app/utils/create_workload_container.py | 16 +- docker-compose.yml | 11 +- docs/curl to add container to pod.md | 8 +- docs/curl to launch containers.md | 35 ++++ websocket_server/events/logs.py | 4 +- websocket_server/events/terminal.py | 2 +- websocket_server/events/vnc.py | 2 +- websocket_server/worker_manager.py | 1 - worker/worker_tasks/container.py | 8 +- 19 files changed, 248 insertions(+), 185 deletions(-) create mode 100644 app/Dockerfile diff --git a/app/Dockerfile b/app/Dockerfile new file mode 100644 index 0000000..4c4ccfd --- /dev/null +++ b/app/Dockerfile @@ -0,0 +1,29 @@ +FROM python:3.11-slim + +# Set environment variables +ENV PYTHONDONTWRITEBYTECODE=1 \ + PYTHONUNBUFFERED=1 \ + POETRY_VERSION=1.7.1 + +# Set working directory +WORKDIR /app + +# Install OS dependencies +RUN apt-get update && apt-get install -y --no-install-recommends \ + build-essential \ + libmariadb-dev \ + curl \ + git \ + && rm -rf /var/lib/apt/lists/* + +# Copy requirements and install +COPY requirements.txt . + +RUN pip install --upgrade pip \ + && pip install -r requirements.txt + +# Copy project files +COPY . . + +# Default entrypoint (can be overridden in docker-compose) +CMD ["bash"] diff --git a/app/__init__.py b/app/__init__.py index 2fd4454..fa01817 100644 --- a/app/__init__.py +++ b/app/__init__.py @@ -18,6 +18,7 @@ from flask_migrate import Migrate from celery import Celery, Task from logger import logger from app.utils.standard_responses import api_response +from flask_cors import CORS # --------------------------------------------------------------------------- # # 0. Environment setup # @@ -29,6 +30,8 @@ logger.info("Environment variables loaded from .env") # 1. Flask application & database # # --------------------------------------------------------------------------- # 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", diff --git a/app/celery_app.py b/app/celery_app.py index 0bf132a..db57e92 100644 --- a/app/celery_app.py +++ b/app/celery_app.py @@ -10,13 +10,13 @@ If you want per-project beat schedules, append them to `celery_app.conf` after import. """ -from logging import getLogger +from logger import logger from app import celery_app as celery # ← single source-of-truth from app import app as flask_app # access to config & logger -log = getLogger("xcloudify.celery_factory") -log.info("Using existing Celery instance from app.__init__") +# log = getLogger("xcloudify.celery_factory") +logger.info("Using existing Celery instance from app.__init__") # ───────────────────────── Optional: add beat jobs ──────────────────── celery.conf.beat_schedule.update( diff --git a/app/controller/api/workload_container_routes.py b/app/controller/api/workload_container_routes.py index 05ec607..f60b5bf 100644 --- a/app/controller/api/workload_container_routes.py +++ b/app/controller/api/workload_container_routes.py @@ -351,6 +351,7 @@ def get_pods(): "container_name": pf.workload.name if pf.workload else None, } for pf in pod.port_forwardings + if pf.deleted != 1 ], "cloudflare_tunnel": { "tunnel_id": tunnel.tunnel_id if tunnel else None, @@ -415,146 +416,97 @@ 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. + Mark a single container for deletion via pod-update. """ try: - workload_uuid = workload_id # Wrap with uuid.UUID(workload_id) if needed + workload_uuid = workload_id 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" - ), + 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) + mapping = ContainerPodContainer.query.filter_by(container_workload_id=_container.id).first() + if not mapping or not mapping.pod: + logger.error(f"Container {_container.id} is not part of a pod.") + return jsonify({"error": "Container is not part of a pod"}), 400 - if websocket_server_response.status_code == 201: - logger.info(f"Task created successfully for container {_container.id}") - _container.set_status("pending-deleted") + send_pod_deletion_update(mapping.pod, [str(_container.id)], delete_pod=False) - # 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 + logger.info(f"Marked container {_container.id} as pending-deleted") + return jsonify({'message': 'Container marked for deletion and pod update sent'}), 200 @api_bp.route('/workloads/pods/', 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 + Delete an entire pod and all its containers via pod-update. """ try: - pod_uuid = (pod_id) + 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 + # Delete 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) + container_ids = [str(mapping.container.id) for mapping in pod.container_mappings] + send_pod_deletion_update(pod, container_ids, delete_pod=True) - # 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 - if not container.deleted: - 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 + + +from app.utils.create_workload_container import build_pod_payload + +def send_pod_deletion_update(pod: ContainerPod, containers_to_delete: list[str], delete_pod: bool = False): + """ + Use the canonical build_pod_payload function to construct the full pod intent, + and override desired_state fields as needed for deletions. + """ + pod_payload = build_pod_payload(pod) + + # Update containers' desired_state to "deleted" where needed + for container in pod_payload["job_details"]["containers"]: + cid = container["container_id"] + if cid in containers_to_delete: + container["desired_state"] = "deleted" + workload = Workload.query.get(cid) + workload.set_status("pending-deleted") + db.session.add(workload) + + # Optionally mark NSController and pod for deletion + if delete_pod: + pod_payload["job_details"]["nscontroller"]["desired_state"] = "deleted" + pod.status = "pending-deleted" + db.session.add(pod) + + ns_id = pod_payload["job_details"]["nscontroller"]["container_id"] + ns = Workload.query.get(ns_id) + if ns: + ns.set_status("pending-deleted") + db.session.add(ns) + + db.session.commit() + + # Send the intent to the worker via WebSocket + logger.info(f"Sending pod-update to WebSocket server: {pod_payload}") + try: + response = requests.post( + app.config["WEBSOCKET_SERVER_URL"], + data=json.dumps(pod_payload), + headers={"Content-Type": "application/json"} + ) + if response.status_code == 201: + logger.info(f"WebSocket task accepted for pod {pod.id}") + else: + logger.error(f"WebSocket task rejected: {response.status_code} - {response.text}") + except Exception as e: + logger.error(f"Exception sending pod-update for pod {pod.id}: {str(e)}") diff --git a/app/controller/api/workload_host_routes.py b/app/controller/api/workload_host_routes.py index 40016bb..bd9f97f 100644 --- a/app/controller/api/workload_host_routes.py +++ b/app/controller/api/workload_host_routes.py @@ -39,8 +39,6 @@ def edit_workload_host(workload_host_id): @api_bp.route('/workload_hosts/', methods=['GET']) def get_workload_host(workload_host_id): workload_host = WorkloadHost.query.filter_by(id=workload_host_id,deleted=0).first_or_404() - - logger.debug(workload_host.to_json()) return jsonify(workload_host.to_json()) @api_bp.route('/workload_hosts/', methods=['DELETE']) diff --git a/app/models/models.py b/app/models/models.py index 9e076cd..8367133 100644 --- a/app/models/models.py +++ b/app/models/models.py @@ -689,10 +689,10 @@ class ContainerPod(BaseModel): ) ) containers = q.all() - logger.debug( - "Pod %s → %d containers (include_deleted=%s)", - self.id, len(containers), include_deleted, - ) + # logger.debug( + # "Pod %s → %d containers (include_deleted=%s)", + # self.id, len(containers), include_deleted, + # ) return containers @classmethod @@ -706,20 +706,19 @@ class ContainerPod(BaseModel): pods = db.session.query(cls).filter_by(deleted=False).all() results = [] - logger.debug( - "Scanning %d pods for active containers (include_deleted=%s)", - len(pods), include_deleted_containers, - ) + # logger.debug( + # "Scanning %d pods for active containers (include_deleted=%s)", + # len(pods), include_deleted_containers, + # ) for pod in pods: containers = pod.active_containers(include_deleted_containers) if containers or include_deleted_containers: results.append({"pod": pod, "containers": containers}) - logger.debug("Returning %d pods after filtering", len(results)) + # logger.debug("Returning %d pods after filtering", len(results)) return results - class Image(BaseModel): __tablename__ = "images" location = Column(String(255), nullable=False) diff --git a/app/requirements.txt b/app/requirements.txt index 1c473eb..e54756f 100644 --- a/app/requirements.txt +++ b/app/requirements.txt @@ -11,4 +11,7 @@ flask_migrate websocket-client streamlit aiohttp -celery \ No newline at end of file +celery +flask_cors +python-dotenv +PyJWT \ No newline at end of file diff --git a/app/tasks/reconcile_online_workers.py b/app/tasks/reconcile_online_workers.py index 763374e..97c1e82 100644 --- a/app/tasks/reconcile_online_workers.py +++ b/app/tasks/reconcile_online_workers.py @@ -16,76 +16,92 @@ from app import celery_app as celery, db, logger, app from app.models.models import WorkloadHost redisclient = redis.from_url(app.config["REDIS_URL"]) +enhanced_debug=False @celery.task(name="tasks.reconcile_online_workers", bind=True) def reconcile_online_workers(self) -> None: """ Reconcile online/offline status of WorkloadHost entries based on Redis liveness heartbeats. """ - - logger.info("Starting worker heartbeat reconciliation task") + if enhanced_debug: logger.info("=== Starting reconcile_online_workers task ===") now = int(time.time()) cutoff = datetime.utcnow() - timedelta(seconds=10) + timeout = app.config["PING_HEARTBEAT_TIMEOUT_SECONDS"] + logger.debug(f"Using PING_HEARTBEAT_TIMEOUT_SECONDS = {timeout}") + logger.debug(f"Cutoff for online host updated_at = {cutoff.isoformat()}") + # ----------------------------------------------------------------------- - # 1. Check all workers currently marked as ONLINE + # 1. Evaluate workers marked as ONLINE # ----------------------------------------------------------------------- online_workers = WorkloadHost.query.filter( WorkloadHost._status == "online", WorkloadHost.deleted == 0, WorkloadHost.updated_at < cutoff ).all() - - logger.info(f"Found {len(online_workers)} workers marked as online") + if enhanced_debug: logger.info(f"[PHASE 1] Checking {len(online_workers)} workers marked as ONLINE") for worker in online_workers: redis_key = f"ws_liveness:{worker.id}" - last_seen = redisclient.get(redis_key) + logger.debug(f"[ONLINE → OFFLINE?] Checking worker {worker.id}") + logger.debug(f" Redis key: {redis_key}") - is_stale = False - if not last_seen: - logger.warning(f"[OFFLINE CHECK] No heartbeat for WorkloadHost {worker.id}") + last_seen = redisclient.get(redis_key) + if last_seen is None: + logger.warning(f" ❌ No Redis key found — assuming stale") is_stale = True else: try: last_seen = int(last_seen) - if (now - last_seen) > app.config["PING_HEARTBEAT_TIMEOUT_SECONDS"]: - is_stale = True - except Exception: - logger.exception(f"[OFFLINE CHECK] Invalid heartbeat data for WorkloadHost {worker.id}") + delta = now - last_seen + logger.debug(f" ✅ Last heartbeat: {last_seen} (delta={delta}s)") + + is_stale = delta > timeout + if is_stale: + logger.warning(f" ⚠️ Stale heartbeat (>{timeout}s) — marking offline") + else: + if enhanced_debug: logger.debug(f" 🟢 Recent heartbeat — staying online") + except Exception as e: + logger.exception(f" ❌ Exception parsing Redis heartbeat value: {e}") is_stale = True if is_stale: - logger.warning(f"WorkloadHost {worker.id} is stale — marking offline") worker.set_status("offline") db.session.commit() + logger.info(f" ✅ Worker {worker.id} marked OFFLINE") # ----------------------------------------------------------------------- - # 2. Check all workers currently marked as OFFLINE + # 2. Evaluate workers marked as OFFLINE # ----------------------------------------------------------------------- offline_workers = WorkloadHost.query.filter( WorkloadHost._status == "offline", WorkloadHost.deleted == 0 ).all() - - logger.info(f"Found {len(offline_workers)} workers marked as offline") + if enhanced_debug: logger.info(f"[PHASE 2] Checking {len(offline_workers)} workers marked as OFFLINE") for worker in offline_workers: redis_key = f"ws_liveness:{worker.id}" + logger.debug(f"[OFFLINE → ONLINE?] Checking worker {worker.id}") + logger.debug(f" Redis key: {redis_key}") + last_seen = redisclient.get(redis_key) + if last_seen is None: + logger.debug(f" ❌ No Redis key found — remaining offline") + continue - is_active = False - if last_seen: - try: - last_seen = int(last_seen) - if (now - last_seen) <= app.config["PING_HEARTBEAT_TIMEOUT_SECONDS"]: - is_active = True - except Exception: - logger.exception(f"[ONLINE CHECK] Invalid heartbeat data for WorkloadHost {worker.id}") + try: + last_seen = int(last_seen) + delta = now - last_seen + logger.debug(f" ✅ Last heartbeat: {last_seen} (delta={delta}s)") - if is_active: - logger.info(f"WorkloadHost {worker.id} is alive again — marking online") - worker.set_status("online") - db.session.commit() + if delta <= timeout: + if enhanced_debug: logger.info(f" 🔁 Worker is back alive — marking ONLINE") + worker.set_status("online") + db.session.commit() + logger.info(f" ✅ Worker {worker.id} marked ONLINE") + else: + logger.debug(f" ⚠️ Heartbeat exists but is stale — staying offline") + except Exception as e: + logger.exception(f" ❌ Exception parsing Redis heartbeat value: {e}") - logger.info("Finished heartbeat reconciliation task.") + if enhanced_debug: logger.info("=== Finished reconcile_online_workers task ===") diff --git a/app/templates/container_logs.html b/app/templates/container_logs.html index 74846a6..534df82 100644 --- a/app/templates/container_logs.html +++ b/app/templates/container_logs.html @@ -54,7 +54,7 @@