Added capability to launch with ENV params and add a container to an existing pod

This commit is contained in:
2025-07-26 00:54:01 +09:30
parent 4cceda071c
commit 98b009ab1c
4 changed files with 541 additions and 370 deletions
+37 -25
View File
@@ -53,7 +53,7 @@ def validate_payload(payload):
uuid.UUID(payload['pod'], version=4)
sanitized_payload['pod'] = payload['pod']
except ValueError:
raise ValueError("'pod' must be a valid UUID.")
raise ValueError(f"'pod' must be a valid UUID. got {type(payload['pod'])} {payload['pod']}")
# Validate container list
if 'containers' not in payload:
@@ -63,36 +63,30 @@ def validate_payload(payload):
sanitized_payload['containers'] = []
# Validate each container in the 'containers' list and build the sanitized version
for container in payload['containers']:
if not isinstance(container, dict):
raise ValueError("Each container must be a dictionary.")
# Check for required keys: 'docker_image' and 'container_name'
if 'docker_image' not in container:
raise ValueError("Container is missing 'docker_image'.")
if 'container_name' not in container:
raise ValueError("Container is missing 'container_name'.")
# Initialize the sanitized container
sanitized_container = {
'docker_image': container['docker_image'],
'container_name': container['container_name']
}
# Add 'cpu_shares' if present and valid
if 'cpu_shares' in container:
if not isinstance(container['cpu_shares'], (int, float)) or container['cpu_shares'] <= 0:
raise ValueError("'cpu_shares' must be a positive number.")
sanitized_container['cpu_shares'] = container['cpu_shares']
# Add 'mem_limit' if present and valid
if 'mem_limit' in container:
if not isinstance(container['mem_limit'], (int, float)) or container['mem_limit'] <= 0:
raise ValueError("'mem_limit' must be a positive number.")
sanitized_container['mem_limit'] = container['mem_limit']
# Add 'ports' if present and valid
if 'ports' in container:
if not isinstance(container['ports'], list):
raise ValueError("'ports' must be a list.")
@@ -104,13 +98,11 @@ def validate_payload(payload):
raise ValueError("Each port mapping must have 'internal' and 'external' keys.")
if not isinstance(port_mapping['internal'], int) or not isinstance(port_mapping['external'], int):
raise ValueError("'internal' and 'external' ports must be integers.")
# Handle use_dns in port mapping
port_data = {
'internal': port_mapping['internal'],
'external': port_mapping['external']
}
if 'use_dns' in port_mapping:
port_data['use_dns'] = bool(port_mapping['use_dns'])
@@ -118,8 +110,7 @@ def validate_payload(payload):
if sanitized_ports:
sanitized_container['ports'] = sanitized_ports
# Add 'networks' if present and valid
if 'networks' in container:
if isinstance(container['networks'], str):
if container['networks'].strip():
@@ -134,11 +125,28 @@ def validate_payload(payload):
if sanitized_networks:
sanitized_container['networks'] = sanitized_networks
elif container['networks'] is None:
pass # allow it to be missing or null
pass
else:
raise ValueError("'networks' must be a string or a list of strings.")
# Add the sanitized container to the sanitized payload
if 'env' in container:
if isinstance(container['env'], dict):
sanitized_env = []
for key, value in container['env'].items():
if not isinstance(key, str) or not isinstance(value, str):
raise ValueError("Environment variable keys and values must be strings.")
sanitized_env.append(f"{key}={value}")
sanitized_container['env'] = sanitized_env
elif isinstance(container['env'], list):
sanitized_env = []
for entry in container['env']:
if not isinstance(entry, str) or '=' not in entry:
raise ValueError("Each environment variable in the list must be a string in 'KEY=VALUE' format.")
sanitized_env.append(entry)
sanitized_container['env'] = sanitized_env
else:
raise ValueError("'env' must be a dictionary or a list of 'KEY=VALUE' strings.")
sanitized_payload['containers'].append(sanitized_container)
return sanitized_payload
@@ -188,7 +196,6 @@ def request_container_workload():
status=202,
)[0]
@api_bp.route('/workloads/containers/status_update/<system_container_id>', methods=['PUT'])
def update_container_workload(system_container_id):
data = request.json
@@ -231,7 +238,9 @@ def update_container_workload(system_container_id):
container = Workload.query.filter(
Workload.id == system_container_id,
or_(Workload.workload_type == "Container", Workload.workload_type == "NSController"),
Workload.deleted == False
# Commenting out the deleted filter becuase when we need handle containers being
# rebuilt we must accept that deleted containers can become not deleted again
# Workload.deleted == False
).first_or_404()
except Exception as e:
@@ -255,7 +264,9 @@ def update_container_workload(system_container_id):
container.soft_delete()
# If container status is "deleted", check if all containers in the same pod are deleted
check_deleted_container(container)
else:
container.restore()
db.session.commit()
return jsonify({"success": True, "message": "Status updated successfully."}), 200
@@ -495,8 +506,9 @@ def delete_pod(pod_id):
container_delete_requests = []
for mapping in pod.container_mappings:
container = mapping.container
container.set_status("pending-deleted")
db.session.add(container)
if not container.deleted:
container.set_status("pending-deleted")
db.session.add(container)
# Add this container to our websocket deletion request
container_delete_requests.append({
+1 -1
View File
@@ -544,7 +544,7 @@ class Workload(BaseModel):
"pending-allocated": ["*"],
"dead": ["*"],
"new": ["*"],
"deleted": ["deleted"],
"deleted": ["deleted","running"],
"pending-deleted": ["running","deleted","dead","stopped"], # A container can go from pending-deleted to running if there is a race between a deletion and a creation event
"stopped": ["deleted"]
}
+414 -306
View File
@@ -1,375 +1,483 @@
#app/utils/create_workload_container.py
import random
from flask import json, request, jsonify
import requests
from app import db, logger, app
from app.models.models import PortForwarding, Workload, WorkloadHost, WorkloadHostPortMapping, VirtualDataCenter, Label, VolumeWorkloadMapping, ContainerPod, ContainerPodContainer, CloudflareDNSRecord, CloudflareTunnel, WorkloadRequest
import uuid
from app.scheduling_filters import ExcludeAllDisabledHosts, MostAvailableCapacity, ExcludeAllOfflineHosts, DockerCapableHosts, ExcludeHostsWithoutNorthSouthIP, RandomizeHostsOrder
from app.services.cloudflare import CloudflareTunnelManager
from werkzeug.exceptions import abort
# app/utils/create_workload_container.py
"""
Create or extend container pods in xCloudify.
This module supports two workflows:
• NEW POD – payload contains 'virtual_data_center'
• ADD TO POD – payload contains 'pod'
Both workflows end by sending a full “intent-based” pod-update
payload to the WebSocket worker. The worker is idempotent and will
ignore duplicates.
"""
import random
import uuid
from datetime import datetime
from typing import List, Dict
import requests
from flask import json
from sqlalchemy.orm import joinedload
from app import db, logger, app
from app.models.models import (
Workload, WorkloadHost, WorkloadHostPortMapping, VirtualDataCenter,
ContainerPod, ContainerPodContainer, PortForwarding, CloudflareTunnel,
CloudflareDNSRecord, WorkloadRequest
)
from app.scheduling_filters import (
ExcludeHostsWithoutNorthSouthIP, DockerCapableHosts, ExcludeAllOfflineHosts,
ExcludeAllDisabledHosts, MostAvailableCapacity, RandomizeHostsOrder
)
from app.services.cloudflare import CloudflareTunnelManager
from app.utils.standard_responses import api_response
def add_container_workload(request_data):
# DONT Validate incoming data
validated_data = request_data.payload
logger.debug(validated_data)
# ─────────────────────────────────────────────────────────────────────────────
# Helper ▸ Validate a single host through the normal filter chain
# ─────────────────────────────────────────────────────────────────────────────
def ensure_pod_host_is_eligible(pod: ContainerPod) -> WorkloadHost:
"""
Re-validate the pod’s host with the standard scheduler filters.
request_vdc = VirtualDataCenter.query.filter_by(id=(validated_data['virtual_data_center'])).first()
if request_vdc is None:
logger.error(f"Virtual Data Center not found with ID: {validated_data['virtual_data_center']}")
return api_response(
success=False,
status=404,
message="Virtual Data Center not found",
error_type="RESOURCE_NOT_FOUND"
)
all_workload_hosts = WorkloadHost.query.filter_by(region_id=request_vdc.region.id).all()
logger.debug(f"Found {len(all_workload_hosts)} hosts eligible for placement")
Returns
-------
WorkloadHost
The same host if it passes every filter.
filters = [ExcludeHostsWithoutNorthSouthIP(), DockerCapableHosts(), ExcludeAllOfflineHosts(), ExcludeAllDisabledHosts(), MostAvailableCapacity(), RandomizeHostsOrder()]
filtered_hosts = all_workload_hosts
for filter in filters:
filtered_hosts = filter.apply(filtered_hosts)
Raises
------
ValueError
If the host fails ANY filter.
"""
hosts = [pod.workload_host]
if not filtered_hosts:
logger.error(f"No hosts found in region {request_vdc.region.name} to place this request after filtering")
filters = [
ExcludeHostsWithoutNorthSouthIP(),
DockerCapableHosts(), # update if supporting libvirt pods
ExcludeAllOfflineHosts(),
ExcludeAllDisabledHosts(),
MostAvailableCapacity(), # harmless on 1-item list
RandomizeHostsOrder(),
]
for f in filters:
hosts = f.apply(hosts)
if not hosts:
raise ValueError(
f"Pod host {pod.workload_host_id} failed filter "
f"{f.__class__.__name__}"
)
return hosts[0]
# ─────────────────────────────────────────────────────────────────────────────
# Public API ▸ entrypoint called by Celery worker or API view
# ─────────────────────────────────────────────────────────────────────────────
def add_container_workload(request_row: WorkloadRequest):
"""
Dispatcher: decide whether we are *creating* a new pod or
*adding* containers to an existing pod, based on the request
payload that has already passed schema validation.
Parameters
----------
request_row : WorkloadRequest
The persisted request row holding the validated payload.
"""
payload = request_row.payload
logger.debug("add_container_workload payload: %s", payload)
try:
if "pod" in payload:
_add_to_existing_pod_flow(request_row, payload)
else:
_create_new_pod_flow(request_row, payload)
except ValueError as exc:
logger.error("Workload build error: %s", exc, exc_info=True)
return api_response(
success=False,
status=400,
message=f"No hosts found in region {request_vdc.region.name} to place this request",
error_type="NO_ELIGIBLE_HOSTS"
message=str(exc),
error_type="WORKLOAD_BUILD_FAILED",
)
selected_host = filtered_hosts[0]
logger.info(f"Selected host: {selected_host}")
return None # success handled inside each flow
# --- Step 1: Create a new Pod ---
new_pod = ContainerPod(
name=f"Pod_{uuid.uuid4()}",
workload_host_id=selected_host.id,
vdc_id=request_vdc.id
# ─────────────────────────────────────────────────────────────────────────────
# Workflow ▸ create brand-new pod
# ─────────────────────────────────────────────────────────────────────────────
def _create_new_pod_flow(request_row: WorkloadRequest, payload: Dict):
"""Implements the original new-pod logic inside a helper."""
vdc: VirtualDataCenter = VirtualDataCenter.query.get(payload["virtual_data_center"])
if not vdc:
raise ValueError(f"Virtual Data Center {payload['virtual_data_center']} not found")
# 1. Host selection via filters
all_hosts = WorkloadHost.query.filter_by(region_id=vdc.region.id).all()
filters = [
ExcludeHostsWithoutNorthSouthIP(),
DockerCapableHosts(),
ExcludeAllOfflineHosts(),
ExcludeAllDisabledHosts(),
MostAvailableCapacity(),
RandomizeHostsOrder(),
]
eligible = all_hosts
for f in filters:
eligible = f.apply(eligible)
if not eligible:
raise ValueError(f"No eligible hosts in region {vdc.region.name}")
host = eligible[0]
logger.info("Selected host %s for new pod", host.id)
# 2. Build DB objects & enqueue task
_build_pod_objects_and_enqueue(
payload=payload,
vdc=vdc,
host=host,
existing_pod=None,
)
db.session.add(new_pod)
db.session.flush() # Flush so we get new_pod.id
logger.info("New-pod flow completed for request %s", request_row.id)
# --- Step 2: Create NSController for the Pod ---
new_nscontroller_name = f"NSCONTROLLER_{uuid.uuid4()}"
nscontroller_launch_params = {
"docker_image": "busybox",
"container_name": new_nscontroller_name,
"command": "sleep infinity",
"networks": [],
"ports": []
}
new_nscontroller = Workload(
name=new_nscontroller_name,
workload_type="NSController",
vdc_id=request_vdc.id,
launch_params=json.dumps(nscontroller_launch_params),
workload_host_id=selected_host.id
# ─────────────────────────────────────────────────────────────────────────────
# Workflow ▸ add containers to existing pod
# ─────────────────────────────────────────────────────────────────────────────
def _add_to_existing_pod_flow(request_row: WorkloadRequest, payload: Dict):
"""Append containers to an existing pod (intent-based update)."""
pod: ContainerPod = ContainerPod.query.options(
joinedload(ContainerPod.workload_host),
joinedload(ContainerPod.nscontroller_workload),
).get(payload["pod"])
if not pod or pod.deleted:
raise ValueError(f"Pod {payload['pod']} not found or deleted")
host = ensure_pod_host_is_eligible(pod)
vdc = pod.vdc
logger.info(
"Appending containers to pod %s on host %s (VDC %s)",
pod.id, host.id, vdc.id
)
db.session.add(new_nscontroller)
db.session.flush() # So we have new_nscontroller.id for the pod linkage
new_pod.nscontroller_workload_id = new_nscontroller.id
db.session.add(new_pod)
db.session.flush()
_build_pod_objects_and_enqueue(
payload=payload,
vdc=vdc,
host=host,
existing_pod=pod,
)
logger.info("Add-to-pod flow completed for request %s", request_row.id)
client_response = []
# Track Cloudflare configurations for bulk setup
ingress_mappings = []
dns_records_to_create = []
# ─────────────────────────────────────────────────────────────────────────────
# Shared builder ▸ writes DB rows, sends WebSocket task
# ─────────────────────────────────────────────────────────────────────────────
def _build_pod_objects_and_enqueue(
*, payload: Dict, vdc: VirtualDataCenter, host: WorkloadHost,
existing_pod: ContainerPod | None
):
"""
Create all required DB objects (Pod, Workloads, PortForwarding,
Cloudflare, etc.) and queue a “pod-update” task.
# --- Step 3: Create all Container Workloads inside the Pod ---
for _container in validated_data['containers']:
logger.debug(f"Creating container {_container['container_name']}")
Parameters
----------
payload : dict
vdc : VirtualDataCenter
host : WorkloadHost
existing_pod : ContainerPod | None
None → create pod & NSController; otherwise reuse.
"""
pod = existing_pod
nscontroller = None
# Build initial container object
new_container = Workload(
name=_container['container_name'],
# ── 1. Pod + NSController bootstrap ────────────────────────────────
if pod is None:
pod = ContainerPod(
name=f"Pod_{uuid.uuid4()}",
workload_host_id=host.id,
vdc_id=vdc.id,
)
db.session.add(pod)
db.session.flush()
ns_name = f"NSCONTROLLER_{uuid.uuid4()}"
ns_lp = {
"docker_image": "busybox",
"container_name": ns_name,
"command": "sleep infinity",
"networks": [],
"ports": [],
}
nscontroller = Workload(
name=ns_name,
workload_type="NSController",
vdc_id=vdc.id,
launch_params=json.dumps(ns_lp),
workload_host_id=host.id,
)
db.session.add(nscontroller)
db.session.flush()
pod.nscontroller_workload_id = nscontroller.id
db.session.add(pod)
else:
nscontroller = Workload.query.get(pod.nscontroller_workload_id)
ns_lp = json.loads(nscontroller.launch_params)
# Track Cloudflare details
ingress_mappings: List[Dict] = []
dns_records_to_create: List[Dict] = []
client_response: List[Dict] = []
# ── 2. Create / map containers ──────────────────────────────────────
for cdef in payload["containers"]:
container = Workload(
name=cdef["container_name"],
workload_type="Container",
vdc_id=request_vdc.id,
launch_params=json.dumps(_container),
workload_host_id=selected_host.id
vdc_id=vdc.id,
launch_params=json.dumps(cdef),
workload_host_id=host.id,
)
db.session.add(new_container)
db.session.flush() # So we have new_container.id
db.session.add(container)
db.session.flush()
# Map container to pod
mapping = ContainerPodContainer(
pod_id=new_pod.id,
container_workload_id=new_container.id
)
db.session.add(mapping)
db.session.add(ContainerPodContainer(
pod_id=pod.id,
container_workload_id=container.id
))
# --- Step 4: Handle ports and WorkloadHostPortMapping ---
if 'ports' in _container:
for port_mapping in _container['ports']:
internal_port = port_mapping['internal']
preferred_external_port = port_mapping['external']
use_dns = port_mapping.get('use_dns', False)
# Ports
for pm in cdef.get("ports", []):
int_port = pm["internal"]
pref_ext = pm["external"]
use_dns = pm.get("use_dns", False)
assigned_external_port = find_available_external_port(
selected_host.id, preferred_port=preferred_external_port
)
ext_port = find_available_external_port(host.id, pref_ext)
logger.info(f"Assigning internal port:{internal_port} -> external port:{assigned_external_port} for container {new_container.id}")
ns_lp["ports"].append({"internal": int_port, "external": ext_port})
# Add this port to NSController launch params
nscontroller_launch_params['ports'].append({
"internal": internal_port,
"external": assigned_external_port
db.session.add(WorkloadHostPortMapping(
workload_host_id=host.id,
container_workload_id=container.id,
internal_port=int_port,
external_port=ext_port,
))
pf = PortForwarding(
name="pf",
pod_id=pod.id,
internal_port=int_port,
external_port=ext_port,
protocol=pm.get("protocol", "tcp"),
ip_address=host.ip_address_northsouth,
workload_id=container.id,
)
db.session.add(pf)
db.session.flush()
if use_dns:
hostname = f"{cdef['container_name']}-{int_port}.hawkvelt.tech"
ingress_mappings.append({
"dns_hostname": hostname,
"local_ip": "127.0.0.1",
"local_port": str(int_port),
})
dns_records_to_create.append({
"hostname": hostname,
"forward_id": pf.id,
})
# Save host-level port mapping
new_port_mapping = WorkloadHostPortMapping(
workload_host_id=selected_host.id,
container_workload_id=new_container.id,
internal_port=internal_port,
external_port=assigned_external_port
)
db.session.add(new_port_mapping)
# Save the pod-level forwarding rule
new_forward = PortForwarding(
name="f",
pod_id=new_pod.id,
internal_port=internal_port,
external_port=assigned_external_port,
protocol=port_mapping.get("protocol", "tcp"),
ip_address=selected_host.ip_address_northsouth,
workload_id=new_container.id,
)
db.session.add(new_forward)
db.session.flush() # Ensure we have new_forward.id
# If use_dns is True for this port, prepare Cloudflare config
if use_dns:
hostname = f"{_container['container_name']}-{internal_port}.hawkvelt.tech"
ingress_mappings.append({
"dns_hostname": hostname,
"local_ip": "127.0.0.1",
"local_port": str(internal_port)
})
dns_records_to_create.append({
"hostname": hostname,
"forward_id": new_forward.id,
"internal_port": internal_port
})
client_response.append({
"pod_id": str(new_pod.id),
"container_id": str(new_container.id),
"container_name": new_container.name,
"pod_id": str(pod.id),
"container_id": str(container.id),
"container_name": container.name,
"container_status": "pending-allocation",
})
# Handle Cloudflare setup if we have any DNS-enabled ports
# ── 3. Cloudflare tunnel & DNS records (intent-based) ───────────────
if ingress_mappings:
logger.debug(f"Setting up Cloudflare tunnel with {len(ingress_mappings)} routes")
api_token = "ri6lIjM-aRJBY_xZ82w0Haew93U6YgZYHi5jby1-"
account_id = "5095a74b62fee53cc5d997c67443bac5"
zone_id = "e2cafdd8929869d5db885d8824f514b5"
cloudflare_manager = CloudflareTunnelManager(api_token, account_id, zone_id, logger)
# Create a single tunnel with multiple ingress rules
tunnel_name = f"tun-{new_pod.id}"
cloudflare_response = cloudflare_manager.setup_tunnel(
tunnel_name=tunnel_name,
ingress_mappings=ingress_mappings
)
logger.debug(cloudflare_response)
cloudflare_tunnel_token = cloudflare_response['token']
# Create tunnel record in our database
new_cloudflare_tunnel = CloudflareTunnel(
token=cloudflare_tunnel_token,
account_id=account_id,
tunnel_id=cloudflare_response['tunnel_id'],
name=cloudflare_response['tunnel_name'],
tunnel_secret=cloudflare_response['tunnel_secret'],
nscontroller_workload_id=new_nscontroller.id
)
db.session.add(new_cloudflare_tunnel)
db.session.flush()
logger.debug(dns_records_to_create)
# Create DNS records for each hostname
for dns_record in dns_records_to_create:
# Create DNS record in Cloudflare
# dns_response = cloudflare_manager.create_dns_record(
# hostname=dns_record['hostname'],
# target=f"{cloudflare_response['tunnel_id']}.cfargotunnel.com"
# )
# TODO - Make sure we are connecting the DNS record in the DB t whjats in Cloudflare so we can delete them later
# The commented code above is not required becuase the DNS records are created with the tunnel, so we need to query
# the setup_tunnel reponse for the DNs info...maybe?
# Create DNS record in our database
for _createdDNSRecord in cloudflare_response['dns_records_created']:
logger.debug(_createdDNSRecord)
if _createdDNSRecord['hostname']==dns_record['hostname']:
logger.debug(f"Match {_createdDNSRecord['hostname']}")
dns_record_id=_createdDNSRecord['response']['id']
api_token = app.config["CLOUDFLARE_API_TOKEN"]
account_id = app.config["CLOUDFLARE_ACCOUNT_ID"]
zone_id = app.config["CLOUDFLARE_ZONE_ID"]
cf_mgr = CloudflareTunnelManager(api_token, account_id, zone_id, logger)
new_dns_record = CloudflareDNSRecord(
tunnel = CloudflareTunnel.query.filter_by(
nscontroller_workload_id=nscontroller.id
).first()
if tunnel:
cf_rsp = cf_mgr.setup_tunnel(
tunnel_name=tunnel.name,
ingress_mappings=ingress_mappings,
existing_tunnel_id=tunnel.tunnel_id,
)
else:
cf_rsp = cf_mgr.setup_tunnel(
tunnel_name=f"tun-{pod.id}",
ingress_mappings=ingress_mappings,
)
tunnel = CloudflareTunnel(
token=cf_rsp["token"],
account_id=account_id,
tunnel_id=cf_rsp["tunnel_id"],
name=cf_rsp["tunnel_name"],
tunnel_secret=cf_rsp["tunnel_secret"],
nscontroller_workload_id=nscontroller.id,
)
db.session.add(tunnel)
db.session.flush()
# cloudflared side-car
side_name = f"cloudflared-sidecar-{pod.id}"
side_lp = {
"docker_image": "cloudflare/cloudflared:latest",
"container_name": side_name,
"command": f"tunnel --no-autoupdate run --token {cf_rsp['token']}",
"network_mode": f"container:{nscontroller.name}",
"restart_policy": "always",
}
sidecar = Workload(
name=side_name,
workload_type="Container",
vdc_id=vdc.id,
launch_params=json.dumps(side_lp),
workload_host_id=host.id,
)
db.session.add(sidecar)
db.session.flush()
db.session.add(ContainerPodContainer(
pod_id=pod.id,
container_workload_id=sidecar.id
))
# DNS records
for rec in dns_records_to_create:
dns_id = next(
(d["response"]["id"] for d in cf_rsp["dns_records_created"]
if d["hostname"] == rec["hostname"]),
None
)
dns_row = CloudflareDNSRecord(
name="dns",
zone_id=zone_id,
dns_record_id=dns_record_id,
hostname=dns_record['hostname'],
dns_record_id=dns_id,
hostname=rec["hostname"],
record_type="CNAME",
content=f"{cloudflare_response['tunnel_id']}.cfargotunnel.com",
tunnel_id=new_cloudflare_tunnel.id
content=f"{tunnel.tunnel_id}.cfargotunnel.com",
tunnel_id=tunnel.id,
)
db.session.add(new_dns_record)
db.session.add(dns_row)
db.session.flush()
logger.debug(f"Created DNS Record - {hostname}")
# Link the port forwarding to the DNS record
port_forward = PortForwarding.query.filter_by(id=dns_record['forward_id']).first()
if port_forward:
logger.debug("Found the associated PortForwarding in the DB")
port_forward.dns_record_id = new_dns_record.id
db.session.add(port_forward)
logger.debug(f"Linked DNS record to PortForwarding {dns_record['forward_id']}")
# link back to PortForwarding
pf = PortForwarding.query.get(rec["forward_id"])
pf.dns_record_id = dns_row.id
db.session.add(pf)
# Create cloudflared sidecar container
sidecar_container_name = f"cloudflared-sidecar-{new_pod.id}"
sidecar_launch_params = {
"docker_image": "cloudflare/cloudflared:latest",
"container_name": sidecar_container_name,
"command": f"tunnel --no-autoupdate run --token {cloudflare_tunnel_token}",
"network_mode": f"container:{new_nscontroller.name}",
"restart_policy": "always"
}
sidecar_workload = Workload(
name=sidecar_container_name,
workload_type="Container",
vdc_id=request_vdc.id,
launch_params=json.dumps(sidecar_launch_params),
workload_host_id=selected_host.id
)
db.session.add(sidecar_workload)
db.session.flush()
sidecar_mapping = ContainerPodContainer(
pod_id=new_pod.id,
container_workload_id=sidecar_workload.id
)
db.session.add(sidecar_mapping)
client_response.append({
"pod_id": str(new_pod.id),
"container_id": str(sidecar_workload.id),
"container_name": sidecar_workload.name,
"container_status": "pending-allocation",
"cloudflare_tunnel": {
"tunnel_id": new_cloudflare_tunnel.tunnel_id,
"hostnames": [dns['hostname'] for dns in dns_records_to_create]
}
})
# --- Step 5: Update NSController launch params with final ports ---
new_nscontroller.launch_params = json.dumps(nscontroller_launch_params)
db.session.add(new_nscontroller)
# --- Step 6: Commit everything ---
# ── 4. Finalise NSController launch params & commit ─────────────────
nscontroller.launch_params = json.dumps(ns_lp)
db.session.add(nscontroller)
db.session.commit()
# --- Step 7: Build and Send Full Pod Payload ---
pod_payload = build_pod_payload(new_pod)
headers = {"Content-Type": "application/json"}
logger.debug(f"Sending full pod payload to websocket server: {pod_payload}")
# ── 5. Send full intent-based pod payload to WebSocket ─────────────
pod_payload = build_pod_payload(pod)
logger.debug("Pod payload ➜ WebSocket: %s", pod_payload)
resp = requests.post(
app.config["WEBSOCKET_SERVER_URL"],
data=json.dumps(pod_payload),
headers={"Content-Type": "application/json"},
)
websocket_server_response = requests.post(app.config['WEBSOCKET_SERVER_URL'], data=json.dumps(pod_payload), headers=headers)
if websocket_server_response.status_code == 201:
logger.info(f"Task created successfully for pod {new_pod.id}")
for container_info in client_response:
container = Workload.query.get((container_info['container_id']))
container.set_status("allocated")
db.session.add(container)
new_nscontroller.set_status("allocated")
db.session.add(new_nscontroller)
db.session.commit()
if resp.status_code == 201:
logger.info("Pod-update queued for pod %s", pod.id)
_set_container_statuses(client_response, "allocated")
nscontroller.set_status("allocated")
else:
logger.error(f"Pod creation request failed for pod {new_pod.id}")
for container_info in client_response:
container = Workload.query.get((container_info['container_id']))
container.set_status("failed-allocation")
db.session.add(container)
new_nscontroller.set_status("failed-allocation")
db.session.add(new_nscontroller)
db.session.commit()
logger.error("Pod-update failed for pod %s · %s", pod.id, resp.text)
_set_container_statuses(client_response, "failed-allocation")
nscontroller.set_status("failed-allocation")
db.session.add(nscontroller)
db.session.commit()
def build_pod_payload(pod):
nscontroller = Workload.query.get(pod.nscontroller_workload_id)
container_mappings = ContainerPodContainer.query.filter_by(pod_id=pod.id).all()
# ─────────────────────────────────────────────────────────────────────────────
# Intent payload builder ▸ always list every container
# ─────────────────────────────────────────────────────────────────────────────
def build_pod_payload(pod: ContainerPod) -> Dict:
"""
Build a complete “intent” payload for `pod-update` that enumerates
*all* containers mapped to the pod. The worker is responsible for
deduplication/idempotency.
"""
ns = Workload.query.get(pod.nscontroller_workload_id)
ns_lp = json.loads(ns.launch_params)
containers = []
for mapping in container_mappings:
container = Workload.query.get(mapping.container_workload_id)
launch_params = json.loads(container.launch_params)
container_payload = {
"container_id": str(container.id),
"docker_image": launch_params.get('docker_image'),
"container_name": container.name,
"desired_state": "running"
for m in ContainerPodContainer.query.filter_by(pod_id=pod.id).all():
w = Workload.query.get(m.container_workload_id)
launch_params = json.loads(w.launch_params)
cont = {
"container_id": str(w.id),
"docker_image": launch_params["docker_image"],
"container_name": w.name,
"desired_state": "running",
}
# If a command is defined in the launch params, add it
if 'command' in launch_params and launch_params['command']:
container_payload['command'] = launch_params['command']
containers.append(container_payload)
if launch_params.get("command"):
cont["command"] = launch_params["command"]
if launch_params.get("env"):
cont["env"] = launch_params["env"]
nscontroller_launch_params = json.loads(nscontroller.launch_params)
containers.append(cont)
pod_payload = {
return {
"worker_id": str(pod.workload_host_id),
"task_type": "pod-update",
"job_details": {
"pod_id": str(pod.id),
"nscontroller": {
"container_id": str(nscontroller.id),
"docker_image": nscontroller_launch_params.get('docker_image'),
"command": nscontroller_launch_params.get('command', 'sleep infinity'),
"networks": nscontroller_launch_params.get('networks', []),
"ports": nscontroller_launch_params.get('ports', [])
"container_id": str(ns.id),
"docker_image": ns_lp["docker_image"],
"command": ns_lp.get("command", "sleep infinity"),
"networks": ns_lp.get("networks", []),
"ports": ns_lp.get("ports", []),
},
"containers": containers
}
"containers": containers,
},
}
return pod_payload
# ─────────────────────────────────────────────────────────────────────────────
# Utility helpers
# ─────────────────────────────────────────────────────────────────────────────
def find_available_external_port(host_id: str, preferred_port: int | None = None,
port_range=(30000, 40000)) -> int:
"""Return an unused external port on *host_id* (tries preferred first)."""
used = {m.external_port for m in
WorkloadHostPortMapping.query.filter_by(workload_host_id=host_id).all()}
def find_available_external_port(host_id, preferred_port=None, port_range=(30000, 40000)):
"""Find a free external port on a given host."""
existing_ports = set(
mapping.external_port for mapping in WorkloadHostPortMapping.query.filter_by(workload_host_id=host_id).all()
)
if preferred_port and preferred_port not in existing_ports:
if preferred_port and preferred_port not in used:
return preferred_port
attempts = 0
while attempts < 1000:
port_candidate = random.randint(port_range[0], port_range[1])
if port_candidate not in existing_ports:
return port_candidate
attempts += 1
for _ in range(1000):
p = random.randint(*port_range)
if p not in used:
return p
raise RuntimeError("No free external port after 1000 attempts")
raise Exception("No available ports found after 1000 attempts")
def _set_container_statuses(client_resp: List[Dict], new_status: str):
"""Bulk update container statuses with audit logging."""
for c in client_resp:
w = Workload.query.get(c["container_id"])
w.set_status(new_status)
db.session.add(w)
+89 -38
View File
@@ -118,6 +118,8 @@ class ContainerTask:
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"]
try:
# Step 1 – create & start
@@ -161,67 +163,116 @@ class ContainerTask:
self.logger.warning(f"Failed to remove container {container.name}: {str(e)}")
def is_container_correct(self, container, container_spec):
"""Check if a running container matches its spec (image, networks, ports, env vars, volumes)."""
"""
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.
"""
try:
config = container.attrs.get('Config', {})
host_config = container.attrs.get('HostConfig', {})
config = container.attrs.get('Config', {})
host_config = container.attrs.get('HostConfig', {})
network_settings = container.attrs.get('NetworkSettings', {})
# === 1. Check Image ===
# ─────────────────── 1️⃣ IMAGE ───────────────────
expected_image = container_spec.get("docker_image")
running_image = container.image.tags[0] if container.image.tags else None
running_image = container.image.tags[0] if container.image.tags else None
if expected_image != running_image:
self.logger.warning(f"Image mismatch: expected {expected_image}, got {running_image}")
self.logger.warning(
"Image mismatch: expected %s, got %s", expected_image, running_image
)
return False
# === 2. Check Networks (only for NSController) ===
# ─────────────────── 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(f"Network mismatch: expected {expected_networks}, got {existing_networks}")
self.logger.warning(
"Network mismatch: expected %s, got %s", expected_networks, existing_networks
)
return False
# === 3. Check Ports (only for NSController) ===
# ─────────────────── 3️⃣ PORTS (NSController only) ───────────────────
if container_spec.get("ports") and container_spec.get("nscontroller_container_id") is None:
bound_ports = set()
ports_config = network_settings.get('Ports', {})
for port_proto, mappings in ports_config.items():
if mappings:
for mapping in mappings:
external_port = int(mapping['HostPort'])
internal_port = int(port_proto.split('/')[0])
bound_ports.add((internal_port, external_port))
expected_ports = set((p['internal'], p['external']) for p in container_spec['ports'])
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']}
if not expected_ports.issubset(bound_ports):
self.logger.warning(f"Ports mismatch: expected {expected_ports}, got {bound_ports}")
self.logger.warning(
"Ports mismatch: expected %s, got %s", expected_ports, bound_ports
)
return False
# === 4. Check Env Vars ===
expected_env = container_spec.get("env", {})
if expected_env:
running_env = {item.split("=")[0]: item.split("=")[1] for item in config.get('Env', []) if "=" in item}
for key, value in expected_env.items():
if running_env.get(key) != str(value):
self.logger.warning(f"Env var mismatch: {key} expected {value}, got {running_env.get(key)}")
return False
# ─────────────────── 4️⃣ ENVIRONMENT VARIABLES ───────────────────
def _list_to_dict(env_list: List[str]) -> dict[str, str]:
"""Convert ['KEY=VALUE', …] → {'KEY': 'VALUE', …} safely."""
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
# === 5. Check Volumes ===
# 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', []))
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️⃣ VOLUMES ───────────────────
expected_volumes = container_spec.get("volumes", [])
if expected_volumes:
running_mounts = container.attrs.get("Mounts", [])
running_volume_paths = set(m['Destination'] for m in running_mounts)
expected_volume_paths = set(v['destination'] for v in expected_volumes if 'destination' in v)
if not expected_volume_paths.issubset(running_volume_paths):
self.logger.warning(f"Volume mismatch: expected mounts {expected_volume_paths}, got {running_volume_paths}")
running_mounts = container.attrs.get("Mounts", [])
running_volume_paths = {m['Destination'] for m in running_mounts}
expected_paths = {v['destination'] for v in expected_volumes
if isinstance(v, dict) and 'destination' in v}
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
return True
# ─────────────────── ✅ ALL CHECKS PASSED ───────────────────
return True
except Exception as e:
self.logger.warning(f"Error validating container {container.name}: {str(e)}")
except Exception as exc:
self.logger.warning(
"Error validating container %s: %s", container.name, exc, exc_info=True
)
return False
def shutdown(self):