Files
3cloud-backend/worker/worker_tasks/container.py
T
JamesBhattarai 7ec2e4d47b Feat: Implement NAT for Private network
Added enable_nat bool for networks to allow nats
Implemented docker network bridge, to allow nat with tenancy seperated as `enable_icc:false`
2026-08-08 17:41:15 +02:00

2244 lines
104 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 shutil
import threading
import base64
from typing import List, Dict, Any, Optional, Callable
from settings import settings
from worker_tasks.ovs_sdn import OVS_SDN
# Import DNS configuration constants
try:
from app.utils.constants import GLOBAL_DNS_CONFIG_BASE_PATH, DNS_CONTAINER_MOUNT_PATH, NSCONTROLLER_NAT_NETWORK
except ImportError:
# Fallback if app.utils.constants is not available in worker environment
GLOBAL_DNS_CONFIG_BASE_PATH = "/var/lib/xcloudify/dns-configs"
DNS_CONTAINER_MOUNT_PATH = "/etc/dnsmasq.d"
NSCONTROLLER_NAT_NETWORK = "xcloudify-wan"
def _should_use_sudo() -> bool:
"""
Determine if sudo should be used for network namespace commands.
Returns False when running inside Docker where sudo is not available.
"""
# Check if /.dockerenv exists (Docker indicator)
if os.path.exists('/.dockerenv'):
return False
# Check if sudo command is available
if shutil.which('sudo') is None:
return False
return True
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
# Only send delta updates if there are no launch failures
status_changes = self.compare_statuses(initial_statuses, final_statuses)
if status_changes and not self.launch_failures:
self.send_delta_updates(status_changes)
elif self.launch_failures:
self.logger.info(
f"Skipping delta status updates due to {len(self.launch_failures)} launch failures"
)
# 7. Send a full-picture status sync for convergence after reconciliation
# Only send status sync if there are no launch failures - the launch_failed
# event is already sent separately and the server will handle the status
if not self.launch_failures:
self.send_full_statuses_for_host(pod_payload, scope="pod")
else:
self.logger.info(
f"Skipping full status sync due to {len(self.launch_failures)} launch failures"
)
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 _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)
nscontroller_launch_failed = False
nscontroller_container = None
for ns_container_spec in nscontroller_specs:
desired_state = ns_container_spec.get("desired_state", "running")
if desired_state == "deleted":
self.logger.debug(
"Processing NSController container_id=%s (action: delete)",
ns_container_spec.get("container_id"),
)
else:
self.logger.debug(
"Processing NSController container_id=%s (action: ensure %s)",
ns_container_spec.get("container_id"),
desired_state,
)
result = self.ensure_container(ns_container_spec, is_nscontroller=True)
# Check if NSController launch failed
container_id = ns_container_spec.get("container_id")
if result is None and any(failure.get("id") == container_id for failure in self.launch_failures):
nscontroller_launch_failed = True
self.logger.error(f"NSController container {container_id} failed to launch. Aborting pod launch.")
# Find the failure details and log them
for failure in self.launch_failures:
if failure.get("id") == container_id:
self.logger.error(f"NSController failure reason: {failure.get('error')}")
break
# Dump NSController console logs for debugging
self._dump_container_console_logs(container_id)
else:
nscontroller_container = result
# ---- Wait for NSController health check to pass before launching dependents
if not nscontroller_launch_failed and nscontroller_container:
health_check_passed = self.wait_for_nscontroller_health(nscontroller_container, nscontroller_container_id)
if not health_check_passed:
nscontroller_launch_failed = True
self.logger.error(f"NSController {nscontroller_container_id} health check failed. Cleaning up container.")
# Add NSController to launch_failures since health check failed
self.launch_failures.append({
"id": nscontroller_container_id,
"error": "NSController health check failed"
})
# Dump NSController console logs for debugging
self._dump_container_console_logs(nscontroller_container_id)
# Clean up the NSController container that failed health check
self._cleanup_failed_nscontroller_container(nscontroller_container, nscontroller_container_id)
# ---- launch regular containers with networking pointer if available
# Only proceed if NSController launch was successful
if not nscontroller_launch_failed:
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"),
)
desired_state = regular_container_spec.get("desired_state", "running")
nscontroller_id = regular_container_spec.get("nscontroller_container_id")
if desired_state == "deleted":
self.logger.debug(
"Processing regular container_id=%s (action: delete)",
regular_container_spec.get("container_id"),
)
else:
self.logger.debug(
"Processing regular container_id=%s (action: ensure %s, nscontroller: %s)",
regular_container_spec.get("container_id"),
desired_state,
nscontroller_id or "none",
)
self.ensure_container(regular_container_spec)
else:
self.logger.warning("Skipping launch of regular containers due to NSController launch failure.")
# Add all regular containers to launch_failures since they were skipped
for regular_container_spec in regular_container_specs:
desired_state = regular_container_spec.get("desired_state", "running")
if desired_state != "deleted":
container_id = regular_container_spec.get("container_id")
if container_id:
self.launch_failures.append({
"id": container_id,
"error": "Skipped due to NSController launch failure"
})
# 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 _dump_container_console_logs(self, container_id):
"""
Dump container console logs to debug for troubleshooting.
Args:
container_id: The system container ID to fetch logs for
"""
try:
container = self.find_container(container_id)
if container:
# Get last 100 lines of logs
logs = container.logs(tail=100, stdout=True, stderr=True).decode("utf-8", errors="replace")
if logs:
# Truncate if too long for logging
log_lines = logs.split('\n')
if len(log_lines) > 20:
logs = '\n'.join(log_lines[:10] + ['...'] + log_lines[-10:])
self.logger.debug(f"Container {container_id} console logs: {logs}")
else:
self.logger.debug(f"Container {container_id} has no console output")
else:
self.logger.debug(f"Container {container_id} not found for console log dump")
except Exception as e:
self.logger.debug(f"Failed to dump console logs for {container_id}: {e}")
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)
# Log the action that will be taken at debug level
if desired_state == "deleted":
self.logger.debug(f"ensure_container container_id={container_id} action=delete")
elif restart_requested:
self.logger.debug(f"ensure_container container_id={container_id} action=restart")
elif desired_state == "stopped":
self.logger.debug(f"ensure_container container_id={container_id} action=stop")
elif desired_state == "running":
if existing:
self.logger.debug(f"ensure_container container_id={container_id} action=ensure_running existing=True")
else:
self.logger.debug(f"ensure_container container_id={container_id} action=ensure_running existing=False")
else:
self.logger.debug(f"ensure_container container_id={container_id} action=unknown desired_state={desired_state}")
# Handle deletion
if desired_state == "deleted":
if existing:
self.logger.info(f"Deleting container {container_id}")
# Get network_ports from container_spec for port cleanup
network_ports = container_spec.get("network_ports", [])
self.delete_container(existing, network_ports)
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
}
if not self.docker_monitor.is_blacklisted(container_id):
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
# Check if OVS bridge is explicitly configured
network_ports = container_spec.get("network_ports", [])
has_ovs_bridge = any(p.get("ovs_bridge") for p in network_ports)
enable_host_nat = bool(container_spec.get("enable_host_nat", False))
if is_nscontroller:
if enable_host_nat:
# Give the NSController a uplink (eth0 + default route +
# Docker's automatic MASQUERADE) so its entrypoint can NAT the
# tenant subnet to the internet.
network_mode = self._ensure_nat_network()
self.logger.info(f"NSController: using network_mode={network_mode} (host NAT enabled)")
elif has_ovs_bridge:
# Use OVS networking - container gets no Docker network, OVS ports attached later
network_mode = "none"
self.logger.info("NSController: using network_mode=none (OVS ports will be added)")
else:
# Use Docker bridge for port publishing
network_mode = "bridge"
self.logger.info("NSController: using network_mode=bridge (Docker port publishing)")
elif 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": settings.get_value("WORKER_ID"),
"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 is_nscontroller:
container_kwargs["cap_add"] = ["NET_ADMIN", "NET_RAW"]
self.logger.info(f"NSController {system_container_id}: adding NET_ADMIN + NET_RAW capabilities")
if enable_host_nat:
container_kwargs["sysctls"] = {"net.ipv4.ip_forward": "1"}
# Add restart policy from container spec if provided, otherwise use default from settings
restart_policy = container_spec.get("restart_policy", settings.get_value("RESTART_POLICY", "always"))
self.logger.info(f"Using restart policy '{restart_policy}' for container {system_container_id}")
container_kwargs["restart_policy"] = {"Name": restart_policy}
# 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" 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"])
# 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})")
else:
self.logger.debug(f"CPU not specified in container spec, not setting")
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}")
# Handle DNS configuration for NSControllers with use_dns flag
use_dns = container_spec.get("use_dns", False)
if use_dns and is_nscontroller:
self.logger.info(f"Handling DNS configuration for NSController {system_container_id}")
# Extract VDC ID from container spec
vdc_id = container_spec.get("vdc_id")
if not vdc_id:
self.logger.error(f"use_dns is true but vdc_id is missing in container spec")
self.launch_failures.append({
"id": system_container_id,
"error": "use_dns requires vdc_id in container specification"
})
return None
# Construct VDC-specific DNS config directory path for tenant isolation
# Each VDC gets its own subdirectory to prevent cross-tenant DNS leakage
dns_config_dir = f"{GLOBAL_DNS_CONFIG_BASE_PATH}/vdc-{vdc_id}"
self.logger.info(f"DNS config directory for VDC {vdc_id}: {dns_config_dir}")
# Create DNS config directory if it doesn't exist
try:
os.makedirs(dns_config_dir, mode=0o755, exist_ok=True)
self.logger.info(f"Ensured DNS config directory exists: {dns_config_dir}")
except Exception as e:
self.logger.error(f"Failed to create DNS config directory {dns_config_dir}: {e}")
self.launch_failures.append({
"id": system_container_id,
"error": f"Failed to create DNS config directory: {str(e)}"
})
return None
# Write initial DNS configuration if provided
dns_config = container_spec.get("dns_config")
if dns_config:
zone_file = dns_config.get("zone_file")
zone_name = dns_config.get("zone_name")
if zone_file and zone_name:
config_filename = f"{vdc_id}.conf"
config_path = os.path.join(dns_config_dir, config_filename)
try:
# Write DNS config atomically using temp file
temp_path = config_path + ".tmp"
with open(temp_path, 'w') as f:
f.write(zone_file)
if not zone_file.endswith('\n'):
f.write('\n')
f.flush()
os.fsync(f.fileno())
# Set proper permissions before rename
os.chmod(temp_path, 0o644)
# Atomic rename
os.rename(temp_path, config_path)
self.logger.info(f"Initial DNS config written to {config_path}")
except Exception as e:
# Clean up temp file if it exists
if os.path.exists(temp_path):
try:
os.unlink(temp_path)
except Exception:
pass
self.logger.error(f"Failed to write initial DNS config: {e}")
# Don't fail container creation, just log the error
# DNS can be updated later via DNS update task
else:
self.logger.debug("dns_config provided but missing zone_file or zone_name")
else:
self.logger.debug("No initial dns_config provided in container spec")
# Add DNS bind mount to container volumes
if "volumes" not in container_kwargs:
container_kwargs["volumes"] = {}
container_kwargs["volumes"][dns_config_dir] = {
"bind": DNS_CONTAINER_MOUNT_PATH,
"mode": "rw" # Read-write: entrypoint writes 00-provider.conf here
}
self.logger.info(
f"Added DNS bind mount: {dns_config_dir} -> {DNS_CONTAINER_MOUNT_PATH} (read-write)"
)
# Add DNS configuration for NSControllers
if use_dns and is_nscontroller:
container_kwargs["dns"] = ["127.0.0.1"]
container_kwargs["dns_opt"] = [] # Prevent Docker from adding Google DNS
container_kwargs["dns_search"] = []
# Add metadata socket mount for NSControllers (for metadata proxy communication)
if is_nscontroller:
if "volumes" not in container_kwargs:
container_kwargs["volumes"] = {}
container_kwargs["volumes"]["/tmp/xcloudify"] = {
"bind": "/tmp/xcloudify",
"mode": "rw"
}
self.logger.info(
f"Added metadata socket mount: /tmp/xcloudify -> /tmp/xcloudify (read-write)"
)
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 (NSControllers skip this - they use OVS ports only)
if not is_nscontroller and network_mode is None:
for network_name in networks:
self.attach_to_network(container, network_name)
# Step 3 – configure network ports if specified
network_ports = container_spec.get("network_ports", [])
if network_ports:
self.logger.info(f"Configuring {len(network_ports)} network ports for container {container.name}")
# Initialize OVS_SDN for port configuration
ovs_sdn = OVS_SDN({}, self.logger)
try:
for port_config in network_ports:
port_name = port_config.get("name")
mac_address = port_config.get("mac_address")
ovs_bridge = port_config.get("ovs_bridge", "br-int")
if port_name and mac_address:
self.logger.info(f"Creating OVS port {port_name} with MAC {mac_address} on bridge {ovs_bridge}")
ovs_sdn.port_mgr.ensure_vm_port(port_name, mac_address, ovs_bridge)
else:
self.logger.warning(f"Skipping port configuration: missing port_name or mac_address")
except Exception as e:
self.logger.error(f"Failed to configure network ports: {e}")
# Record failure
self.launch_failures.append({
"id": system_container_id,
"error": f"Network port configuration failed: {str(e)}"
})
# Clean up the container that was already created
try:
self.logger.info(f"Cleaning up container {container.name} due to port configuration failure")
container.stop(timeout=5)
container.remove(force=True)
except Exception as cleanup_error:
self.logger.warning(f"Failed to clean up container {container.name}: {cleanup_error}")
return None
# Step 4 – Attach OVS ports to container's network namespace (for NSControllers with OVS bridge)
if is_nscontroller and has_ovs_bridge and network_ports:
try:
self.attach_ovs_ports_to_container_namespace(container, network_ports)
except Exception as e:
self.logger.error(f"Failed to attach OVS ports to container namespace: {e}")
self.logger.error(traceback.format_exc())
# Record failure
self.launch_failures.append({
"id": system_container_id,
"error": f"Failed to attach OVS ports to container namespace: {str(e)}"
})
# Clean up the container that was already created
try:
self.logger.info(f"Cleaning up container {container.name} due to port attachment failure")
container.stop(timeout=5)
container.remove(force=True)
except Exception as cleanup_error:
self.logger.warning(f"Failed to clean up container {container.name}: {cleanup_error}")
return None
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 wait_for_nscontroller_health(self, container, container_id, timeout=120, check_interval=2):
"""
Wait for NSController container health check to pass before launching dependent containers.
This method polls the container's health status and waits until:
1. The container is running
2. The Docker health check reports "healthy"
Args:
container: The Docker container object
container_id: The system container ID for logging
timeout: Maximum time to wait in seconds (default: 120)
check_interval: Time between health checks in seconds (default: 5)
Returns:
bool: True if health check passed, False if timeout or container failed
"""
self.logger.info(f"Waiting for NSController {container_id} health check to pass (timeout: {timeout}s)")
start_time = time.time()
last_status = None
while True:
elapsed = time.time() - start_time
if elapsed > timeout:
self.logger.error(f"NSController {container_id} health check timed out after {elapsed:.1f}s")
return False
try:
# Reload container to get latest status
container.reload()
current_status = container.status
# Check if container is still running
if current_status != "running":
self.logger.warning(f"NSController {container_id} status changed to {current_status}, expected 'running'")
return False
# Check Docker health check status
health = container.attrs.get("State", {}).get("Health", {})
health_status = health.get("Status", "unknown")
failing_streak = health.get("FailingStreak", 0)
if health_status != last_status:
self.logger.info(f"NSController {container_id} health status: {health_status} (elapsed: {elapsed:.1f}s)")
last_status = health_status
if health_status == "healthy":
self.logger.info(f"NSController {container_id} health check passed after {elapsed:.1f}s")
return True
# Check for health check failures
if failing_streak > 0:
self.logger.debug(f"NSController {container_id} health check failing streak: {failing_streak}")
# Log health check log if available
health_log = health.get("Log", [])
if health_log:
last_check = health_log[-1] if health_log else {}
output = last_check.get("Output", "")
if output:
self.logger.debug(f"NSController {container_id} last health check output: {output.strip()}")
except docker.errors.NotFound:
self.logger.error(f"NSController {container_id} container not found during health check")
return False
except Exception as e:
self.logger.warning(f"Error checking NSController {container_id} health: {e}")
self.logger.debug(f"Waiting for NSController {container_id} health check ({elapsed:.1f}s elapsed)")
time.sleep(check_interval)
def _cleanup_failed_nscontroller_container(self, container, container_id):
"""
Clean up an NSController container that failed to become healthy.
This method stops and removes the container, cleans up network namespace
symlinks, and injected files.
Args:
container: The Docker container object to clean up
container_id: The system container ID for logging
"""
if not container:
self.logger.warning(f"No container object provided for cleanup of {container_id}")
return
self.logger.info(f"Cleaning up failed NSController container {container_id}")
# Use the delete_container method which handles full cleanup
# Note: We don't have access to the original network_ports spec here,
# but the namespace cleanup will handle the OVS port detachment
self.delete_container(container, network_ports=None)
self.logger.info(f"Successfully cleaned up failed NSController container {container_id}")
def cleanup_container_namespace(self, container_id: str):
"""Clean up the network namespace symlink created for OVS port integration."""
# The symlink is created at /run/netns/<container_id>
netns_path = f"/run/netns/{container_id}"
sudo_prefix = ['sudo'] if _should_use_sudo() else []
try:
# Remove the symlink
if os.path.islink(netns_path):
os.unlink(netns_path)
self.logger.info(f"Cleaned up network namespace symlink: {netns_path}")
else:
self.logger.debug(f"No namespace symlink found at {netns_path}")
except Exception as e:
self.logger.warning(f"Failed to cleanup namespace symlink {netns_path}: {e}")
def delete_container(self, container, network_ports=None):
"""Stop and remove a Docker container, and clean up network ports if specified."""
system_container_id = container.labels.get("system_container_id")
docker_id = container.id
# Clean up network namespace symlink before deleting the container
self.cleanup_container_namespace(system_container_id)
# Clean up network ports before deleting the container
if network_ports:
self.logger.info(f"Cleaning up {len(network_ports)} network ports for container {container.name}")
try:
ovs_sdn = OVS_SDN({}, self.logger)
for port_config in network_ports:
port_name = port_config.get("name")
ovs_bridge = port_config.get("ovs_bridge", "br-int")
if port_name:
self.logger.info(f"Deleting OVS port: {port_name} from bridge {ovs_bridge}")
ovs_sdn.port_mgr.delete_port(port_name, ovs_bridge)
except Exception as e:
self.logger.warning(f"Failed to clean up network ports: {e}")
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
# ─────────────────── 2b️⃣ NETWORK MODE (NSController only) ───────────────────
# NSControllers use the NAT uplink network if host NAT is enabled,
# else network_mode="none" if OVS ports will be added, else "bridge"
# for Docker port publishing (matching launch_container logic)
if container_spec.get("workload_type") == "NSController":
network_ports = container_spec.get("network_ports", [])
has_ovs_bridge = any(p.get("ovs_bridge") for p in network_ports)
if container_spec.get("enable_host_nat"):
expected_network_mode = NSCONTROLLER_NAT_NETWORK
elif has_ovs_bridge:
expected_network_mode = "none"
else:
expected_network_mode = "bridge"
running_network_mode = host_config.get("NetworkMode", "")
if running_network_mode != expected_network_mode:
self.logger.warning(
"NSController network mode mismatch: expected '%s', got '%s'",
expected_network_mode, running_network_mode
)
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" in container_spec:
expected_cpu = int(container_spec["cpu"])
# 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
# ─────────────────── 7️⃣ NSController Network Mode (Regular containers only) ───────────────────
# For regular containers, check if they're connected to the correct NSController
if container_spec.get("nscontroller_container_id"):
nscontroller_container_id = container_spec["nscontroller_container_id"]
self.logger.debug(f"Checking NSController network mode for container {container_spec.get('container_id')}")
# Find the current NSController container
nscontroller = self.find_container(nscontroller_container_id)
if nscontroller:
# Get the current network mode from running container
running_network_mode = host_config.get("NetworkMode", "")
expected_network_mode = f"container:{nscontroller.id}"
self.logger.debug(f"Running network mode: {running_network_mode}")
self.logger.debug(f"Expected network mode: {expected_network_mode}")
# If the container is not using the current NSController, it needs recreation
if running_network_mode != expected_network_mode:
self.logger.warning(
"NSController mismatch: container is connected to %s, should be connected to %s",
running_network_mode, expected_network_mode
)
return False
else:
# NSController not found - this will be handled by the main reconciliation logic
self.logger.warning(f"NSController {nscontroller_container_id} not found for container {container_spec.get('container_id')}")
# ─────────────────── ✅ 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
"""
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'],
},
'managed_by': settings.get_value("WORKER_ID")
}
# Put event in the queue for the worker client
if not self.docker_monitor.is_blacklisted(container_id):
self.docker_monitor.event_queue.put(event_data)
self.logger.info(f"Sent delta update for container {container_id}: {change['from']} -> {change['to']}")
def send_full_statuses_for_host(self, pod_payload, scope="host"):
"""
Emit a full-picture status snapshot for convergence after reconciliation.
Also handles detection and reporting of containers that are expected to exist
but are missing on the host.
scope:
- "host": all containers managed by this worker (label managed_by=<WORKER_ID>)
- "pod": only containers present in pod_payload
"""
if not self.docker_monitor or not hasattr(self.docker_monitor, 'event_queue'):
self.logger.warning("No docker_monitor or event_queue available for full status sync")
return
try:
worker_id = settings.get_value("WORKER_ID")
except Exception:
worker_id = None
allowed_ids = None
if scope == "pod":
allowed_ids = {
c.get("container_id")
for c in (pod_payload or {}).get("containers", [])
if c.get("container_id")
}
try:
containers = self.docker_client.containers.list(
all=True,
filters={"label": f"managed_by={worker_id}"} if worker_id else None
)
except Exception as e:
self.logger.error(f"Failed to list containers for full status sync: {e}")
return
# Create a set of container IDs that actually exist on the host
existing_container_ids = set()
sent = 0
skipped = 0
for c in containers:
system_id = (c.labels or {}).get("system_container_id")
if scope == "pod" and allowed_ids is not None and system_id not in allowed_ids:
continue
if not system_id:
skipped += 1
continue
# Add to existing container IDs
existing_container_ids.add(system_id)
name = (c.labels or {}).get("name", c.name)
st = (c.status or "").lower()
# Map Docker status -> accepted event types for websocket_server
if st == "running":
event_type = "running"
elif st in ("exited", "dead"):
event_type = "die"
else:
# Skip unsupported/neutral states to avoid confusing the API
skipped += 1
continue
event_data = {
'source': 'docker',
'event_type': event_type,
'container_id': c.id,
'system_container_id': system_id,
'container_name': name,
'timestamp': time.time(),
'details': {
'status': event_type,
'Type': 'container',
'Action': event_type,
},
'managed_by': worker_id
}
if not self.docker_monitor.is_blacklisted(system_id):
self.docker_monitor.event_queue.put(event_data)
self.logger.debug(f"Added container status event to queue: {event_type} for container {c.id}")
sent += 1
# Handle containers that are expected to exist according to the server but are missing on the host
if scope == "pod" and pod_payload:
# Get the list of container IDs that should exist
server_container_ids = set()
containers_data = pod_payload.get("containers", [])
for container_data in containers_data:
container_id = container_data.get("container_id")
if container_id:
server_container_ids.add(container_id)
# Find missing containers
missing_container_ids = server_container_ids - existing_container_ids
for missing_container_id in missing_container_ids:
self.logger.info(f"Reporting missing container {missing_container_id} to server")
# Find the container data in the workloads
missing_container_data = None
for container_data in containers_data:
if container_data.get("container_id") == missing_container_id:
missing_container_data = container_data
break
# Use container name from data if available, otherwise use container ID
container_name = missing_container_data.get("container_name", missing_container_id) if missing_container_data else missing_container_id
# Create event data for the missing container
event_data = {
'source': 'docker',
'event_type': 'absent', # Report as 'die' since it's missing
'container_id': missing_container_id,
'system_container_id': missing_container_id,
'container_name': container_name,
'timestamp': time.time(),
'details': {
'status': 'absent',
'id': missing_container_id,
'from': missing_container_data.get("docker_image", "unknown") if missing_container_data else "unknown",
'Type': 'container',
'Action': 'absent',
},
'scope': 'local',
'time': int(time.time()),
'timeNano': int(time.time() * 1000000000)
}
# Put event in the queue for the worker client
if self.docker_monitor and hasattr(self.docker_monitor, 'event_queue'):
self.docker_monitor.event_queue.put(event_data)
self.logger.debug(f"Added missing container event to queue: die for container {container_name}")
sent += 1
self.logger.info(f"Full status sync sent: scope={scope} sent={sent} skipped={skipped}")
def shutdown(self):
"""Gracefully shutdown Docker connection."""
try:
self.docker_client.close()
except Exception:
pass
def _ensure_nat_network(self) -> str:
"""
Ensure the NAT uplink docker bridge network exists on this host.
NSControllers launched here (instead of network_mode="none"/"bridge")
get a real eth0 + default route + Docker's automatic MASQUERADE for
that subnet, which their entrypoint uses to NAT the tenant network's
egress traffic out to the internet.
"""
try:
self.docker_client.networks.get(NSCONTROLLER_NAT_NETWORK)
except docker.errors.NotFound:
self.logger.info(f"Creating NAT uplink network '{NSCONTROLLER_NAT_NETWORK}'")
try:
self.docker_client.networks.create(
NSCONTROLLER_NAT_NETWORK,
driver="bridge",
options={
"com.docker.network.bridge.name": NSCONTROLLER_NAT_NETWORK,
# NSControllers for different tenants share this uplink network;
# they shouldn't be able to reach each other directly through it.
"com.docker.network.bridge.enable_icc": "false",
},
)
except docker.errors.APIError as e:
# Another concurrent launch on this host may have just created it.
if not self._docker_network_exists(NSCONTROLLER_NAT_NETWORK):
raise
self.logger.debug(f"NAT uplink network '{NSCONTROLLER_NAT_NETWORK}' created concurrently: {e}")
return NSCONTROLLER_NAT_NETWORK
def _docker_network_exists(self, network_name: str) -> bool:
try:
self.docker_client.networks.get(network_name)
return True
except docker.errors.NotFound:
return False
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)}")
def attach_ovs_ports_to_container_namespace(self, container, network_ports):
"""
Attach OVS ports to a container's network namespace.
This method moves OVS ports created on the host into the container's network
namespace and configures them with the specified MAC address and IP address.
Args:
container: The Docker container object (must be running with network_mode="none")
network_ports: List of port configurations, each containing:
- name: port name (e.g., "port-645611")
- mac_address: MAC address (e.g., "fa:16:4a:6e:7b:d1")
- ip_address: IP address (e.g., "10.0.178.12")
- ovs_bridge: OVS bridge name (e.g., "br-int")
"""
container_name = container.name
system_container_id = container.labels.get("system_container_id", container_name)
self.logger.info(f"Attaching {len(network_ports)} OVS ports to container namespace: {container_name}")
# Determine sudo prefix based on environment
sudo_prefix = ['sudo'] if _should_use_sudo() else []
# Get the container's PID using docker inspect
try:
result = subprocess.run(
sudo_prefix + ["docker", "inspect", "--format='{{.State.Pid}}'", system_container_id],
capture_output=True,
text=True
)
if result.returncode != 0:
raise RuntimeError(f"Failed to get container PID: {result.stderr}")
# Extract PID from output (remove any quotes and whitespace)
container_pid = result.stdout.strip().strip("'").strip('"')
if not container_pid or not container_pid.isdigit():
raise RuntimeError(f"Invalid container PID: {result.stdout}")
self.logger.debug(f"Container {container_name} PID: {container_pid}")
except Exception as e:
self.logger.error(f"Failed to get container PID: {e}")
raise
# Create /run/netns directory if it doesn't exist
netns_dir = "/run/netns"
try:
os.makedirs(netns_dir, exist_ok=True)
self.logger.debug(f"Ensured /run/netns directory exists")
except Exception as e:
self.logger.error(f"Failed to create /run/netns directory: {e}")
raise
# Create symlink at /run/netns/<container_id> pointing to /proc/<pid>/ns/net
netns_path = f"{netns_dir}/{system_container_id}"
proc_ns_path = f"/proc/{container_pid}/ns/net"
# Remove any existing symlink
try:
if os.path.islink(netns_path):
os.unlink(netns_path)
self.logger.debug(f"Removed existing symlink at {netns_path}")
except OSError:
pass
# Create the symlink
try:
os.symlink(proc_ns_path, netns_path)
self.logger.debug(f"Created namespace symlink: {netns_path} -> {proc_ns_path}")
except Exception as e:
self.logger.error(f"Failed to create namespace symlink: {e}")
raise
# Use system_container_id as the namespace name for ip netns commands
ns_name = system_container_id
for port_config in network_ports:
port_name = port_config.get("name")
mac_address = port_config.get("mac_address")
ip_address = port_config.get("ip_address")
if not port_name or not mac_address or not ip_address:
self.logger.warning(f"Skipping port configuration: missing required fields (name={port_name}, mac={mac_address}, ip={ip_address})")
continue
self.logger.info(f"Attaching OVS port {port_name} to container {container_name}")
self.logger.debug(f" MAC: {mac_address}, IP: {ip_address}")
try:
# Determine sudo prefix based on environment
sudo_prefix = ['sudo'] if _should_use_sudo() else []
# Step 1: Remove any existing symlink for this namespace
result = subprocess.run(
sudo_prefix + ["ip", "netns", "del", ns_name],
capture_output=True,
text=True
)
if result.returncode != 0 and "No such file or directory" not in result.stderr:
self.logger.warning(f"Failed to delete existing namespace symlink: {result.stderr}")
# Step 2: Create symlink to make namespace accessible via ip netns
result = subprocess.run(
sudo_prefix + ["ln", "-s", proc_ns_path, netns_path],
capture_output=True,
text=True
)
if result.returncode != 0:
raise RuntimeError(f"Failed to create namespace symlink: {result.stderr}")
self.logger.debug(f"Created namespace symlink: {netns_path} -> {proc_ns_path}")
# Step 3: Move the OVS port into the namespace
result = subprocess.run(
sudo_prefix + ["ip", "link", "set", port_name, "netns", ns_name],
capture_output=True,
text=True
)
if result.returncode != 0:
raise RuntimeError(f"Failed to move port to namespace: {result.stderr}")
self.logger.debug(f"Moved port {port_name} to namespace {ns_name}")
# Step 4: Configure the port inside the container namespace
# Set MAC address
result = subprocess.run(
sudo_prefix + ["ip", "netns", "exec", ns_name, "ip", "link", "set", "address", mac_address, "dev", port_name],
capture_output=True,
text=True
)
if result.returncode != 0:
raise RuntimeError(f"Failed to set MAC address: {result.stderr}")
self.logger.debug(f"Set MAC address {mac_address} on port {port_name}")
# Bring the port up
result = subprocess.run(
sudo_prefix + ["ip", "netns", "exec", ns_name, "ip", "link", "set", "dev", port_name, "up"],
capture_output=True,
text=True
)
if result.returncode != 0:
raise RuntimeError(f"Failed to bring port up: {result.stderr}")
self.logger.debug(f"Brough port {port_name} up")
prefix = port_config.get("subnet_mask")
if not prefix:
self.logger.warning(
f"No subnet_mask for port {port_name}; defaulting to /24"
)
prefix = 24
result = subprocess.run(
sudo_prefix + ["ip", "netns", "exec", ns_name, "ip", "addr", "add", f"{ip_address}/{prefix}", "dev", port_name],
capture_output=True,
text=True
)
if result.returncode != 0:
raise RuntimeError(f"Failed to add IP address: {result.stderr}")
self.logger.debug(f"Added IP address {ip_address}/{prefix} to port {port_name}")
# Ensure loopback is up inside the namespace
result = subprocess.run(
sudo_prefix + ["ip", "netns", "exec", ns_name, "ip", "link", "set", "lo", "up"],
capture_output=True,
text=True
)
if result.returncode != 0:
self.logger.warning(f"Failed to bring loopback up: {result.stderr}")
else:
self.logger.debug("Loopback interface up in namespace")
self.logger.info(f"Successfully attached and configured OVS port {port_name}")
except Exception as e:
self.logger.error(f"Failed to attach OVS port {port_name}: {e}")
raise
self.logger.info(f"Completed attaching {len(network_ports)} OVS ports to container {container_name}")
# 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)
def handle_reconcile_and_delete(self, job_details):
"""
Handle the reconcile_and_delete task by deleting containers that are not in the expected list.
This method:
1. Gets the list of expected container IDs from job_details
2. Finds all containers running on this host with managed_by=<WORKER_ID>
3. Deletes any container that is NOT in the expected list
Args:
job_details (dict): Contains expected_container_ids list and the
authoritative flag saying whether that list can be trusted
Returns:
dict: Result of the reconciliation operation
"""
expected_container_ids = set(job_details.get("expected_container_ids", []))
deleted_containers = []
failed_deletions = []
# An empty expected set means "delete everything", so only act on it when the
# server confirms it actually read the set from the API. Without this, an API
# blip during reconciliation would wipe every container on the host.
if not job_details.get("authoritative", False):
self.logger.warning(
"Skipping reconcile_and_delete: expected container set is not authoritative"
)
return {
"success": True,
"skipped": True,
"reason": "non_authoritative_expected_set",
"deleted_containers": [],
"failed_deletions": [],
"expected_count": len(expected_container_ids),
"deleted_count": 0,
"failed_count": 0,
}
self.logger.info(f"Starting reconcile_and_delete: expected {len(expected_container_ids)} containers")
try:
# Get all containers managed by this worker
worker_id = settings.get_value("WORKER_ID")
all_containers = self.docker_client.containers.list(all=True, filters={"label": f"managed_by={worker_id}"})
self.logger.info(f"Found {len(all_containers)} containers managed by worker {worker_id}")
for container in all_containers:
system_container_id = container.labels.get("system_container_id")
# Skip containers that are in the expected list
if system_container_id in expected_container_ids:
self.logger.debug(f"Keeping expected container: {system_container_id}")
continue
# Delete containers that are not expected
self.logger.info(f"Deleting unexpected container: {system_container_id}")
try:
container.stop(timeout=5)
container.remove(force=True)
deleted_containers.append(system_container_id)
self.logger.info(f"Successfully deleted container: {system_container_id}")
except Exception as e:
self.logger.error(f"Failed to delete container {system_container_id}: {e}")
failed_deletions.append({
"container_id": system_container_id,
"error": str(e)
})
result = {
"success": len(failed_deletions) == 0,
"deleted_containers": deleted_containers,
"failed_deletions": failed_deletions,
"expected_count": len(expected_container_ids),
"deleted_count": len(deleted_containers),
"failed_count": len(failed_deletions)
}
self.logger.info(f"Reconcile and delete completed: {len(deleted_containers)} deleted, {len(failed_deletions)} failed")
return result
except Exception as e:
self.logger.error(f"Error during reconcile_and_delete: {e}")
return {
"success": False,
"error": str(e),
"deleted_containers": deleted_containers,
"failed_deletions": failed_deletions
}
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