handle container launch failure reporting
This commit is contained in:
@@ -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}")
|
||||
+21
-3
@@ -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})
|
||||
|
||||
@@ -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": [<system_container_id>, ...] # 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:
|
||||
|
||||
Reference in New Issue
Block a user