Files
3cloud-backend/worker/worker_tasks/container.py
T

1307 lines
57 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.
import time
import traceback
import docker
import uuid
import pty
import os
import subprocess
import select
import threading
import base64
from typing import List, Dict, Any, Optional, Callable
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, docker_monitor=None):
self.docker_client = docker.from_env()
self.logger = logger
# Track container IDs that fail to launch so the caller can notify the server
self.launch_failures: list[dict[str, str]] = [] # ← CHANGED
# Docker monitor reference will be set in handle_pod_update
self.docker_monitor = docker_monitor
# --------------------------------------------------------------------- #
# Orchestrator-facing entry-point
# --------------------------------------------------------------------- #
def handle_pod_update_with_reconciliation(self, pod_payload):
"""
Handle a pod-update task with blacklist-based noise reduction and status reconciliation.
Returns
-------
dict
{
"success": bool,
"response": {},
"launch_failures": [<system_container_id>, ...] # when any fail
}
"""
self.logger.info(f"Starting reconciliation-based pod update for Pod {pod_payload['pod_id']}")
# 1. Capture initial container statuses
initial_statuses = self.capture_container_statuses(pod_payload)
# 2. Blacklist all containers to prevent event processing during update
self.blacklist_containers(pod_payload)
try:
# 3. Process pod update using existing logic
result = self._process_pod_update(pod_payload)
# 4. Unblacklist containers
self.unblacklist_containers(pod_payload)
# 5. Capture final statuses
final_statuses = self.capture_container_statuses(pod_payload)
# 6. Compare statuses and send delta updates
status_changes = self.compare_statuses(initial_statuses, final_statuses)
if status_changes:
self.send_delta_updates(status_changes)
self.logger.info(f"Reconciliation-based pod update completed for Pod {pod_payload['pod_id']}")
return result
except Exception as e:
# Ensure containers are unblacklisted even if update fails
self.logger.error(f"Pod update failed, ensuring containers are unblacklisted: {str(e)}")
self.logger.error(traceback.format_exc())
self.unblacklist_containers(pod_payload)
raise
def handle_pod_update(self, pod_payload):
"""
Handle a pod-update task. This method now delegates to the reconciliation-based approach.
Returns
-------
dict
{
"success": bool,
"response": {},
"launch_failures": [<system_container_id>, ...] # when any fail
}
"""
# Use the new reconciliation-based approach
return self.handle_pod_update_with_reconciliation(pod_payload)
def _process_pod_update(self, pod_payload):
"""
Internal method that processes the pod update using the original logic.
This is the core pod update logic extracted from the original handle_pod_update method.
"""
response_payload = {"success": False, "response": {}}
self.launch_failures = [] # reset in case the instance is re-used
# Check for lifecycle operations in the payload
lifecycle_ops = []
for container in pod_payload.get("containers", []):
container_id = container.get("container_id")
if container.get("restart"):
lifecycle_ops.append(f"restart container {container_id}")
elif container.get("desired_state") == "stopped":
lifecycle_ops.append(f"stop container {container_id}")
elif container.get("desired_state") == "running":
lifecycle_ops.append(f"ensure container {container_id} is running")
if lifecycle_ops:
self.logger.info(f"Processing pod update for Pod {pod_payload['pod_id']} with lifecycle operations: {', '.join(lifecycle_ops)}")
else:
self.logger.info(f"Processing pod update for Pod {pod_payload['pod_id']}")
# Reset launch failures for this operation
self.launch_failures = []
container_specs = pod_payload["containers"]
# Process any storage volumes in the payload
for container in container_specs:
if 'storage' in container and container['storage']:
self._process_storage_volumes(container['storage'])
nscontroller_specs: list[dict] = []
regular_container_specs: list[dict] = []
for container_spec in container_specs:
if container_spec.get("workload_type") == "NSController":
nscontroller_specs.append(container_spec)
else:
regular_container_specs.append(container_spec)
nscontroller_container_id: str | None = None
if len(nscontroller_specs) == 0:
self.logger.info("No NSController containers found in this payload.")
elif len(nscontroller_specs) == 1:
nscontroller_container_id = nscontroller_specs[0].get("container_id")
self.logger.info(
"Identified single NSController container_id=%s to initialize networking.",
nscontroller_container_id,
)
else:
# Multiple NSControllers present—log a warning and pick the first deterministically.
nscontroller_container_id = nscontroller_specs[0].get("container_id")
self.logger.warning(
"Multiple NSController containers detected (%d). Proceeding with the first container_id=%s.",
len(nscontroller_specs),
nscontroller_container_id,
)
# ---- launch NSController(s) first (no injection of nscontroller_container_id)
for ns_container_spec in nscontroller_specs:
self.logger.debug(
"Launching NSController container_id=%s",
ns_container_spec.get("container_id"),
)
self.ensure_container(ns_container_spec, is_nscontroller=True)
# ---- launch regular containers with networking pointer if available
for regular_container_spec in regular_container_specs:
if nscontroller_container_id:
regular_container_spec["nscontroller_container_id"] = nscontroller_container_id
self.logger.debug(
"Injecting nscontroller_container_id=%s into container_id=%s",
nscontroller_container_id,
regular_container_spec.get("container_id"),
)
self.logger.debug(
"Launching regular container_id=%s",
regular_container_spec.get("container_id"),
)
self.ensure_container(regular_container_spec)
# Set success to False if there were any launch failures
response_payload["success"] = len(self.launch_failures) == 0
response_payload["launch_failures"] = self.launch_failures
if self.launch_failures:
self.logger.warning(f"Pod update completed with {len(self.launch_failures)} failures")
for failure in self.launch_failures:
self.logger.warning(f" - Container {failure['id']}: {failure['error']}")
else:
self.logger.info(f"Pod update completed successfully")
return response_payload
def shutdown(self):
"""Gracefully shutdown Docker connection."""
try:
self.docker_client.close()
except Exception:
pass
# --------------------------------------------------------------------- #
# 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"]
desired_state = container_spec.get("desired_state", "running")
restart_requested = container_spec.get("restart", False)
existing = self.find_container(container_id)
# Handle deletion
if desired_state == "deleted":
if existing:
self.logger.info(f"Deleting container {container_id}")
self.delete_container(existing)
return None
# Handle restart request
if restart_requested and existing:
self.logger.info(f"Restarting container {container_id} as requested")
success = self.restart_container(existing)
if not success:
self.logger.warning(f"Failed to restart container {container_id}")
return existing
# Handle stopped state
if desired_state == "stopped" and existing:
if existing.status != "exited":
self.logger.info(f"Stopping container {container_id} as requested")
success = self.stop_container(existing)
if not success:
self.logger.warning(f"Failed to stop container {container_id}")
else:
self.logger.info(f"Container {container_id} already stopped")
return existing
# Handle start request for stopped container
if desired_state == "running" and existing and existing.status == "exited":
self.logger.info(f"Starting stopped container {container_id}")
success = self.start_container(existing)
if not success:
self.logger.warning(f"Failed to start container {container_id}")
# If we couldn't start it, try to recreate it
self.logger.info(f"Attempting to recreate container {container_id}")
try:
# First try to remove the invalid container
try:
self.delete_container(existing)
except Exception as e:
self.logger.warning(f"Error cleaning up invalid container {container_id}: {str(e)}")
# Then create a new one
return self.launch_container(container_spec, is_nscontroller=(container_spec.get("nscontroller_container_id") is None))
except Exception as e:
self.logger.error(f"Failed to recreate container {container_id}: {str(e)}")
return existing
# Handle non-existent container
if not existing:
self.logger.info(f"Launching new container {container_id}")
return self.launch_container(container_spec, is_nscontroller)
# Handle configuration drift
if not self.is_container_correct(existing, container_spec):
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 pull_image_with_progress(self, image_name: str, container_id: str, enable_progress_tracking: bool = False) -> Dict[str, Any]:
"""
Pull a Docker image with optional progress tracking.
Args:
image_name: The Docker image name to pull
container_id: The system container ID for tracking
enable_progress_tracking: Whether to send detailed progress events
Returns:
Dict with pull results including status, duration, bytes downloaded, etc.
"""
pull_start_time = time.time()
pull_result = {
"status": "idle",
"source": "unknown",
"duration_seconds": 0.0,
"bytes_downloaded": 0,
"layers_completed": 0,
"total_layers": 0,
"eta_seconds": None,
"error": None
}
try:
self.logger.info(f"Starting pull for image {image_name} (container: {container_id})")
# Check if image already exists locally
try:
self.docker_client.images.get(image_name)
pull_result["status"] = "completed"
pull_result["source"] = "cached"
pull_result["duration_seconds"] = time.time() - pull_start_time
self.logger.info(f"Image {image_name} already cached locally")
return pull_result
except docker.errors.ImageNotFound:
pass # Image not cached, proceed with pull
# Set initial status
pull_result["status"] = "pulling"
pull_result["source"] = "remote"
# Progress tracking variables
last_progress_update = time.time()
progress_tracking_enabled = enable_progress_tracking # Start with initial setting
layer_events_by_id = {} # Track latest event per layer
stop_progress = threading.Event()
def send_progress_event(progress_data: Dict[str, Any]):
"""Send progress update via WebSocket if tracking is enabled"""
if not progress_tracking_enabled:
return
try:
# Send event through docker_monitor if available
if self.docker_monitor and hasattr(self.docker_monitor, 'event_queue'):
event_data = {
'source': 'docker',
'event_type': 'pull_progress',
'container_id': container_id,
'system_container_id': container_id,
'timestamp': time.time(),
'pull_progress': progress_data
}
self.docker_monitor.event_queue.put(event_data)
self.logger.debug(f"Sent pull progress event for {container_id}: {progress_data['bytes_downloaded']} bytes, {progress_data['layers_completed']}/{progress_data['layers_total']} layers")
except Exception as e:
self.logger.warning(f"Failed to send progress event: {e}")
def compute_overall_progress() -> tuple:
"""Sum current bytes and count completed layers."""
overall_current_bytes = 0
layers_completed = 0
for latest_event in layer_events_by_id.values():
progress_detail = latest_event.get("progressDetail", {})
current_bytes = progress_detail.get("current")
total_bytes = progress_detail.get("total")
status = latest_event.get("status")
if current_bytes is not None:
overall_current_bytes += int(current_bytes)
# Count as completed if status indicates completion or current == total
if status in ["Download complete", "Pull complete"] or \
(current_bytes is not None and total_bytes is not None and current_bytes == total_bytes):
layers_completed += 1
return overall_current_bytes, layers_completed
def human_readable_bytes(byte_count) -> str:
"""Convert a byte count into a human-readable string."""
if byte_count is None:
return "?"
unit_labels = ["B", "KB", "MB", "GB", "TB"]
value = float(byte_count)
unit_index = 0
while value >= 1024.0 and unit_index < len(unit_labels) - 1:
value /= 1024.0
unit_index += 1
return f"{value:.1f} {unit_labels[unit_index]}"
def log_progress_snapshot():
"""Log a concise, rolling snapshot of pull progress."""
status_counts = {}
for latest_event in layer_events_by_id.values():
status_text = latest_event.get("status") or "Unknown"
status_counts[status_text] = status_counts.get(status_text, 0) + 1
overall_current_bytes, layers_completed = compute_overall_progress()
# Update pull result
pull_result["bytes_downloaded"] = overall_current_bytes or 0
pull_result["layers_completed"] = layers_completed
pull_result["total_layers"] = len(layer_events_by_id)
self.logger.debug(f"Pull progress: {human_readable_bytes(overall_current_bytes or 0)} downloaded - Layers: {layers_completed}/{len(layer_events_by_id)}")
# Start the pull operation
self.logger.debug(f"Pulling image {image_name} with initial progress tracking: {enable_progress_tracking}")
# Use the existing API client for progress tracking
pull_response = self.docker_client.api.pull(image_name, stream=True, decode=True)
# Process the pull stream
for stream_event in pull_response:
if stop_progress.is_set():
break
current_time = time.time()
pull_duration = current_time - pull_start_time
# Enable progress tracking if pull takes longer than 30 seconds
if not progress_tracking_enabled and pull_duration > 10.0:
self.logger.info(f"Pull taking longer than 10s ({pull_duration:.1f}s), enabling detailed progress tracking")
progress_tracking_enabled = True
# Handle errors emitted by the engine
if "error" in stream_event:
error_message = stream_event.get("error") or "Unknown error"
self.logger.error(f"Pull error: {error_message}")
raise RuntimeError(error_message)
# Track per-layer status
layer_identifier = stream_event.get("id")
if layer_identifier:
layer_events_by_id[layer_identifier] = stream_event
# Periodic snapshot and event sending
if progress_tracking_enabled and (current_time - last_progress_update) >= 5.0:
last_progress_update = current_time
log_progress_snapshot()
# Send progress event
overall_current_bytes, layers_completed = compute_overall_progress()
send_progress_event({
"bytes_downloaded": overall_current_bytes or 0,
"layers_completed": layers_completed,
"layers_total": len(layer_events_by_id),
"status": "pulling"
})
# Pull completed successfully
pull_end_time = time.time()
pull_result["status"] = "completed"
pull_result["duration_seconds"] = pull_end_time - pull_start_time
# Send final progress event if tracking was ever enabled
if progress_tracking_enabled:
send_progress_event({
"bytes_downloaded": pull_result["bytes_downloaded"],
"layers_completed": pull_result["layers_completed"],
"layers_total": pull_result["total_layers"],
"status": "completed"
})
self.logger.info(f"Successfully pulled image {image_name} in {pull_result['duration_seconds']:.2f}s")
except Exception as e:
pull_end_time = time.time()
pull_result["status"] = "failed"
pull_result["duration_seconds"] = pull_end_time - pull_start_time
pull_result["error"] = str(e)
# Send failure event
if enable_progress_tracking:
send_progress_event({
"bytes_downloaded": pull_result["bytes_downloaded"],
"layers_completed": pull_result["layers_completed"],
"layers_total": pull_result["total_layers"],
"status": "failed",
"error": str(e)
})
self.logger.error(f"Failed to pull image {image_name}: {e}")
finally:
# Clean up progress thread if it exists
if 'stop_progress' in locals():
stop_progress.set()
return pull_result
def launch_container(self, container_spec, is_nscontroller=False):
"""Launch a new container from spec. Returns the Container or None on failure."""
self.logger.debug("Launching container")
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")
# Handle image pulling with progress tracking
if docker_image:
# Normalize image name
if ":" not in docker_image:
docker_image = f"{docker_image}:latest"
try:
# Start pull with progress tracking - the method will handle the 30-second logic internally
self.logger.debug(f"Starting image pull for {docker_image}")
pull_result = self.pull_image_with_progress(
docker_image,
system_container_id,
enable_progress_tracking=False # Start without tracking, enable after 30s if needed
)
# Handle pull results
if pull_result["status"] == "failed":
self.logger.error(f"Failed to pull image {docker_image}: {pull_result.get('error')}")
self.launch_failures.append({
"id": system_container_id,
"error": f"Image pull failed: {pull_result.get('error')}"
})
return None
self.logger.info(f"Image pull completed: {docker_image} ({pull_result['source']}) in {pull_result['duration_seconds']:.2f}s")
except Exception as e:
self.logger.error(f"Error during image pull for {docker_image}: {e}")
self.launch_failures.append({
"id": system_container_id,
"error": f"Image pull error: {str(e)}"
})
return None
# Ports (only for NSControllers, normal containers use shared namespace)
if is_nscontroller:
self.logger.debug("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,
}
# Note: Image is already pulled above, so containers.run() will use the local image
if ports:
container_kwargs["ports"] = ports
if "command" in container_spec:
container_kwargs["command"] = container_spec["command"]
if "env" in container_spec:
container_kwargs["environment"] = container_spec["env"]
# Add CPU and memory limits if specified
if "cpu_shares" in container_spec:
# Docker expects cpu_shares between 2-262144, with 1024 being default (1 CPU)
# Convert our simple value to Docker's expected range
cpu_value = int(container_spec["cpu_shares"])
# Scale to Docker's range - 1 CPU share = 1024 in Docker's scale
docker_cpu_shares = max(2, min(262144, cpu_value * 1024))
container_kwargs["cpu_shares"] = docker_cpu_shares
self.logger.debug(f"Setting CPU shares to {docker_cpu_shares} (from {cpu_value})")
if "mem_limit" in container_spec:
# Convert to string with 'm' suffix for Docker API
container_kwargs["mem_limit"] = f"{container_spec['mem_limit']}m"
self.logger.debug("pre-storage")
# Handle storage volumes if specified
if "storage" in container_spec:
if "volumes" not in container_kwargs:
container_kwargs["volumes"] = {}
for storage_entry in container_spec["storage"]:
# Gather and validate common fields
volume_type: str = storage_entry.get("volume_type", "").strip().lower()
container_mount_point: str | None = storage_entry.get("mount_point")
if volume_type not in {"nfs", "local"}:
self.logger.warning(
"Skipping storage entry due to unsupported volume_type",
extra={"volume_type": volume_type, "entry": storage_entry}
)
continue
if not container_mount_point or not isinstance(container_mount_point, str):
self.logger.error(
"Storage entry missing or invalid 'mount_point'; skipping",
extra={"volume_type": volume_type, "entry": storage_entry}
)
continue
# Default access mode
is_read_only: bool = bool(storage_entry.get("read_only", False))
# Compute expanded_path based on volume type
if volume_type == "nfs":
# NFS: mount location is SHAREDFS_ROOT/<volume_id>
# expanded_path is set by the _process_storage_volumes function where is maps the local NFS dir and the volume ID to a specific folder.. Which could be different on every host
expanded_path = storage_entry.get("expanded_path")
if not expanded_path:
self.logger.error(
"NFS storage requires 'expanded_path'; skipping entry",
extra={"entry": storage_entry}
)
continue
self.logger.info(
"Resolved NFS volume host path using volume_id under SHAREDFS_ROOT",
extra={
"volume_type": volume_type,
"volume_id": storage_entry.get("volume_id"),
"sharedfs_root": storage_entry.get("expanded_path"),
"container_mount_point": container_mount_point,
"read_only": is_read_only,
},
)
elif volume_type == "local":
# LOCAL: mount location is the provided volume_path directly
local_volume_path = storage_entry.get("volume_path")
if not local_volume_path or not isinstance(local_volume_path, str):
self.logger.error(
"Local storage requires 'volume_path'; skipping entry",
extra={"entry": storage_entry}
)
continue
# Normalize leading slashes
expanded_path = local_volume_path.strip()
self.logger.info(
"Resolved Local volume host path using provided volume_path",
extra={
"volume_type": volume_type,
"expanded_path": expanded_path,
"container_mount_point": container_mount_point,
"read_only": is_read_only,
},
)
# Build Docker volume binding mapping
# Note: We use the dict form (preferred by docker-py) instead of f-string format.
access_mode: str = "ro" if is_read_only else "rw"
container_kwargs["volumes"][expanded_path] = {
"bind": container_mount_point,
"mode": access_mode,
}
self.logger.debug(
"Added volume binding to container kwargs",
extra={
"bind": container_mount_point,
"mode": access_mode,
"current_volumes_count": len(container_kwargs["volumes"]),
},
)
# Handle injected files
if "injected_files" in container_spec:
self.logger.debug(f"Processing injected_files for container {system_container_id}: {container_spec['injected_files']}")
# Create temp directory for injected files
temp_dir = f"/tmp/xcloudify_injected/{system_container_id}"
os.makedirs(temp_dir, exist_ok=True)
for injected_file in container_spec["injected_files"]:
filename = injected_file["filename"]
content_b64 = injected_file["content"]
permissions = injected_file["permissions"]
# Decode base64
try:
content = base64.b64decode(content_b64)
except Exception as e:
self.logger.error(f"Failed to decode base64 content for {filename}: {e}")
continue
# Write to temp file
temp_file_path = os.path.join(temp_dir, os.path.basename(filename))
try:
with open(temp_file_path, 'wb') as f:
f.write(content)
# Set permissions
os.chmod(temp_file_path, int(permissions, 8))
self.logger.debug(f"Injected file {filename} written to {temp_file_path} with permissions {permissions}")
# Add bind mount
if "volumes" not in container_kwargs:
container_kwargs["volumes"] = {}
container_kwargs["volumes"][temp_file_path] = {
"bind": filename,
"mode": "ro" # Read-only, as injected files shouldn't be modified
}
except Exception as e:
self.logger.error(f"Failed to write injected file {filename}: {e}")
continue
else:
self.logger.debug(f"No injected_files found in container_spec for {system_container_id}")
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(
{"id": system_container_id, "error": str(e)}
)
return None
def find_container(self, system_container_id):
"""Find a container by its system_container_id label."""
try:
for container in self.docker_client.containers.list(all=True):
if container.labels.get("system_container_id") == system_container_id:
# Verify the container is still valid by checking its status
try:
# This will fail if the container is no longer valid
container.reload()
return container
except docker.errors.NotFound:
self.logger.warning(f"Container with system_id {system_container_id} exists in list but is invalid, skipping")
continue
except Exception as e:
self.logger.error(f"Error finding container {system_container_id}: {str(e)}")
return None
def stop_container(self, container):
"""Stop a Docker container without removing it."""
system_container_id = container.labels.get("system_container_id")
docker_id = container.id
try:
container.stop(timeout=10)
self.logger.info(f"Stopped container {container.name}")
return True
except Exception as e:
self.logger.error(f"Failed to stop container {container.name}: {str(e)}")
# Add to launch failures so it's reported back to the server
if system_container_id:
self.launch_failures.append(
{"id": system_container_id, "error": f"Failed to stop: {str(e)}"}
)
return False
def start_container(self, container):
"""Start a stopped Docker container."""
system_container_id = container.labels.get("system_container_id")
docker_id = container.id
self.logger.debug("Checking for docker_monitor")
try:
# First reload the container to make sure it's still valid
try:
container.reload()
except docker.errors.NotFound:
# Container exists in our DB but not in Docker
error_msg = "Container exists in database but not in Docker daemon"
self.logger.error(f"Failed to start container {container.name}: {error_msg}")
if system_container_id:
self.launch_failures.append(
{"id": system_container_id, "error": f"Failed to start: {error_msg}"}
)
return False
# Now try to start it
container.start()
self.logger.info(f"Started container {container.name}")
return True
except Exception as e:
self.logger.error(f"Failed to start container {container.name}: {str(e)}")
# Add to launch failures so it's reported back to the server
if system_container_id:
self.launch_failures.append(
{"id": system_container_id, "error": f"Failed to start: {str(e)}"}
)
return False
def restart_container(self, container):
"""Restart a Docker container."""
system_container_id = container.labels.get("system_container_id")
docker_id = container.id
try:
container.restart(timeout=10)
self.logger.info(f"Restarted container {container.name}")
return True
except Exception as e:
self.logger.error(f"Failed to restart container {container.name}: {str(e)}")
# Add to launch failures so it's reported back to the server
if system_container_id:
self.launch_failures.append(
{"id": system_container_id, "error": f"Failed to restart: {str(e)}"}
)
return False
def delete_container(self, container):
"""Stop and remove a Docker container."""
system_container_id = container.labels.get("system_container_id")
docker_id = container.id
try:
container.stop(timeout=5)
self.logger.info(f"Stopped container {container.name}")
except Exception as e:
self.logger.warning(f"Failed to stop container {container.name}: {str(e)}")
try:
container.remove(force=True)
self.logger.info(f"Removed container {container.name}")
except Exception as e:
self.logger.warning(f"Failed to remove container {container.name}: {str(e)}")
# Clean up injected files temp directory
temp_dir = f"/tmp/xcloudify_injected/{system_container_id}"
try:
if os.path.exists(temp_dir):
import shutil
shutil.rmtree(temp_dir)
self.logger.debug(f"Cleaned up injected files temp directory: {temp_dir}")
except Exception as e:
self.logger.warning(f"Failed to clean up injected files temp directory {temp_dir}: {e}")
def is_container_correct(self, container, container_spec):
self.logger.debug("Inside is_container_correct")
"""
Determine whether the running Docker container still matches the desired
launch specification.
Parameters
----------
container : docker.models.containers.Container
The live Docker container object.
container_spec : dict
The desired launch specification sent from the orchestrator.
Returns
-------
bool
True → container matches spec (no action required).
False → drift detected; caller should recreate the container.
"""
# I have removed this becuase it SHOUD be idempotent, but also it's causing NSControllers not to refresh their ports when adding a new container to an existing pod
# # If this is just a lifecycle operation (restart/stop/start), don't check for drift
# if container_spec.get("restart") or container_spec.get("desired_state") in ["stopped", "running"]:
# # Only check if the container exists with the right ID
# if container and container.labels.get("system_container_id") == container_spec.get("container_id"):
# return True
try:
config = container.attrs.get('Config', {})
host_config = container.attrs.get('HostConfig', {})
network_settings = container.attrs.get('NetworkSettings', {})
# ─────────────────── 1️⃣ IMAGE ───────────────────
# Normalize tags to include ':latest' if omitted
expected_image = container_spec.get("docker_image")
if expected_image and ":" not in expected_image:
expected_image += ":latest"
running_image = container.image.tags[0] if container.image.tags else None
if expected_image != running_image:
self.logger.warning(
"Image mismatch: expected %s, got %s", expected_image, running_image
)
return False
# ─────────────────── 2️⃣ NETWORKS (NSController only) ───────────────────
if container_spec.get("networks") and container_spec.get("nscontroller_container_id") is None:
existing_networks = set(network_settings.get('Networks', {}).keys())
expected_networks = set(container_spec.get("networks", []))
if not expected_networks.issubset(existing_networks):
self.logger.warning(
"Network mismatch: expected %s, got %s", expected_networks, existing_networks
)
return False
# ─────────────────── 3️⃣ PORTS (NSController only) ───────────────────
self.logger.debug("Checking ports")
if container_spec.get("ports") and container_spec.get("nscontroller_container_id") is None:
self.logger.debug("has ports and nscontroller_container_id is none")
bound_ports = {
(int(proto.split('/')[0]), int(mapping['HostPort']))
for proto, mappings in network_settings.get('Ports', {}).items() if mappings
for mapping in mappings
}
expected_ports = {(p['internal'], p['external']) for p in container_spec['ports']}
self.logger.debug(f"Bound ports{bound_ports}")
self.logger.debug(f"Expected ports{expected_ports}")
if not expected_ports.issubset(bound_ports):
self.logger.warning(
"Ports mismatch: expected %s, got %s", expected_ports, bound_ports
)
return False
# ─────────────────── 4️⃣ ENVIRONMENT VARIABLES ───────────────────
def _list_to_dict(env_list) -> dict[str, str]:
"""Convert ['KEY=VALUE', …] → {'KEY': 'VALUE', …} safely."""
result = {}
if not env_list:
return result
for item in env_list:
if not isinstance(item, str) or '=' not in item:
# ignore malformed entries, they will be caught as drift
continue
key, value = item.split('=', 1)
result[key.strip()] = value.strip()
return result
# normalise expected → dict
expected_env_raw = container_spec.get("env", {})
if isinstance(expected_env_raw, list):
expected_env = _list_to_dict(expected_env_raw)
elif isinstance(expected_env_raw, dict):
# ensure all values are str for reliable comparison
expected_env = {k: str(v) for k, v in expected_env_raw.items()}
else:
expected_env = {}
# normalise running → dict
running_env = _list_to_dict(config.get('Env') or [])
for key, expected_val in expected_env.items():
running_val = running_env.get(key)
if running_val != expected_val:
self.logger.warning(
"Env var mismatch: %s expected=%s got=%s",
key, expected_val, running_val
)
return False
# ─────────────────── 5️⃣ RESOURCE LIMITS ───────────────────
# Check CPU shares
if "cpu_shares" in container_spec:
expected_cpu = int(container_spec["cpu_shares"])
# Scale to Docker's range - 1 CPU share = 1024 in Docker's scale
expected_docker_cpu = max(2, min(262144, expected_cpu * 1024))
running_cpu = host_config.get("CpuShares", 0)
if running_cpu != expected_docker_cpu:
self.logger.warning(
"CPU shares mismatch: expected %s (Docker value: %s), got %s",
expected_cpu, expected_docker_cpu, running_cpu
)
return False
# Check memory limit
if "mem_limit" in container_spec:
expected_mem = int(container_spec["mem_limit"]) * 1024 * 1024 # Convert MB to bytes
running_mem = host_config.get("Memory", 0)
# Allow small difference due to conversion rounding
if abs(running_mem - expected_mem) > 1024 * 1024: # 1MB tolerance
self.logger.warning(
"Memory limit mismatch: expected %s MB, got %s bytes",
container_spec["mem_limit"], running_mem
)
return False
# ─────────────────── 6️⃣ VOLUMES ───────────────────
expected_volumes = container_spec.get("storage", [])
if expected_volumes:
running_mounts = container.attrs.get("Mounts", [])
running_volume_paths = {m['Destination'] for m in running_mounts}
# Build expected paths, handling NFS volumes specially
expected_paths = set()
for v in expected_volumes:
if isinstance(v, dict) and 'mount_point' in v:
expected_paths.add(v['mount_point'])
if not expected_paths.issubset(running_volume_paths):
self.logger.warning(
"Volume mismatch: expected mounts %s, got %s",
expected_paths, running_volume_paths
)
return False
# ─────────────────── ✅ ALL CHECKS PASSED ───────────────────
return True
except Exception as exc:
self.logger.warning(
"Error validating container %s: %s", container.name, exc, exc_info=True
)
return False
def _process_storage_volumes(self, storage_volumes):
"""
Process storage volumes before container creation.
For NFS volumes, create the directory structure if it doesn't exist.
Args:
storage_volumes: List of volume configurations
"""
from settings import settings
import os
for volume in storage_volumes:
if volume.get('volume_type') == 'nfs':
# Get the shared filesystem root from settings
sharedfs_root = settings.get_value("SHAREDFS_ROOT", "/mnt/shared")
# Get the volume path
# volume_path = volume.get('volume_path')
# This is ignored
# if not volume_path:
# self.logger.warning(f"NFS volume missing path: {volume}")
# continue
# Construct the full path
expanded_path = os.path.join(sharedfs_root, volume.get('volume_id').lstrip('/'))
self.logger.info(f"Processing NFS volume at path: {expanded_path}")
volume['expanded_path']=expanded_path
# Create the directory if it doesn't exist
if not os.path.exists(expanded_path):
try:
self.logger.info(f"Creating NFS directory: {expanded_path}")
os.makedirs(expanded_path, exist_ok=True)
except OSError as e:
self.logger.error(f"Failed to create NFS directory {expanded_path}: {e}")
else:
self.logger.debug(f"NFS directory already exists: {expanded_path}")
def capture_container_statuses(self, pod_payload):
"""
Capture current status of all containers in the pod.
Args:
pod_payload: The pod update payload containing container specifications
Returns:
dict: Mapping of container_id to current status
"""
statuses = {}
# Get all containers in the pod from the payload
for container_spec in pod_payload.get("containers", []):
container_id = container_spec.get("container_id")
if container_id:
# Get current status from Docker
container = self.find_container(container_id)
if container:
statuses[container_id] = container.status
self.logger.debug(f"Captured status for container {container_id}: {container.status}")
else:
statuses[container_id] = "absent"
self.logger.debug(f"Container {container_id} not found, marking as absent")
self.logger.info(f"Captured initial statuses for {len(statuses)} containers")
return statuses
def blacklist_containers(self, pod_payload):
"""
Add all containers in the pod update to the DockerMonitor blacklist
to prevent event processing during the update.
Args:
pod_payload: The pod update payload containing container specifications
"""
if not self.docker_monitor:
self.logger.warning("No docker_monitor available for blacklisting")
return
blacklisted_count = 0
for container_spec in pod_payload.get("containers", []):
container_id = container_spec.get("container_id")
if container_id:
self.docker_monitor.add_to_blacklist(container_id)
blacklisted_count += 1
self.logger.info(f"Blacklisted {blacklisted_count} containers during pod update")
def unblacklist_containers(self, pod_payload):
"""
Remove all containers in the pod update from the DockerMonitor blacklist
after the update is complete.
Args:
pod_payload: The pod update payload containing container specifications
"""
if not self.docker_monitor:
self.logger.warning("No docker_monitor available for unblacklisting")
return
unblacklisted_count = 0
for container_spec in pod_payload.get("containers", []):
container_id = container_spec.get("container_id")
if container_id:
self.docker_monitor.remove_from_blacklist(container_id)
unblacklisted_count += 1
self.logger.info(f"Unblacklisted {unblacklisted_count} containers after pod update")
def compare_statuses(self, initial_statuses, final_statuses):
"""
Compare initial and final statuses to find changes.
Args:
initial_statuses: Container ID to initial status mapping
final_statuses: Container ID to final status mapping
Returns:
dict: Mapping of container_id to status changes
"""
changes = {}
all_container_ids = set(initial_statuses.keys()) | set(final_statuses.keys())
for container_id in all_container_ids:
initial_status = initial_statuses.get(container_id, "absent")
final_status = final_statuses.get(container_id, "absent")
if initial_status != final_status:
changes[container_id] = {
"from": initial_status,
"to": final_status
}
self.logger.debug(f"Status change detected for container {container_id}: {initial_status} -> {final_status}")
self.logger.info(f"Found {len(changes)} status changes out of {len(all_container_ids)} containers")
return changes
def send_delta_updates(self, status_changes):
"""
Send only the changed statuses to the server via the event queue.
Args:
status_changes: Container ID to status change mapping
"""
if not status_changes:
self.logger.info("No status changes to report")
return
if not self.docker_monitor or not hasattr(self.docker_monitor, 'event_queue'):
self.logger.warning("No docker_monitor or event_queue available for sending delta updates")
return
for container_id, change in status_changes.items():
# Get container name for the event
container_name = container_id # fallback to ID
container = self.find_container(container_id)
if container:
container_name = container.labels.get("name", container.name)
# Create event data for the status change
event_data = {
'source': 'docker',
'event_type': change['to'], # Use final status as event type
'container_id': container_id,
'system_container_id': container_id,
'container_name': container_name,
'timestamp': time.time(),
'details': {
'status': change['to'],
'previous_status': change['from'],
'Type': 'container',
'Action': change['to'],
'Actor': {
'ID': container_id,
'Attributes': {
'managed_by': 'worker_agent',
'name': container_name,
'system_container_id': container_id
}
}
},
'managed_by': 'worker_agent'
}
# Put event in the queue for the worker client
self.docker_monitor.event_queue.put(event_data)
self.logger.info(f"Sent delta update for container {container_id}: {change['from']} -> {change['to']}")
def shutdown(self):
"""Gracefully shutdown Docker connection."""
try:
self.docker_client.close()
except Exception:
pass
def attach_to_network(self, container, network_name):
"""Attach an existing container to a Docker network."""
try:
network = self.docker_client.networks.get(network_name)
network.connect(container)
self.logger.info(f"Attached container {container.name} to network {network_name}")
except docker.errors.NotFound:
self.logger.error(f"Network {network_name} not found when attaching container {container.name}")
except Exception as e:
self.logger.error(f"Failed to attach container {container.name} to network {network_name}: {str(e)}")
# Single log file request
def handle_container_log_request(self, data):
"""
Handles a container-log request and delegates to ContainerTask.
"""
self.logger.info(f"Handling container-log {data}")
container_name = data.get("container_name")
lines = int(data.get("lines", 100))
request_id = data.get("request_id")
try:
"""Fetch the last N lines of logs from a container."""
try:
container = self.docker_client.containers.get(container_name)
logs = container.logs(tail=lines, stdout=True, stderr=True).decode("utf-8")
result= {
"success": True,
"container_name": container_name,
"logs": logs
}
except Exception as e:
self.logger.error(f"Error fetching logs for {container_name}: {str(e)}")
result= {
"success": False,
"container_name": container_name,
"logs": f"Error: {str(e)}"
}
return result
except Exception as e:
self.logger.error(f"Error in container-log handler: {e}")
# Log streaming section
def stream_logs(self, container_name, on_log, on_error, cancel_flag):
"""
Stream logs from a container. This runs in a thread.
- on_log(line): callback for each log line.
- on_error(err): callback for error.
- cancel_flag: callable that returns True if stream should stop.
"""
try:
container = self.docker_client.containers.get(container_name)
for line in container.logs(stream=True, follow=True, tail=100):
if cancel_flag():
self.logger.info(f"Log stream for container {container_name} cancelled.")
break
on_log(line.decode("utf-8").strip())
except Exception as e:
self.logger.exception(f"Error streaming logs from {container_name}: {e}")
on_error(str(e))
def start_terminal(self, logger, container_name, command="/bin/bash"):
master_fd, slave_fd = pty.openpty()
process = subprocess.Popen(
["docker", "exec", "-it", container_name, command],
stdin=slave_fd, stdout=slave_fd, stderr=slave_fd,
universal_newlines=True
)
return TerminalSession(logger, process, master_fd)
class TerminalSession:
def __init__(self, logger, process, master_fd):
self.process = process
self.master_fd = master_fd
logger.debug(f"Starting TerminalSession")
def write(self, data):
os.write(self.master_fd, data.encode())
def stream_output(self, logger, on_output, cancel_flag):
logger.debug(f"Starting stream_output")
try:
while True:
if cancel_flag():
break
r, _, _ = select.select([self.master_fd], [], [], 0.1)
if r:
output = os.read(self.master_fd, 1024).decode("utf-8", errors="ignore")
logger.debug(f"output - {output}")
on_output(output)
finally:
self.close()
def close(self):
try:
os.close(self.master_fd)
except Exception:
pass
try:
self.process.terminate()
except Exception:
pass