From b282d4cf891dddd6b8b353eaeff2083b925d200a Mon Sep 17 00:00:00 2001 From: Cory Hawklvelt Date: Mon, 30 Jun 2025 14:36:39 +0930 Subject: [PATCH] handle container launch failure reporting --- websocket_server/events/docker.py | 107 +++++++-------------- worker/workerClient.py | 24 ++++- worker/worker_tasks/container.py | 148 +++++++++++++++++------------- 3 files changed, 140 insertions(+), 139 deletions(-) diff --git a/websocket_server/events/docker.py b/websocket_server/events/docker.py index 8e081ed..01ce4b2 100644 --- a/websocket_server/events/docker.py +++ b/websocket_server/events/docker.py @@ -18,26 +18,22 @@ def register_socketio_handlers(socketio): @socketio.on("docker_event") def handle_docker_event(data): - """ - Handle Docker events received from the worker. - Extract relevant data and log it to the console. - """ + """Process docker-related status events arriving from workers.""" try: - # Extract relevant data from the event worker_id = data.get("worker_id") - event_type = data.get("type") + event_type = data.get("type") # e.g. docker_start, docker_die, docker_launch_failed details = data.get("details", {}) container_id = details.get("container_id") system_container_id = details.get("system_container_id") container_name = details.get("container_name") unix_timestamp = details.get("timestamp") - timestamp=datetime.utcfromtimestamp(unix_timestamp).isoformat() + timestamp = datetime.utcfromtimestamp(unix_timestamp).isoformat() status = details.get("details", {}).get("status") image = details.get("details", {}).get("from") action = details.get("details", {}).get("Action") - # Log the extracted data + # ----------------- logging boilerplate ----------------------- # logger.info(f"Received Docker event from worker {worker_id}:") logger.info(f" Event Type: {event_type}") logger.info(f" Docker Container ID: {container_id}") @@ -47,73 +43,36 @@ def register_socketio_handlers(socketio): logger.info(f" Status: {status}") logger.info(f" Image: {image}") logger.info(f" Action: {action}") - - # Optionally, log the full details for debugging purposes logger.debug(f"Full event details: {data}") - if event_type=="docker_start": - #Send the update to the API Server - # system_container_id - # worker_id - # state="started" - # timestamp - action_URL=f"workloads/containers/status_update/{system_container_id}" - logger.debug(f"Sending task update payload to API server for Container ID {system_container_id}") + # ------------------- status-update calls --------------------- # + if event_type in ( + "docker_start", + "docker_destroy", + "docker_die", + "docker_launch_failed", # ← NEW + ): + action_URL = f"workloads/containers/status_update/{system_container_id}" + new_status_map = { + "docker_start": "running", + "docker_destroy": "deleted", + "docker_die": "dead", + "docker_launch_failed": "launch_failed", # NEW + } + payload = { + "new_status": new_status_map[event_type], + "system_container_id": system_container_id, + "worker_id": worker_id, + "timestamp": timestamp, + } + headers = {"Content-Type": "application/json"} + websocket_server_response = requests.put( + f"{api_server_url}/{action_URL}", + data=json.dumps(payload), + headers=headers, + ) + logger.debug(websocket_server_response) + logger.info("Update complete") - payload = { - "new_status": "running", - "system_container_id": system_container_id, - "worker_id": worker_id, - "timestamp": timestamp - } - - headers = {"Content-Type": "application/json"} - websocket_server_response = requests.put(f"{api_server_url}/{action_URL}", data=json.dumps(payload), headers=headers) - logger.debug( websocket_server_response) - logger.info("Update complete") - - if event_type=="docker_destroy": - #Send the update to the API Server - # system_container_id - # worker_id - # state="started" - # timestamp - action_URL=f"workloads/containers/status_update/{system_container_id}" - logger.debug(f"Sending task update payload to API server for Container ID {system_container_id}") - - payload = { - "new_status": "deleted", - "system_container_id": system_container_id, - "worker_id": worker_id, - "timestamp": timestamp - } - - headers = {"Content-Type": "application/json"} - websocket_server_response = requests.put(f"{api_server_url}/{action_URL}", data=json.dumps(payload), headers=headers) - logger.debug( websocket_server_response) - logger.info("Update complete") - - if event_type=="docker_die": - #Send the update to the API Server - # system_container_id - # worker_id - # state="started" - # timestamp - action_URL=f"workloads/containers/status_update/{system_container_id}" - logger.debug(f"Sending task update payload to API server for Container ID {system_container_id}") - - payload = { - "new_status": "dead", - "system_container_id": system_container_id, - "worker_id": worker_id, - "timestamp": timestamp - } - - headers = {"Content-Type": "application/json"} - websocket_server_response = requests.put(f"{api_server_url}/{action_URL}", data=json.dumps(payload), headers=headers) - logger.debug( websocket_server_response) - logger.info("Update complete") - except Exception as e: - logger.error(f"Error processing Docker event: {e}") - + logger.error(f"Error processing Docker event: {e}") \ No newline at end of file diff --git a/worker/workerClient.py b/worker/workerClient.py index 20887e2..71b4593 100644 --- a/worker/workerClient.py +++ b/worker/workerClient.py @@ -14,7 +14,7 @@ import json import asyncio import socketio import threading - +import time class WorkerClient: def __init__(self, event_queue=None, docker_monitor=None): """Initializes the WorkerClient using global settings.""" @@ -104,7 +104,7 @@ class WorkerClient: if not self.joined_server: await self.sio.emit("join_request", {"worker_id": self.worker_id, "worker_secret": self.worker_secret}) logger.info(f"Worker {self.worker_id} asked to join the server.") - + async def handle_task(self, data): """Handle incoming tasks from the server.""" if not self.joined_server: @@ -126,6 +126,7 @@ class WorkerClient: result = None try: + # ----------------------- Task dispatch ------------------------ # if task_type == "report": result = await loop.run_in_executor(None, functools.partial(ReportTask("", logger).Execute)) elif task_type == "ping": @@ -145,13 +146,30 @@ class WorkerClient: else: raise ValueError(f"Unknown task type: {task_type}") + # --------------------- ACK back to server --------------------- # await self.send_task_result(task_id, result, task_worker_id) + # ---------------- Emit launch-failure events ------------------ # + if task_type == "pod-update" and result and result.get("launch_failures"): + for failed_id in result["launch_failures"]: + await self.send_event({ + "source": "docker", + "event_type": "launch_failed", + "container_id": failed_id, + "system_container_id": failed_id, + "container_name": failed_id, + "timestamp": time.time(), + "details": { + "status": "launch_failed", + "from": "", + "Action": "launch_failed" + } + }) + except Exception as e: logger.error(f"Error processing task {task_id}: {e}") await self.send_task_result(task_id, {"success": False, "response": str(e)}, task_worker_id) - async def send_task_result(self, task_id, result, worker_id): """Send task result back to the server.""" await self.sio.emit("ack", {"task_id": task_id, "worker_id": worker_id, "result": result}) diff --git a/worker/worker_tasks/container.py b/worker/worker_tasks/container.py index 294cf27..e344cb4 100644 --- a/worker/worker_tasks/container.py +++ b/worker/worker_tasks/container.py @@ -4,15 +4,39 @@ import pty import os import subprocess import select +from typing import List + class ContainerTask: + """ + Helper class that reconciles desired container state with the local Docker + daemon and emits events back to the Worker for upstream processing. + """ + def __init__(self, logger): self.docker_client = docker.from_env() - self.logger=logger + self.logger = logger + # Track container IDs that fail to launch so the caller can notify the server + self.launch_failures: List[str] = [] + # --------------------------------------------------------------------- # + # Orchestrator-facing entry-point + # --------------------------------------------------------------------- # def handle_pod_update(self, pod_payload): - """Handle a pod-update task.""" + """ + Handle a pod-update task. + + Returns + ------- + dict + { + "success": bool, + "response": {}, + "launch_failures": [, ...] # when any fail + } + """ response_payload = {"success": False, "response": {}} + self.launch_failures = [] # reset in case the instance is re-used self.logger.info(f"Handling pod update for Pod {pod_payload['pod_id']}") @@ -22,16 +46,18 @@ class ContainerTask: # Always ensure the NSController first self.ensure_container(nscontroller_spec, is_nscontroller=True) - # Then process all other containers for container_spec in container_specs: - container_spec['nscontroller_container_id'] = nscontroller_spec['container_id'] + container_spec["nscontroller_container_id"] = nscontroller_spec["container_id"] self.ensure_container(container_spec) response_payload["success"] = True - + response_payload["launch_failures"] = self.launch_failures return response_payload + # --------------------------------------------------------------------- # + # Internal helpers + # --------------------------------------------------------------------- # def ensure_container(self, container_spec, is_nscontroller=False): """Ensure a container matches the desired state (create, update, or delete).""" container_id = container_spec["container_id"] @@ -50,13 +76,67 @@ class ContainerTask: return self.launch_container(container_spec, is_nscontroller) if not self.is_container_correct(existing, container_spec): - self.logger.warning(f"Config drift detected for container {container_id}, recreating...") + self.logger.warning(f"Config drift detected for container {container_id}, recreating…") self.delete_container(existing) return self.launch_container(container_spec, is_nscontroller) self.logger.info(f"Container {container_id} already running and correct.") return existing + def launch_container(self, container_spec, is_nscontroller=False): + """Launch a new container from spec. Returns the Container or None on failure.""" + docker_image = container_spec.get("docker_image") + container_name = container_spec.get("container_name", f"container_{uuid.uuid4()}") + ports = None + networks = container_spec.get("networks", []) + system_container_id = container_spec.get("container_id") + + # Ports (only for NSControllers, normal containers use shared namespace) + if is_nscontroller: + ports = {f"{p['internal']}/tcp": p["external"] for p in container_spec.get("ports", [])} + + # Networking mode + network_mode = None + if not is_nscontroller and container_spec.get("nscontroller_container_id"): + nscontroller = self.find_container(container_spec["nscontroller_container_id"]) + if nscontroller: + network_mode = f"container:{nscontroller.id}" + + labels = { + "managed_by": "worker_agent", + "system_container_id": system_container_id, + } + + container_kwargs = { + "image": docker_image, + "name": system_container_id, # keep docker name == system ID + "detach": True, + "labels": labels, + "network_mode": network_mode, + } + if ports: + container_kwargs["ports"] = ports + if "command" in container_spec: + container_kwargs["command"] = container_spec["command"] + + try: + # Step 1 – create & start + container = self.docker_client.containers.run(**container_kwargs) + self.logger.info(f"Launched container {container_name} from image {docker_image}") + + # Step 2 – attach to extra networks if needed + if is_nscontroller or (not is_nscontroller and network_mode is None): + for network_name in networks: + self.attach_to_network(container, network_name) + + return container + + except Exception as e: + self.logger.error(f"Failed to launch container {container_name}: {str(e)}") + # record failure so WorkerClient can emit docker_launch_failed + self.launch_failures.append(system_container_id) + return None + def find_container(self, system_container_id): """Find a container by its system_container_id label.""" for container in self.docker_client.containers.list(all=True): @@ -78,62 +158,6 @@ class ContainerTask: except Exception as e: self.logger.warning(f"Failed to remove container {container.name}: {str(e)}") - def launch_container(self, container_spec, is_nscontroller=False): - """Launch a new container from spec.""" - docker_image = container_spec.get("docker_image") - container_name = container_spec.get("container_name", f"container_{uuid.uuid4()}") - ports = None - networks = container_spec.get("networks", []) - system_container_id = container_spec.get("container_id") - - # Ports (only for NSControllers, normal containers use shared namespace) - if is_nscontroller: - ports = {} - for port in container_spec.get("ports", []): - ports[f"{port['internal']}/tcp"] = port['external'] - - # Networking mode - network_mode = None - if not is_nscontroller and container_spec.get("nscontroller_container_id"): - nscontroller = self.find_container(container_spec["nscontroller_container_id"]) - if nscontroller: - network_mode = f"container:{nscontroller.id}" - - labels = { - "managed_by": "worker_agent", - "system_container_id": system_container_id - } - - # Build initial container args - container_kwargs = { - "image": docker_image, - "name": system_container_id, - "detach": True, - "labels": labels, - "network_mode": network_mode, - } - - if ports: - container_kwargs["ports"] = ports - - if "command" in container_spec: - container_kwargs["command"] = container_spec["command"] - - try: - # Step 1: Create container - container = self.docker_client.containers.run(**container_kwargs) - self.logger.info(f"Launched container {container_name} from image {docker_image}") - - # Step 2: Attach to extra networks if needed (only for NSController or normal if network_mode None) - if is_nscontroller or (not is_nscontroller and network_mode is None): - for network_name in networks: - self.attach_to_network(container, network_name) - - return container - except Exception as e: - self.logger.error(f"Failed to launch container {container_name}: {str(e)}") - return None - def is_container_correct(self, container, container_spec): """Check if a running container matches its spec (image, networks, ports, env vars, volumes).""" try: