Merge branch 'main' into v1.01/internal-auth
This commit is contained in:
@@ -61,6 +61,7 @@ def validate_network_data(data, is_update=False):
|
||||
'ipv6_cidr': str,
|
||||
'ipv6_gateway': str,
|
||||
'ovs_bridge': str,
|
||||
'encapsulation': str,
|
||||
}
|
||||
for field, field_type in optional_fields.items():
|
||||
if field in data and not isinstance(data[field], field_type):
|
||||
@@ -137,7 +138,8 @@ def add_network():
|
||||
nscontroller_ip=validated_data.get('nscontroller_ip'),
|
||||
ipv6_cidr=validated_data.get('ipv6_cidr'),
|
||||
ipv6_gateway=validated_data.get('ipv6_gateway'),
|
||||
ovs_bridge=validated_data.get('ovs_bridge')
|
||||
ovs_bridge=validated_data.get('ovs_bridge'),
|
||||
**({'encapsulation': validated_data['encapsulation']} if 'encapsulation' in validated_data else {}),
|
||||
)
|
||||
db.session.add(instance)
|
||||
db.session.commit()
|
||||
|
||||
@@ -78,8 +78,7 @@ def delete_pod(pod_id):
|
||||
- If the pod has no assigned host (unscheduled), soft-delete the pod and associated resources locally
|
||||
without attempting to build a worker payload.
|
||||
|
||||
PROTECTION: Pods containing NSControllers cannot be deleted.
|
||||
NSControllers provide critical networking services and must be manually removed from the database if necessary.
|
||||
The pod's NSController is deleted along with the rest of the pod.
|
||||
"""
|
||||
try:
|
||||
pod_uuid = pod_id
|
||||
@@ -94,16 +93,6 @@ def delete_pod(pod_id):
|
||||
|
||||
try:
|
||||
pod = ContainerPod.query.filter_by(id=pod_uuid).first_or_404()
|
||||
ns_controller = Workload.query.filter_by(pod_id=pod_id, workload_type="NSController", deleted=False).first()
|
||||
if ns_controller:
|
||||
logger.warning(f"Attempted to delete pod {pod_id} that contains NSController {ns_controller.id} - this is protected")
|
||||
return api_response(
|
||||
success=False,
|
||||
message="Cannot delete pods containing NSController. NSControllers provide critical networking services.",
|
||||
status=403,
|
||||
error_type="NSCONTROLLER_PROTECTED",
|
||||
error_details={"pod_id": str(pod.id), "nscontroller_id": str(ns_controller.id)}
|
||||
)
|
||||
|
||||
# Delete host port mappings (defensive; unscheduled pods typically won't have any)
|
||||
all_containers_in_pod=Workload.query.filter_by(pod_id=pod.id).all()
|
||||
@@ -233,15 +222,9 @@ def process_pod_deletion(pod: ContainerPod, containers_to_delete: list[str], del
|
||||
cid = container["container_id"]
|
||||
if cid in containers_to_delete:
|
||||
workload = Workload.query.get(cid)
|
||||
if workload and workload.workload_type == "NSController":
|
||||
logger.warning(
|
||||
"SAFETY BLOCK: refusing pod deletion processing for NSController %s",
|
||||
cid,
|
||||
)
|
||||
continue
|
||||
|
||||
container["desired_state"] = "deleted"
|
||||
|
||||
|
||||
cleanup_workload_related_records(cid, logger=logger)
|
||||
|
||||
api_token = app.config["CLOUDFLARE_API_TOKEN"]
|
||||
@@ -275,26 +258,14 @@ def process_pod_deletion(pod: ContainerPod, containers_to_delete: list[str], del
|
||||
# Collect network ports for SDN update
|
||||
container_networks = collect_and_cleanup_network_ports(cid)
|
||||
all_networks.extend(container_networks)
|
||||
|
||||
mark_workload_pending_deleted(workload)
|
||||
|
||||
# Optionally mark pod for deletion. Never delete protected NSControllers.
|
||||
if delete_pod:
|
||||
ns_controller = Workload.query.filter_by(
|
||||
pod_id=pod.id,
|
||||
workload_type="NSController",
|
||||
deleted=False,
|
||||
).first()
|
||||
|
||||
if ns_controller:
|
||||
logger.warning(
|
||||
f"Skipping pod pending-deleted transition for pod {pod.id} because "
|
||||
f"protected NSController {ns_controller.id} is present."
|
||||
)
|
||||
else:
|
||||
pod.status = "pending-deleted"
|
||||
db.session.add(pod)
|
||||
|
||||
mark_workload_pending_deleted(workload)
|
||||
|
||||
# Optionally mark pod for deletion
|
||||
if delete_pod:
|
||||
pod.status = "pending-deleted"
|
||||
db.session.add(pod)
|
||||
|
||||
db.session.commit()
|
||||
|
||||
# Trigger SDN updates to remove network ports
|
||||
|
||||
@@ -384,6 +384,7 @@ class WorkloadHost(BaseModel):
|
||||
"port_name": port.name,
|
||||
"ratelimit_out": 1000,
|
||||
"is_nscontroller": workload.workload_type == "NSController",
|
||||
"encapsulation": (port.network.encapsulation if port.network else None) or "VXLAN",
|
||||
}
|
||||
if port.network and port.network.ipv4_gateway:
|
||||
net_payload["gateway_ip"] = port.network.ipv4_gateway
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
from abc import ABC, abstractmethod
|
||||
from typing import Any, Dict, List, Optional
|
||||
from app.models.models import WorkloadHost
|
||||
from app.models.models import WorkloadHost, WorkloadHostOVSBridge
|
||||
from logger import logger
|
||||
import ipaddress
|
||||
import random
|
||||
@@ -704,3 +704,68 @@ class AffinityFilter(SchedulingFilter):
|
||||
)
|
||||
|
||||
return filtered_hosts
|
||||
|
||||
class OVSBridgeFilter(SchedulingFilter):
|
||||
"""
|
||||
Keep only hosts that expose every OVS bridge this workload's networks need.
|
||||
|
||||
"""
|
||||
|
||||
def __init__(self, required_bridges):
|
||||
if isinstance(required_bridges, str):
|
||||
required_bridges = [required_bridges]
|
||||
self.required_bridges = list(dict.fromkeys(b for b in (required_bridges or []) if b))
|
||||
|
||||
def apply(self, hosts, requirements=None):
|
||||
if not self.required_bridges or not hosts:
|
||||
logger.debug(
|
||||
f"Skipping OVSBridgeFilter (required bridges: {self.required_bridges}, "
|
||||
f"input hosts: {len(hosts)})"
|
||||
)
|
||||
return hosts
|
||||
|
||||
from app import db
|
||||
|
||||
host_ids = [str(host.id) for host in hosts]
|
||||
logger.debug(
|
||||
f"Starting OVSBridgeFilter. Required bridges: {self.required_bridges}\n"
|
||||
f"Input hosts ({len(hosts)}): {host_ids}"
|
||||
)
|
||||
|
||||
rows = (
|
||||
db.session.query(
|
||||
WorkloadHostOVSBridge.workload_host_id,
|
||||
WorkloadHostOVSBridge.ovs_bridge_name,
|
||||
)
|
||||
.filter(
|
||||
WorkloadHostOVSBridge.workload_host_id.in_(host_ids),
|
||||
WorkloadHostOVSBridge.ovs_bridge_name.in_(self.required_bridges),
|
||||
WorkloadHostOVSBridge.deleted == False,
|
||||
)
|
||||
.all()
|
||||
)
|
||||
|
||||
bridges_by_host = {}
|
||||
for workload_host_id, ovs_bridge_name in rows:
|
||||
bridges_by_host.setdefault(str(workload_host_id), set()).add(ovs_bridge_name)
|
||||
|
||||
required = set(self.required_bridges)
|
||||
filtered_hosts = []
|
||||
for host in hosts:
|
||||
host_bridges = bridges_by_host.get(str(host.id), set())
|
||||
if required.issubset(host_bridges):
|
||||
filtered_hosts.append(host)
|
||||
else:
|
||||
logger.debug(
|
||||
f"Host {host.id} ({host.hostname}) rejected by OVSBridgeFilter: "
|
||||
f"missing bridge(s) {sorted(required - host_bridges)}"
|
||||
)
|
||||
|
||||
filtered_ids = [str(host.id) for host in filtered_hosts]
|
||||
logger.debug(
|
||||
f"Finished OVSBridgeFilter. "
|
||||
f"Input: {len(hosts)}, Filtered out: {len(hosts) - len(filtered_hosts)}, Remaining: {len(filtered_hosts)}\n"
|
||||
f"Removed hosts: {list(set(host_ids) - set(filtered_ids))}\n"
|
||||
f"Remaining hosts: {filtered_ids}"
|
||||
)
|
||||
return filtered_hosts
|
||||
|
||||
@@ -128,56 +128,37 @@ def process_container_deletion(self, container_id: str) -> None:
|
||||
]
|
||||
|
||||
if len(regular_containers) == 0:
|
||||
# No regular containers left. If a protected NSController exists in this pod,
|
||||
# we must NOT delete the pod or NSController.
|
||||
nscontroller = Workload.query.filter_by(
|
||||
workload_type="NSController",
|
||||
# No regular containers left: delete the entire pod, including the
|
||||
# NSController and any other sidecars.
|
||||
should_delete_entire_pod = True
|
||||
logger.info(f"Last regular container deleted from pod {pod.id}, marking entire pod for deletion")
|
||||
|
||||
# Mark all sidecars for deletion (the NSController is a sidecar)
|
||||
sidecars = Workload.query.filter_by(
|
||||
pod_id=pod.id,
|
||||
deleted=False,
|
||||
).first()
|
||||
|
||||
if nscontroller:
|
||||
should_delete_entire_pod = False
|
||||
logger.warning(
|
||||
f"Last regular container deleted from pod {pod.id}, but NSController {nscontroller.id} "
|
||||
f"is protected. Skipping pod/NSController deletion."
|
||||
)
|
||||
AuditEntry.log_event(
|
||||
object=pod,
|
||||
action="pod_deletion_skipped_nscontroller_protected",
|
||||
description=f"Skipped pod deletion because protected NSController {nscontroller.id} is present"
|
||||
)
|
||||
else:
|
||||
# No regular containers and no NSController left: delete the entire pod.
|
||||
should_delete_entire_pod = True
|
||||
logger.info(f"Last regular container deleted from pod {pod.id}, marking entire pod for deletion")
|
||||
|
||||
# Mark all sidecars for deletion
|
||||
sidecars = Workload.query.filter_by(
|
||||
pod_id=pod.id,
|
||||
is_sidecar=True
|
||||
).all()
|
||||
for sidecar in sidecars:
|
||||
mark_workload_pending_deleted(sidecar)
|
||||
containers_to_delete.append(str(sidecar.id))
|
||||
|
||||
AuditEntry.log_event(
|
||||
object=sidecar,
|
||||
action="sidecar_deletion_queued",
|
||||
description=f"Sidecar marked for deletion as part of pod deletion"
|
||||
)
|
||||
|
||||
# Mark the pod itself for deletion
|
||||
pod.status = "pending-deleted"
|
||||
db.session.add(pod)
|
||||
is_sidecar=True
|
||||
).all()
|
||||
for sidecar in sidecars:
|
||||
mark_workload_pending_deleted(sidecar)
|
||||
containers_to_delete.append(str(sidecar.id))
|
||||
|
||||
AuditEntry.log_event(
|
||||
object=pod,
|
||||
action="pod_deletion_queued",
|
||||
description=f"Pod marked for deletion as no regular containers remain"
|
||||
object=sidecar,
|
||||
action="sidecar_deletion_queued",
|
||||
description=f"Sidecar marked for deletion as part of pod deletion"
|
||||
)
|
||||
|
||||
db.session.commit()
|
||||
# Mark the pod itself for deletion
|
||||
pod.status = "pending-deleted"
|
||||
db.session.add(pod)
|
||||
|
||||
AuditEntry.log_event(
|
||||
object=pod,
|
||||
action="pod_deletion_queued",
|
||||
description=f"Pod marked for deletion as no regular containers remain"
|
||||
)
|
||||
|
||||
db.session.commit()
|
||||
|
||||
# Build pod-update payload reflecting current database state
|
||||
if should_delete_entire_pod:
|
||||
|
||||
@@ -23,6 +23,65 @@ from app.utils.failover_utils import (
|
||||
evaluate_region_pause_state
|
||||
)
|
||||
|
||||
|
||||
def _get_running_vms_on_host(host):
|
||||
"""
|
||||
Return non-deleted VirtualMachine workloads currently placed on this host.
|
||||
"""
|
||||
return Workload.query.filter(
|
||||
and_(
|
||||
Workload.workload_host_id == host.id,
|
||||
Workload.workload_type == "VirtualMachine",
|
||||
Workload.deleted == False
|
||||
)
|
||||
).all()
|
||||
|
||||
|
||||
def _hold_vms_on_host(host, reason):
|
||||
"""
|
||||
Mark running VMs on a failed host as waiting-host while preserving their
|
||||
host assignment.
|
||||
|
||||
"""
|
||||
vms = _get_running_vms_on_host(host)
|
||||
for vm in vms:
|
||||
if vm.status == "running":
|
||||
vm.set_status("waiting-host")
|
||||
vm.last_error_reason = reason
|
||||
# Do NOT clear workload_host_id - preserve association for recovery
|
||||
db.session.add(vm)
|
||||
AuditEntry.log_event(
|
||||
object=vm,
|
||||
action="vm_host_hold",
|
||||
description=f"VM {vm.id} held on failed host {host.id} ({reason})"
|
||||
)
|
||||
return vms
|
||||
|
||||
|
||||
def _reallocate_vms_on_host(host):
|
||||
"""
|
||||
Prepare running VMs on a failed host for relocation.
|
||||
|
||||
Clears the host assignment and resets retry counters so that
|
||||
retry_waiting_host_vms performs fresh placement (and dispatch) on its next tick.
|
||||
"""
|
||||
vms = _get_running_vms_on_host(host)
|
||||
for vm in vms:
|
||||
if vm.status == "running":
|
||||
vm.set_status("waiting-host")
|
||||
vm.workload_host_id = None
|
||||
vm.allocation_attempts = 0 # Reset for fresh retry
|
||||
vm.next_retry_at = None # Trigger immediate retry
|
||||
vm.last_error_reason = "host_failure_reschedule"
|
||||
db.session.add(vm)
|
||||
AuditEntry.log_event(
|
||||
object=vm,
|
||||
action="vm_host_failure_reschedule",
|
||||
description=f"VM {vm.id} prepared for rescheduling due to host {host.id} failure"
|
||||
)
|
||||
return vms
|
||||
|
||||
|
||||
def _reconcile_offline_host(host, region, log_primary_offline_event):
|
||||
if log_primary_offline_event:
|
||||
delay = region.host_failure_reallocation_delay
|
||||
@@ -44,7 +103,6 @@ def _reconcile_offline_host(host, region, log_primary_offline_event):
|
||||
|
||||
# Check if region should be paused
|
||||
paused, pause_reason = should_region_be_paused(region)
|
||||
|
||||
if paused:
|
||||
# Region is in pause state - suppress reallocation
|
||||
logger.warning(f"Region {region.id} is paused (reason: {pause_reason}), suppressing reallocation for host {host.id}")
|
||||
@@ -90,7 +148,10 @@ def _reconcile_offline_host(host, region, log_primary_offline_event):
|
||||
action="reallocation_suppressed",
|
||||
description=f"Pod {pod.id} reallocation suppressed due to region pause ({pause_reason})"
|
||||
)
|
||||
|
||||
|
||||
# Hold VMs on the failed host (retry task skips paused regions)
|
||||
vms = _hold_vms_on_host(host, f"region_paused_{pause_reason}")
|
||||
|
||||
db.session.commit()
|
||||
|
||||
# Audit the pause decision
|
||||
@@ -104,7 +165,8 @@ def _reconcile_offline_host(host, region, log_primary_offline_event):
|
||||
details={
|
||||
"pause_reason": pause_reason,
|
||||
"failure_count": failure_count,
|
||||
"affected_pods": len(pods)
|
||||
"affected_pods": len(pods),
|
||||
"affected_vms": len(vms)
|
||||
}
|
||||
)
|
||||
|
||||
@@ -148,7 +210,10 @@ def _reconcile_offline_host(host, region, log_primary_offline_event):
|
||||
action="single_host_hold",
|
||||
description=f"Pod {pod.id} held on failed host due to single-host region"
|
||||
)
|
||||
|
||||
|
||||
# Hold VMs on the failed host (no alternate host to relocate to)
|
||||
vms = _hold_vms_on_host(host, "single_host_failure")
|
||||
|
||||
db.session.commit()
|
||||
|
||||
# Audit the single-host decision
|
||||
@@ -162,6 +227,7 @@ def _reconcile_offline_host(host, region, log_primary_offline_event):
|
||||
details={
|
||||
"failure_count": failure_count,
|
||||
"affected_pods": len(pods),
|
||||
"affected_vms": len(vms),
|
||||
"active_host_count": 1
|
||||
}
|
||||
)
|
||||
@@ -203,7 +269,10 @@ def _reconcile_offline_host(host, region, log_primary_offline_event):
|
||||
action="relocation_disabled",
|
||||
description=f"Pod {pod.id} held due to region.relocate_on_host_failure=False"
|
||||
)
|
||||
|
||||
|
||||
# Hold VMs on the failed host (retry task skips relocation-disabled regions)
|
||||
vms = _hold_vms_on_host(host, "relocation_disabled")
|
||||
|
||||
db.session.commit()
|
||||
|
||||
# Audit the disabled relocation decision
|
||||
@@ -216,7 +285,8 @@ def _reconcile_offline_host(host, region, log_primary_offline_event):
|
||||
decision="relocation_disabled",
|
||||
details={
|
||||
"failure_count": failure_count,
|
||||
"affected_pods": len(pods)
|
||||
"affected_pods": len(pods),
|
||||
"affected_vms": len(vms)
|
||||
}
|
||||
)
|
||||
|
||||
@@ -267,7 +337,10 @@ def _reconcile_offline_host(host, region, log_primary_offline_event):
|
||||
action="no_online_hosts",
|
||||
description=f"Pod {pod.id} held due to no online hosts in region {region.id}"
|
||||
)
|
||||
|
||||
|
||||
# Hold VMs on the failed host (no online host to relocate to)
|
||||
vms = _hold_vms_on_host(host, "no_online_hosts")
|
||||
|
||||
db.session.commit()
|
||||
|
||||
# Audit the no online hosts decision
|
||||
@@ -281,6 +354,7 @@ def _reconcile_offline_host(host, region, log_primary_offline_event):
|
||||
details={
|
||||
"failure_count": failure_count,
|
||||
"affected_pods": len(pods),
|
||||
"affected_vms": len(vms),
|
||||
"online_host_count": 0
|
||||
}
|
||||
)
|
||||
@@ -386,9 +460,11 @@ def _reconcile_offline_host(host, region, log_primary_offline_event):
|
||||
action=f"host_failure_{decision}",
|
||||
description=f"Pod {pod.id} {decision} due to host {host.id} failure"
|
||||
)
|
||||
|
||||
|
||||
reallocated_vms = _reallocate_vms_on_host(host)
|
||||
|
||||
db.session.commit()
|
||||
|
||||
|
||||
# Audit the reallocation decisions
|
||||
log_failure_event(
|
||||
region=region,
|
||||
@@ -396,15 +472,19 @@ def _reconcile_offline_host(host, region, log_primary_offline_event):
|
||||
event_type="reallocation_evaluation",
|
||||
window_count=failure_count,
|
||||
threshold=region.max_host_failures,
|
||||
decision="proceed_to_reschedule" if reallocation_decisions else "mixed_outcomes",
|
||||
decision="proceed_to_reschedule" if (reallocation_decisions or reallocated_vms) else "mixed_outcomes",
|
||||
details={
|
||||
"failure_count": failure_count,
|
||||
"affected_pods": len(pods),
|
||||
"affected_vms": len(reallocated_vms),
|
||||
"reallocation_decisions": reallocation_decisions
|
||||
}
|
||||
)
|
||||
|
||||
logger.info(f"Reallocation evaluation complete for host {host.id}: {len(reallocation_decisions)} pods prepared for rescheduling")
|
||||
|
||||
logger.info(
|
||||
f"Reallocation evaluation complete for host {host.id}: "
|
||||
f"{len(reallocation_decisions)} pods and {len(reallocated_vms)} VMs prepared for rescheduling"
|
||||
)
|
||||
|
||||
@celery.task(name="tasks.handle_host_downtime", bind=True, max_retries=1)
|
||||
def handle_host_downtime(self, host_id: str) -> None:
|
||||
|
||||
@@ -66,6 +66,7 @@ def retry_waiting_host_vms(self) -> None:
|
||||
ExcludeAllOfflineHosts, ExcludeAllDisabledHosts,
|
||||
LibvirtCapableHosts, MostAvailableCapacity, RandomizeHostsOrder,
|
||||
HasSufficientResources, GPUCapableHosts, ExcludeHosts,
|
||||
OVSBridgeFilter,
|
||||
)
|
||||
from app.compute.controller.server_group_affinity import apply_affinity_filters
|
||||
|
||||
@@ -131,6 +132,27 @@ def retry_waiting_host_vms(self) -> None:
|
||||
set_waiting_host_and_next_retry_vm(vm, reason="no-eligible-hosts-insufficient-gpu")
|
||||
continue
|
||||
|
||||
from app.models.network import Network, NetworkPort
|
||||
|
||||
ports = NetworkPort.query.filter_by(workload_id=vm.id).all()
|
||||
this_vm_networks = list({p.network_id for p in ports if p.network_id})
|
||||
networks_by_id = {
|
||||
nid: Network.query.get(nid) for nid in this_vm_networks
|
||||
}
|
||||
required_ovs_bridges = [
|
||||
net.ovs_bridge for net in networks_by_id.values() if net and net.ovs_bridge
|
||||
]
|
||||
|
||||
if required_ovs_bridges:
|
||||
filtered_hosts = OVSBridgeFilter(required_ovs_bridges).apply(filtered_hosts)
|
||||
if not filtered_hosts:
|
||||
logger.info(
|
||||
"No hosts carrying OVS bridge(s) %s for VM %s retry, backing off",
|
||||
required_ovs_bridges, vm.id,
|
||||
)
|
||||
set_waiting_host_and_next_retry_vm(vm, reason="no-eligible-hosts-missing-ovs-bridge")
|
||||
continue
|
||||
|
||||
# Exclude hosts already tried this cycle; if that empties an otherwise-eligible
|
||||
# pool, every candidate has failed → terminal (bounded by retry policy).
|
||||
eligible = filtered_hosts
|
||||
@@ -158,27 +180,35 @@ def retry_waiting_host_vms(self) -> None:
|
||||
db.session.commit()
|
||||
logger.info("VM %s retry: selected host %s (%s)", vm.id, selected_host.id, selected_host.hostname)
|
||||
|
||||
# Collect ports and volumes from existing DB records
|
||||
# Collect volumes from existing DB records
|
||||
from app.models.models import VolumeWorkloadMapping, Volume
|
||||
from app.models.network import NetworkPort
|
||||
|
||||
vol_mappings = VolumeWorkloadMapping.query.filter_by(workload_id=vm.id).all()
|
||||
volumes = [Volume.query.get(m.volume_id) for m in vol_mappings if m.volume_id]
|
||||
volumes = [v for v in volumes if v]
|
||||
|
||||
ports = NetworkPort.query.filter_by(workload_id=vm.id).all()
|
||||
this_vm_ports = ports
|
||||
this_vm_networks = list({p.network_id for p in ports if p.network_id})
|
||||
|
||||
from app.models.network import Network
|
||||
host_bridges = [b for b in selected_host.ovs_bridges if not b.deleted]
|
||||
this_vm_port_ovs_bridges = {}
|
||||
missing_bridge_ports = []
|
||||
for port in ports:
|
||||
net = Network.query.get(port.network_id) if port.network_id else None
|
||||
net = networks_by_id.get(port.network_id) if port.network_id else None
|
||||
bridge = (net.ovs_bridge if net else None)
|
||||
if not bridge and selected_host.ovs_bridges:
|
||||
bridge = selected_host.ovs_bridges[0].ovs_bridge_name
|
||||
if bridge:
|
||||
this_vm_port_ovs_bridges[port.id] = bridge
|
||||
if not bridge and host_bridges:
|
||||
bridge = host_bridges[0].ovs_bridge_name
|
||||
if not bridge:
|
||||
missing_bridge_ports.append(port.id)
|
||||
continue
|
||||
this_vm_port_ovs_bridges[port.id] = bridge
|
||||
|
||||
if missing_bridge_ports:
|
||||
logger.error(
|
||||
"No ovs_bridge resolvable for port(s) %s on host %s for VM %s",
|
||||
missing_bridge_ports, selected_host.id, vm.id,
|
||||
)
|
||||
set_waiting_host_and_next_retry_vm(vm, reason="no-ovs-bridge-for-port")
|
||||
continue
|
||||
|
||||
# GPU PCI addresses
|
||||
gpu_pci_addresses = []
|
||||
|
||||
@@ -8,7 +8,7 @@ from app.models.models import Volume, Workload, WorkloadHost, VirtualDataCenter,
|
||||
from app.models.models import VmMetadata # metadata PR
|
||||
import uuid
|
||||
from app.models.network import Network
|
||||
from app.scheduling_filters import ExcludeAllDisabledHosts, LibvirtCapableHosts, HasSufficientResources, MostAvailableCapacity, ExcludeAllOfflineHosts, ExcludeHostsWithoutNorthSouthIP, RandomizeHostsOrder, GPUCapableHosts, CapableHostsFilter
|
||||
from app.scheduling_filters import ExcludeAllDisabledHosts, LibvirtCapableHosts, HasSufficientResources, MostAvailableCapacity, ExcludeAllOfflineHosts, ExcludeHostsWithoutNorthSouthIP, RandomizeHostsOrder, GPUCapableHosts, CapableHostsFilter, OVSBridgeFilter
|
||||
from app.compute.controller.server_group_affinity import apply_affinity_filters
|
||||
|
||||
from werkzeug.exceptions import abort
|
||||
@@ -74,7 +74,8 @@ def dispatch_vm_to_worker(vm: Workload, selected_host: WorkloadHost, vm_config:
|
||||
"""
|
||||
dns_task_id = None
|
||||
try:
|
||||
dns_task_id = send_dns_updates_for_vdc(request_vdc.id, logger)
|
||||
dns_task_ids = send_dns_updates_for_vdc(request_vdc.id, logger)
|
||||
dns_task_id = dns_task_ids.get(selected_host.id)
|
||||
except Exception as dns_exc:
|
||||
logger.error(f"DNS update after VM port creation failed: {dns_exc}")
|
||||
|
||||
@@ -230,6 +231,16 @@ def add_VirtualMachine_workload(request_data) -> str:
|
||||
db.session.commit()
|
||||
logger.info(f"Created VM workload {new_VirtualMachine.id} ({vm_vcpu} vCPU, {vm_memory} MB RAM)")
|
||||
|
||||
this_vm_network_objs = []
|
||||
for network in VirtualMachine.get('networks', []):
|
||||
network_obj = Network.query.filter_by(id=network['id'], deleted=False).first()
|
||||
if not network_obj:
|
||||
logger.error(f"Network {network['id']} not found for VM {new_VirtualMachine.id}")
|
||||
return "error"
|
||||
this_vm_network_objs.append(network_obj)
|
||||
|
||||
required_ovs_bridges = [n.ovs_bridge for n in this_vm_network_objs if n.ovs_bridge]
|
||||
|
||||
# ── Host selection ──
|
||||
vm_candidate_hosts, server_group = apply_affinity_filters(list(filtered_hosts), validated_data)
|
||||
new_VirtualMachine.server_group_id = getattr(server_group, 'id', None)
|
||||
@@ -260,6 +271,16 @@ def add_VirtualMachine_workload(request_data) -> str:
|
||||
overall_outcome = "waiting-host"
|
||||
continue
|
||||
|
||||
if required_ovs_bridges:
|
||||
vm_candidate_hosts = OVSBridgeFilter(required_ovs_bridges).apply(vm_candidate_hosts)
|
||||
if not vm_candidate_hosts:
|
||||
logger.warning(
|
||||
f"No hosts carrying OVS bridge(s) {required_ovs_bridges} for VM {new_VirtualMachine.id}"
|
||||
)
|
||||
set_waiting_host_and_next_retry_vm(new_VirtualMachine, reason="no-eligible-hosts-missing-ovs-bridge")
|
||||
overall_outcome = "waiting-host"
|
||||
continue
|
||||
|
||||
selected_host = vm_candidate_hosts[0]
|
||||
new_VirtualMachine.workload_host_id = selected_host.id
|
||||
db.session.add(new_VirtualMachine)
|
||||
@@ -314,16 +335,10 @@ def add_VirtualMachine_workload(request_data) -> str:
|
||||
# ── Network ports ──
|
||||
this_vm_ports = []
|
||||
this_vm_networks = []
|
||||
this_vm_network_objs = []
|
||||
this_vm_port_ovs_bridges = {}
|
||||
|
||||
for network in VirtualMachine.get('networks', []):
|
||||
network_obj = Network.query.filter_by(id=network['id'], deleted=False).first()
|
||||
if not network_obj:
|
||||
logger.error(f"Network {network['id']} not found for VM {new_VirtualMachine.id}")
|
||||
return "error"
|
||||
for network_obj in this_vm_network_objs:
|
||||
this_vm_networks.append(network_obj.id)
|
||||
this_vm_network_objs.append(network_obj)
|
||||
try:
|
||||
port = network_obj.create_port(db.session, workload_id=new_VirtualMachine.id, use_dhcp_range=True)
|
||||
except Exception as e:
|
||||
@@ -332,8 +347,10 @@ def add_VirtualMachine_workload(request_data) -> str:
|
||||
this_vm_ports.append(port)
|
||||
|
||||
bridge = network_obj.ovs_bridge
|
||||
if not bridge and selected_host.ovs_bridges:
|
||||
bridge = selected_host.ovs_bridges[0].ovs_bridge_name
|
||||
if not bridge:
|
||||
host_bridges = [b for b in selected_host.ovs_bridges if not b.deleted]
|
||||
if host_bridges:
|
||||
bridge = host_bridges[0].ovs_bridge_name
|
||||
if not bridge:
|
||||
logger.error(f"No ovs_bridge for network {network_obj.id} on host {selected_host.id}")
|
||||
return "error"
|
||||
|
||||
@@ -19,19 +19,10 @@ from typing import List, Optional
|
||||
def mark_workload_pending_deleted(workload: Workload) -> None:
|
||||
"""
|
||||
Mark a workload as pending-deleted and handle initial resource cleanup.
|
||||
|
||||
|
||||
Args:
|
||||
workload: The Workload object to mark for deletion
|
||||
"""
|
||||
# Hard safety rule: NSControllers are protected control-plane workloads
|
||||
# and must never be marked pending-deleted by generic deletion flows.
|
||||
if workload.workload_type == "NSController":
|
||||
logger.warning(
|
||||
"SAFETY BLOCK: refusing to mark NSController %s pending-deleted",
|
||||
workload.id,
|
||||
)
|
||||
return
|
||||
|
||||
# Mark the workload as pending-deleted
|
||||
workload.set_status("pending-deleted")
|
||||
db.session.add(workload)
|
||||
@@ -72,10 +63,10 @@ def mark_workload_pending_deleted(workload: Workload) -> None:
|
||||
def cleanup_workload_resources(workload_id: str) -> None:
|
||||
"""
|
||||
Perform comprehensive cleanup of all resources associated with a workload.
|
||||
|
||||
|
||||
This function handles the actual deletion of resources that were previously
|
||||
marked as pending-deleted.
|
||||
|
||||
|
||||
Args:
|
||||
workload_id: The ID of the workload to clean up
|
||||
"""
|
||||
@@ -84,15 +75,6 @@ def cleanup_workload_resources(workload_id: str) -> None:
|
||||
logger.warning(f"Workload {workload_id} not found during cleanup")
|
||||
return
|
||||
|
||||
# Hard safety rule: NSControllers are protected and should not be deleted
|
||||
# or have their network resources removed by generic cleanup.
|
||||
if workload.workload_type == "NSController":
|
||||
logger.warning(
|
||||
"SAFETY BLOCK: refusing cleanup_workload_resources for NSController %s",
|
||||
workload.id,
|
||||
)
|
||||
return
|
||||
|
||||
api_token = app.config.get("CLOUDFLARE_API_TOKEN")
|
||||
account_id = app.config.get("CLOUDFLARE_ACCOUNT_ID")
|
||||
zone_id = app.config.get("CLOUDFLARE_ZONE_ID")
|
||||
@@ -208,50 +190,30 @@ def cleanup_cloudflare_resources_for_workload(workload: Workload, cf_manager: Op
|
||||
|
||||
def mark_pod_pending_deleted(pod: ContainerPod) -> List[str]:
|
||||
"""
|
||||
Mark all containers in a pod as pending-deleted.
|
||||
|
||||
Mark all containers in a pod, including its NSController, as pending-deleted.
|
||||
|
||||
Args:
|
||||
pod: The ContainerPod to mark for deletion
|
||||
|
||||
|
||||
Returns:
|
||||
List of container IDs that were marked for deletion
|
||||
"""
|
||||
container_ids = []
|
||||
|
||||
|
||||
# Get all containers in the pod
|
||||
all_containers_in_pod = Workload.query.filter_by(pod_id=pod.id).all()
|
||||
|
||||
has_protected_nscontroller = False
|
||||
|
||||
for container in all_containers_in_pod:
|
||||
if not container:
|
||||
continue
|
||||
|
||||
if container.workload_type == "NSController":
|
||||
has_protected_nscontroller = True
|
||||
logger.warning(
|
||||
"Skipping NSController %s during pod %s pending-delete transition",
|
||||
container.id,
|
||||
pod.id,
|
||||
)
|
||||
continue
|
||||
|
||||
container_ids.append(str(container.id))
|
||||
mark_workload_pending_deleted(container)
|
||||
|
||||
# Mark the pod itself as pending-deleted only when no protected
|
||||
# NSController is present.
|
||||
if has_protected_nscontroller:
|
||||
logger.warning(
|
||||
"Skipping pod %s pending-delete transition because protected NSController exists",
|
||||
pod.id,
|
||||
)
|
||||
else:
|
||||
pod.status = "pending-deleted"
|
||||
db.session.add(pod)
|
||||
|
||||
pod.status = "pending-deleted"
|
||||
db.session.add(pod)
|
||||
db.session.commit()
|
||||
|
||||
|
||||
return container_ids
|
||||
|
||||
|
||||
|
||||
@@ -4,7 +4,6 @@ Consolidates duplicate code between delete_container_workload and delete_pod fun
|
||||
"""
|
||||
|
||||
from flask import g
|
||||
from app import logger
|
||||
from app.models.models import db, WorkloadResourceUsage, VolumeWorkloadMapping, AuditEntry, Workload
|
||||
from app.models.network import NetworkPort
|
||||
|
||||
@@ -12,21 +11,13 @@ from app.models.network import NetworkPort
|
||||
def collect_and_cleanup_network_ports(workload_id: str) -> list[str]:
|
||||
"""
|
||||
Collect network IDs and soft-delete network ports for a workload.
|
||||
|
||||
|
||||
Args:
|
||||
workload_id: The UUID of the workload
|
||||
|
||||
|
||||
Returns:
|
||||
List of network IDs that were associated with the workload
|
||||
"""
|
||||
workload = Workload.query.filter_by(id=workload_id).first()
|
||||
if workload and workload.workload_type == "NSController":
|
||||
logger.warning(
|
||||
"SAFETY BLOCK: refusing to cleanup network ports for NSController %s",
|
||||
workload_id,
|
||||
)
|
||||
return []
|
||||
|
||||
ports = NetworkPort.query.filter_by(workload_id=workload_id).all()
|
||||
network_ids = []
|
||||
for port in ports:
|
||||
|
||||
+41
-25
@@ -6,7 +6,7 @@ and propagating updates to NSController containers via WebSocket tasks.
|
||||
|
||||
Pattern mirrors sdn_helpers.py for consistency with existing codebase.
|
||||
"""
|
||||
from typing import List, Dict, Optional
|
||||
from typing import List, Dict
|
||||
from datetime import datetime
|
||||
import os
|
||||
import re
|
||||
@@ -322,10 +322,12 @@ def generate_dhcp_config_for_vdc(vdc_id: str, logger_instance=None) -> str:
|
||||
except ValueError:
|
||||
pass
|
||||
|
||||
net_tag = f"net-{net.id}"
|
||||
|
||||
if has_dynamic_range:
|
||||
# Full dynamic pool + lease time (12 h)
|
||||
lines.append(
|
||||
f"dhcp-range={net.dhcp_range_start},{net.dhcp_range_end}{subnet_mask},12h"
|
||||
f"dhcp-range=set:{net_tag},{net.dhcp_range_start},{net.dhcp_range_end}{subnet_mask},12h"
|
||||
)
|
||||
else:
|
||||
# Static-only pool: dnsmasq will respond only to hosts with explicit
|
||||
@@ -335,18 +337,27 @@ def generate_dhcp_config_for_vdc(vdc_id: str, logger_instance=None) -> str:
|
||||
if net.ipv4_cidr:
|
||||
try:
|
||||
network_addr = str(_ip.ip_network(net.ipv4_cidr, strict=False).network_address)
|
||||
lines.append(f"dhcp-range={network_addr},static{subnet_mask}")
|
||||
lines.append(f"dhcp-range=set:{net_tag},{network_addr},static{subnet_mask}")
|
||||
except ValueError:
|
||||
log.warning(f"Could not parse ipv4_cidr '{net.ipv4_cidr}' for network {net.name}, skipping static range")
|
||||
|
||||
# Default gateway
|
||||
# Default gateway (scoped to this network's clients only)
|
||||
if net.ipv4_gateway:
|
||||
lines.append(f"dhcp-option=option:router,{net.ipv4_gateway}")
|
||||
lines.append(f"dhcp-option=tag:{net_tag},option:router,{net.ipv4_gateway}")
|
||||
|
||||
# DNS server — prefer nscontroller_ip (itself), fall back to ipv4_dns_servers
|
||||
dns_server = net.nscontroller_ip or net.ipv4_dns_servers
|
||||
if dns_server:
|
||||
lines.append(f"dhcp-option=option:dns-server,{dns_server}")
|
||||
lines.append(f"dhcp-option=tag:{net_tag},option:dns-server,{dns_server}")
|
||||
|
||||
metadata_nexthop = net.nscontroller_ip
|
||||
if metadata_nexthop:
|
||||
routes = f"169.254.169.254/32,{metadata_nexthop}"
|
||||
if net.ipv4_gateway:
|
||||
routes += f",0.0.0.0/0,{net.ipv4_gateway}"
|
||||
lines.append(
|
||||
f"dhcp-option=tag:{net_tag},option:classless-static-route,{routes}"
|
||||
)
|
||||
|
||||
# ── Per-port static DHCP assignments ────────────────────────────────
|
||||
# Every NetworkPort in this network has a pre-assigned MAC and IP in
|
||||
@@ -590,27 +601,36 @@ def _generate_custom_records(vdc_id: str, zone_name: str, log) -> List[str]:
|
||||
return records
|
||||
|
||||
|
||||
def send_dns_updates_for_vdc(vdc_id: str, logger_instance=None) -> Optional[int]:
|
||||
def send_dns_updates_for_vdc(vdc_id: str, logger_instance=None) -> Dict[str, int]:
|
||||
"""
|
||||
Send DNS zone updates to all NSControllers in a VDC.
|
||||
|
||||
|
||||
Mirrors send_sdn_updates_for_networks pattern from sdn_helpers.py.
|
||||
Creates dns-update tasks for each affected worker/host.
|
||||
Uses WebSocket server for task dispatch.
|
||||
|
||||
|
||||
Args:
|
||||
vdc_id: Virtual Data Center ID
|
||||
logger_instance: Optional logger override (defaults to app logger)
|
||||
|
||||
|
||||
Returns:
|
||||
True if at least one update was sent successfully, False otherwise
|
||||
|
||||
Mapping of host_id -> queued dns-update task id. Callers that must
|
||||
sequence work behind a DNS refresh (e.g. booting a VM) should depend on
|
||||
the task for that workload's own host: a VM's DHCP/DNS is served by the
|
||||
NSController on the host it boots on, so no other host's task is a
|
||||
prerequisite.
|
||||
|
||||
Raises:
|
||||
ValueError: If VDC not found or zone generation fails
|
||||
|
||||
|
||||
WHY: DNS updates must be sent to all NSControllers in the VDC because
|
||||
each NSController runs dnsmasq and needs the complete zone file to
|
||||
resolve pod DNS queries within its network namespace.
|
||||
|
||||
Each host's task is emitted independently. They are deliberately NOT chained
|
||||
depends_on one another: hosts serve DNS only for the VMs on themselves, so
|
||||
serialising them buys nothing and lets one unreachable host stall the
|
||||
updates for every other host in the VDC.
|
||||
"""
|
||||
log = logger_instance or logger
|
||||
log.info(f"Broadcasting DNS updates for VDC: {vdc_id}")
|
||||
@@ -655,8 +675,7 @@ def send_dns_updates_for_vdc(vdc_id: str, logger_instance=None) -> Optional[int]
|
||||
# Send to each host via WebSocket server
|
||||
# WHY: WebSocket server pattern matches existing SDN update mechanism
|
||||
# for consistency and reliability
|
||||
success_count = 0
|
||||
last_task_id: Optional[int] = None
|
||||
dns_task_ids: Dict[str, int] = {}
|
||||
# Use environment variable for WebSocket server URL, with fallback to default
|
||||
websocket_server_url = app.config["WEBSOCKET_SERVER_URL"]
|
||||
headers = {"Content-Type": "application/json"}
|
||||
@@ -666,7 +685,7 @@ def send_dns_updates_for_vdc(vdc_id: str, logger_instance=None) -> Optional[int]
|
||||
if not host:
|
||||
log.warning(f"Host {host_id} not found or deleted")
|
||||
continue
|
||||
|
||||
|
||||
# Get NSController container IDs for this host
|
||||
container_ids = nscontroller_map.get(host_id, [])
|
||||
|
||||
@@ -683,9 +702,6 @@ def send_dns_updates_for_vdc(vdc_id: str, logger_instance=None) -> Optional[int]
|
||||
}
|
||||
}
|
||||
|
||||
if last_task_id is not None:
|
||||
task_payload["depends_on"] = last_task_id
|
||||
|
||||
log.info(
|
||||
f"Sending DNS update to host {host.id} with "
|
||||
f"{len(container_ids)} NSController container(s) via {websocket_server_url}"
|
||||
@@ -698,11 +714,11 @@ def send_dns_updates_for_vdc(vdc_id: str, logger_instance=None) -> Optional[int]
|
||||
headers=headers,
|
||||
timeout=5
|
||||
)
|
||||
|
||||
|
||||
if response.status_code == 201:
|
||||
last_task_id = response.json().get("task_id")
|
||||
success_count += 1
|
||||
log.info(f"DNS update task {last_task_id} queued for host {host.id}")
|
||||
task_id = response.json().get("task_id")
|
||||
dns_task_ids[host.id] = task_id
|
||||
log.info(f"DNS update task {task_id} queued for host {host.id}")
|
||||
else:
|
||||
log.error(
|
||||
f"Failed to queue DNS task for host {host.id}: "
|
||||
@@ -711,5 +727,5 @@ def send_dns_updates_for_vdc(vdc_id: str, logger_instance=None) -> Optional[int]
|
||||
except Exception as e:
|
||||
log.exception(f"Error sending DNS update to host {host.id}: {str(e)}")
|
||||
|
||||
log.info(f"DNS update sent to {success_count}/{len(host_ids)} hosts")
|
||||
return last_task_id
|
||||
log.info(f"DNS update sent to {len(dns_task_ids)}/{len(host_ids)} hosts")
|
||||
return dns_task_ids
|
||||
@@ -41,6 +41,7 @@ def dispatch_nscontroller_to_host(nsc_workload: Workload, network: Network, host
|
||||
"vni": network.vni,
|
||||
"ovs_bridge": bridge,
|
||||
"gateway": network.ipv4_gateway,
|
||||
"encapsulation": network.encapsulation,
|
||||
}]
|
||||
|
||||
nsc_container = {
|
||||
|
||||
+2
-1
@@ -18,4 +18,5 @@ python-dotenv
|
||||
psutil
|
||||
celery
|
||||
cryptography
|
||||
watchdog
|
||||
watchdog
|
||||
flask_cors
|
||||
@@ -76,11 +76,19 @@ def register_socketio_handlers(socketio):
|
||||
|
||||
# Check if this is a reconcile_and_delete task completion
|
||||
if task.task_type == "reconcile_and_delete" and result.get("success"):
|
||||
logger.info(f"[{worker_id}] Reconcile and delete task completed successfully, initiating host online transition")
|
||||
logger.info(f"[{worker_id}] Reconcile and delete task completed successfully, initiating VM reconciliation")
|
||||
# Import the function here to avoid circular imports
|
||||
from websocket_server.worker_manager import handle_reconcile_and_delete_completion
|
||||
handle_reconcile_and_delete_completion(worker_id)
|
||||
|
||||
# VM reconciliation is the last step before the host goes online. Handled even
|
||||
# on failure so a libvirt error cannot leave the host stuck in 'reconciling'.
|
||||
if task.task_type == "virtual-machine-reconcile":
|
||||
if not result.get("success"):
|
||||
logger.warning(f"[{worker_id}] VM reconciliation failed, bringing host online without a VM sweep")
|
||||
from websocket_server.worker_manager import handle_vm_reconcile_completion
|
||||
handle_vm_reconcile_completion(worker_id, result)
|
||||
|
||||
# Find workers that have tasks waiting on this task so we can wake them up.
|
||||
# Needed for cross-worker depends_on (e.g. VM create waiting on a DNS update
|
||||
# that ran on a different host).
|
||||
|
||||
@@ -215,10 +215,118 @@ def send_reconcile_and_delete_task(worker_id, expected_containers):
|
||||
|
||||
def handle_reconcile_and_delete_completion(worker_id):
|
||||
"""
|
||||
Handle the completion of the reconcile_and_delete task by moving host to online status.
|
||||
Handle the completion of the reconcile_and_delete task by reconciling VMs.
|
||||
|
||||
The host stays in 'reconciling' until the VM sweep finishes; it is moved online
|
||||
from handle_vm_reconcile_completion.
|
||||
"""
|
||||
logger.info(f"[{worker_id}] Reconcile and delete completed, moving host to online status")
|
||||
update_host_status(worker_id, "online")
|
||||
logger.info(f"[{worker_id}] Reconcile and delete completed, starting VM reconciliation")
|
||||
send_vm_reconcile_task(worker_id)
|
||||
|
||||
|
||||
def send_vm_reconcile_task(worker_id):
|
||||
"""
|
||||
Ask the worker to report the libvirt domains it currently has defined.
|
||||
"""
|
||||
try:
|
||||
task = Task(
|
||||
worker_id=worker_id,
|
||||
task_type="virtual-machine-reconcile",
|
||||
job_details=json.dumps({}),
|
||||
status="pending"
|
||||
)
|
||||
|
||||
session = Session()
|
||||
session.add(task)
|
||||
session.commit()
|
||||
session.close()
|
||||
|
||||
logger.debug(f"[{worker_id}] Created virtual-machine-reconcile task")
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"[{worker_id}] Error creating virtual-machine-reconcile task: {e}")
|
||||
# Never strand the host in 'reconciling' just because we could not ask for the report
|
||||
update_host_status(worker_id, "online")
|
||||
|
||||
|
||||
def send_vm_delete_task(worker_id, virtual_machine_id):
|
||||
"""
|
||||
Queue a virtual-machine-delete task for a single domain on this worker.
|
||||
"""
|
||||
try:
|
||||
task = Task(
|
||||
worker_id=worker_id,
|
||||
task_type="virtual-machine-delete",
|
||||
job_details=json.dumps({
|
||||
"virtual_machine_id": virtual_machine_id,
|
||||
"desired_state": "deleted",
|
||||
}),
|
||||
status="pending"
|
||||
)
|
||||
|
||||
session = Session()
|
||||
session.add(task)
|
||||
session.commit()
|
||||
session.close()
|
||||
|
||||
logger.debug(f"[{worker_id}] Created virtual-machine-delete task for {virtual_machine_id}")
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"[{worker_id}] Error creating virtual-machine-delete task for {virtual_machine_id}: {e}")
|
||||
|
||||
|
||||
def handle_vm_reconcile_completion(worker_id, result):
|
||||
"""
|
||||
Decide which of the worker's reported domains are orphans and queue deletes.
|
||||
|
||||
A domain is an orphan when the workload table says it should no longer be here —
|
||||
typically because failover moved the VM to another host while this one was
|
||||
unreachable, leaving the original still running and writing to shared storage.
|
||||
|
||||
Domains with no matching workload row are left strictly alone: this may be a BYO
|
||||
host whose own VMs predate enrolment, and destroying those would be unrecoverable.
|
||||
"""
|
||||
from websocket_server.config import get_db_session
|
||||
from app.models.models import Workload
|
||||
|
||||
domains = (result or {}).get("domains", [])
|
||||
orphans = []
|
||||
|
||||
try:
|
||||
with get_db_session() as session:
|
||||
for domain in domains:
|
||||
name = domain.get("name")
|
||||
if not name:
|
||||
continue
|
||||
|
||||
workload = session.query(Workload).filter(
|
||||
Workload.id == name,
|
||||
Workload.workload_type == "VirtualMachine",
|
||||
).first()
|
||||
|
||||
if not workload:
|
||||
logger.debug(f"[{worker_id}] Domain {name} is not a known workload, leaving it alone")
|
||||
continue
|
||||
|
||||
if workload.deleted or str(workload.workload_host_id) != str(worker_id):
|
||||
orphans.append(name)
|
||||
|
||||
for virtual_machine_id in orphans:
|
||||
send_vm_delete_task(worker_id, virtual_machine_id)
|
||||
|
||||
if orphans:
|
||||
logger.info(
|
||||
f"[{worker_id}] VM reconciliation queued deletes for {len(orphans)} "
|
||||
f"orphaned domain(s): {orphans}"
|
||||
)
|
||||
else:
|
||||
logger.info(f"[{worker_id}] VM reconciliation found no orphans ({len(domains)} domains reported)")
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"[{worker_id}] Error during VM reconciliation: {e}")
|
||||
|
||||
finally:
|
||||
update_host_status(worker_id, "online")
|
||||
|
||||
|
||||
def worker_dispatch_flag_check():
|
||||
|
||||
+28
-18
@@ -282,24 +282,22 @@ class WorkerClient:
|
||||
logger.warning(f"[RECONNECT] Starting reconnection process | sio.connected={self.sio.connected}")
|
||||
reconnect_delay = 1 # Start with 1 second delay
|
||||
max_reconnect_delay = 30 # Maximum delay of 30 seconds
|
||||
|
||||
while self.reconnection_in_progress:
|
||||
|
||||
while self.reconnection_in_progress and self.running:
|
||||
try:
|
||||
logger.info(f"[RECONNECT] Attempt #{int(reconnect_delay)} | delay: {reconnect_delay}s | sio.connected={self.sio.connected}")
|
||||
await asyncio.sleep(reconnect_delay)
|
||||
|
||||
# Try to connect
|
||||
if not self.sio.connected:
|
||||
logger.info(f"[RECONNECT] Calling sio.connect({self.server_url})")
|
||||
await self.sio.connect(self.server_url, socketio_path=self.socketio_path)
|
||||
|
||||
if self.sio.connected:
|
||||
logger.info(f"[RECONNECT] ✓ Verified connected | sio.connected={self.sio.connected}")
|
||||
self.reconnection_in_progress = False
|
||||
break
|
||||
else:
|
||||
logger.warning(f"[RECONNECT] Already connected! sio.connected={self.sio.connected}")
|
||||
|
||||
# If we get here, we've successfully reconnected
|
||||
logger.info(f"[RECONNECT] ✓ Success | sio.connected={self.sio.connected}")
|
||||
self.reconnection_in_progress = False
|
||||
# The on_connect handler will send the join request
|
||||
break
|
||||
raise ConnectionError("connect() returned without establishing a connection")
|
||||
except Exception as e:
|
||||
logger.error(f"[RECONNECT] ✗ Failed: {e}")
|
||||
# Increase delay exponentially, up to max_reconnect_delay
|
||||
@@ -349,10 +347,14 @@ class WorkerClient:
|
||||
result = await loop.run_in_executor(None, functools.partial(ContainerTask(logger,self.docker_monitor).handle_pod_update_with_reconciliation, job_details))
|
||||
elif task_type == "virtual-machine-create":
|
||||
from worker_tasks.libvirt import LibvirtVirtualMachineTask
|
||||
result = await loop.run_in_executor(None, functools.partial(LibvirtVirtualMachineTask(job_details, logger).execute))
|
||||
# result = await loop.run_in_executor(None, functools.partial(LibvirtVirtualMachineTask(job_details, logger).execute))
|
||||
result = await loop.run_in_executor(None,lambda:LibvirtVirtualMachineTask(job_details,logger).execute())
|
||||
elif task_type == "virtual-machine-delete":
|
||||
from worker_tasks.libvirt import LibvirtVirtualMachineTask
|
||||
result = await loop.run_in_executor(None, functools.partial(LibvirtVirtualMachineTask(job_details, logger).execute))
|
||||
elif task_type == "virtual-machine-reconcile":
|
||||
from worker_tasks.libvirt import reconcile_virtual_machines
|
||||
result = await loop.run_in_executor(None, functools.partial(reconcile_virtual_machines, logger))
|
||||
elif task_type == "container-reconcile":
|
||||
result = await loop.run_in_executor(None, functools.partial(ContainerTask(logger,self.docker_monitor).reconcile_all_containers))
|
||||
elif task_type == "reconcile_and_delete":
|
||||
@@ -514,17 +516,25 @@ class WorkerClient:
|
||||
try:
|
||||
logger.info(f"[START] Worker {self.worker_id} connecting to {self.server_url}...")
|
||||
await self.sio.connect(self.server_url, socketio_path=self.socketio_path)
|
||||
|
||||
logger.info(f"[START] Initial connection established | sio.connected={self.sio.connected}")
|
||||
|
||||
# Start the event queue processor if we have an event queue
|
||||
|
||||
|
||||
event_processor = None
|
||||
if self.event_queue:
|
||||
event_processor = asyncio.create_task(self.process_event_queue())
|
||||
await self.sio.wait()
|
||||
event_processor.cancel()
|
||||
|
||||
else:
|
||||
await self.sio.wait()
|
||||
try:
|
||||
while self.running:
|
||||
await self.sio.wait()
|
||||
if not self.running:
|
||||
break
|
||||
logger.info("[START] Disconnected — waiting for reconnect() to bring the session back")
|
||||
while self.reconnection_in_progress and self.running:
|
||||
await asyncio.sleep(1)
|
||||
logger.info(f"[START] Reconnection settled | sio.connected={self.sio.connected}")
|
||||
finally:
|
||||
if event_processor:
|
||||
event_processor.cancel()
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"[START] Error: {e}")
|
||||
|
||||
@@ -1955,15 +1955,20 @@ class ContainerTask:
|
||||
raise RuntimeError(f"Failed to bring port up: {result.stderr}")
|
||||
self.logger.debug(f"Brough port {port_name} up")
|
||||
|
||||
# Add IP address
|
||||
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}/24", "dev", port_name],
|
||||
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}/24 to port {port_name}")
|
||||
self.logger.debug(f"Added IP address {ip_address}/{prefix} to port {port_name}")
|
||||
|
||||
# Ensure loopback is up inside the namespace
|
||||
result = subprocess.run(
|
||||
@@ -2064,17 +2069,36 @@ class ContainerTask:
|
||||
3. Deletes any container that is NOT in the expected list
|
||||
|
||||
Args:
|
||||
job_details (dict): Contains expected_container_ids list
|
||||
|
||||
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")
|
||||
|
||||
@@ -531,3 +531,39 @@ class LibvirtVirtualMachineTask:
|
||||
on_error(str(exc))
|
||||
|
||||
self.logger.info(f"[vm-log-stream] Stream ended for VM '{vm_id}'")
|
||||
|
||||
|
||||
def reconcile_virtual_machines(logger):
|
||||
"""
|
||||
Report every libvirt domain defined on this host.
|
||||
|
||||
Deliberately report-only. Unlike containers, domains carry no marker saying we
|
||||
created them — an enrolled BYO host may already run VMs that predate it joining
|
||||
the region — so the worker cannot tell an orphan from a stranger. The server
|
||||
decides what to delete by checking the reported names against the workload
|
||||
table, and issues virtual-machine-delete for the ones it owns.
|
||||
|
||||
Returns:
|
||||
dict: {"success": bool, "domains": [{"name": str, "active": bool}, ...]}
|
||||
"""
|
||||
conn = None
|
||||
try:
|
||||
conn = libvirt.open("qemu:///system")
|
||||
if conn is None:
|
||||
raise SchedulingError(ErrorType.LIBVIRT_CONNECTION_FAILED, "Failed to open connection to libvirt")
|
||||
|
||||
domains = [
|
||||
{"name": dom.name(), "active": bool(dom.isActive())}
|
||||
for dom in conn.listAllDomains(0)
|
||||
]
|
||||
|
||||
logger.info(f"virtual-machine-reconcile: reporting {len(domains)} domains")
|
||||
return {"success": True, "domains": domains}
|
||||
|
||||
except libvirt.libvirtError as e:
|
||||
logger.error(f"Error listing domains during virtual-machine-reconcile: {e}")
|
||||
return {"success": False, "response": str(e), "domains": []}
|
||||
|
||||
finally:
|
||||
if conn:
|
||||
conn.close()
|
||||
|
||||
@@ -17,6 +17,17 @@ class FlowBuilder:
|
||||
self.tap_resolver = tap_resolver
|
||||
self.meter_mgr = meter_mgr
|
||||
|
||||
@staticmethod
|
||||
def _is_flat(vm_networks: list) -> bool:
|
||||
"""Whether a VNI group is a flat/passthrough network.
|
||||
|
||||
Compared case- and whitespace-insensitively: the encapsulation column
|
||||
defaults to "VXLAN" (upper-case), so a flat network is typically stored
|
||||
as "FLAT". An exact-case check against "flat" would silently fall
|
||||
through to the VXLAN pipeline.
|
||||
"""
|
||||
return str(vm_networks[0].get("encapsulation", "")).strip().lower() == "flat"
|
||||
|
||||
def build_flows(self, payload: dict) -> tuple:
|
||||
"""
|
||||
Build intended flows, meter IDs, and response payload.
|
||||
@@ -45,8 +56,18 @@ class FlowBuilder:
|
||||
for port in processed_ports:
|
||||
by_vni.setdefault(port["vni"], []).append(port)
|
||||
|
||||
flat_bridges = set()
|
||||
non_flat_bridges = set()
|
||||
|
||||
for vni, vm_networks in by_vni.items():
|
||||
if self._is_flat(vm_networks):
|
||||
bridge = vm_networks[0]["bridge"]
|
||||
flat_bridges.add(bridge)
|
||||
flows += self._build_flat_dhcp_flows(vm_networks, bridge)
|
||||
flows += self._build_flat_metadata_flows(vm_networks, bridge)
|
||||
continue
|
||||
bridge = vm_networks[0]["bridge"]
|
||||
non_flat_bridges.add(bridge)
|
||||
vni_peers = [p for p in peers if p["vni"] == vni]
|
||||
flows += self._build_inbound_flows(vm_networks, vni_peers, bridge, vxlan_ports[bridge])
|
||||
flows += self._build_outbound_flows(vm_networks, vni, bridge)
|
||||
@@ -57,14 +78,19 @@ class FlowBuilder:
|
||||
flows += self._build_metadata_flows(vm_networks, vni, bridge)
|
||||
|
||||
for vni, vm_networks in by_vni.items():
|
||||
if self._is_flat(vm_networks):
|
||||
continue
|
||||
bridge = vm_networks[0]["bridge"]
|
||||
vni_peers = [p for p in peers if p["vni"] == vni]
|
||||
flows += self._build_unicast_flows(vm_networks, vni_peers, vni, bridge, vxlan_ports[bridge])
|
||||
flows += self._build_broadcast_flows(vm_networks, vni_peers, vni, bridge, vxlan_ports[bridge])
|
||||
|
||||
for bridge in bridges:
|
||||
for bridge in non_flat_bridges:
|
||||
flows += self._build_default_flows(bridge)
|
||||
|
||||
for bridge in flat_bridges:
|
||||
flows += self._build_flat_passthrough_flow(bridge)
|
||||
|
||||
return flows, response_payload, meter_ids, bridges
|
||||
|
||||
def _process_local_ports(self, local_ports, bridges, vxlan_ports) -> list:
|
||||
@@ -135,6 +161,7 @@ class FlowBuilder:
|
||||
"gateway_ip": network.get("gateway_ip"),
|
||||
"ratelimit_in": network.get("ratelimit_in"),
|
||||
"ratelimit_out": network.get("ratelimit_out"),
|
||||
"encapsulation": network.get("encapsulation", "VXLAN"),
|
||||
}
|
||||
|
||||
processed.append(port_data)
|
||||
@@ -148,6 +175,7 @@ class FlowBuilder:
|
||||
"ip": network["ip"],
|
||||
"vni": network["vni"],
|
||||
"ovs_bridge": bridge,
|
||||
"encapsulation": network.get("encapsulation", "VXLAN"),
|
||||
})
|
||||
|
||||
if first_port_data is not None:
|
||||
@@ -493,6 +521,131 @@ class FlowBuilder:
|
||||
{"cookie": COOKIE_MANAGED, "table": 40, "priority": 0, "match": "", "actions": "resubmit(,41)", "bridge": bridge},
|
||||
]
|
||||
|
||||
def _build_flat_dhcp_flows(self, vm_networks, bridge: str) -> list:
|
||||
"""Pin DHCP to the same-host NSController on a flat/passthrough network.
|
||||
|
||||
A flat network has no reg0/VNI pipeline — regular traffic is switched
|
||||
by the NORMAL catch-all straight out the physical uplink. But DHCP
|
||||
can't be left to NORMAL: a broadcast DISCOVER would flood out to the
|
||||
WAN and the VM could be answered by whatever DHCP server lives
|
||||
upstream. So we keep just the two table-0 direct-delivery rules from
|
||||
the VXLAN path (VM->NSController request, NSController->VM reply),
|
||||
sitting above the priority-0 NORMAL flow. The VM's request goes only
|
||||
to the local NSController port (never the uplink), and everything else
|
||||
the VM sends still falls through to NORMAL.
|
||||
|
||||
Unlike _build_dhcp_flows we omit its generic reg0-resubmit fallback —
|
||||
there is no pipeline to resubmit into on a flat bridge; a broadcast
|
||||
DHCP reply from the NSController is simply flooded back by NORMAL.
|
||||
|
||||
Returns [] when there is no same-host NSController on this network, in
|
||||
which case the VM has no local DHCP server and falls back to NORMAL.
|
||||
"""
|
||||
nsctl = next((p for p in vm_networks if p["is_nscontroller"]), None)
|
||||
if not nsctl or not nsctl["flow_port_name"]:
|
||||
return []
|
||||
|
||||
ns_port = nsctl["flow_port_name"]
|
||||
flows = []
|
||||
for vm in vm_networks:
|
||||
if vm["is_nscontroller"]:
|
||||
continue
|
||||
vm_port = vm["flow_port_name"]
|
||||
if not vm_port:
|
||||
continue
|
||||
|
||||
# VM -> NSController (DHCP request): unicast to the NSController
|
||||
# port only, so the broadcast never reaches the WAN uplink.
|
||||
flows.append({
|
||||
"cookie": COOKIE_MANAGED,
|
||||
"table": 0,
|
||||
"priority": PRIORITY_DHCP_REQUEST,
|
||||
"match": f'in_port="{vm_port}",udp,tp_dst=67',
|
||||
"actions": f'output:"{ns_port}"',
|
||||
"bridge": bridge,
|
||||
})
|
||||
# NSController -> VM (DHCP reply).
|
||||
flows.append({
|
||||
"cookie": COOKIE_MANAGED,
|
||||
"table": 0,
|
||||
"priority": PRIORITY_DHCP_REPLY,
|
||||
"match": f'in_port="{ns_port}",udp,tp_src=67,dl_dst={vm["mac"]}',
|
||||
"actions": f'output:"{vm_port}"',
|
||||
"bridge": bridge,
|
||||
})
|
||||
|
||||
return flows
|
||||
|
||||
def _build_flat_metadata_flows(self, vm_networks, bridge: str) -> list:
|
||||
"""Metadata service (169.254.169.254) for a flat network via same-host NSController.
|
||||
|
||||
Mirrors _build_metadata_flows but without the reg0/VNI pipeline. On a
|
||||
flat bridge there is no table-1 stage, so the ARP responder for the
|
||||
metadata IP lives in table 0, scoped to each VM's in_port (so we never
|
||||
answer a metadata ARP arriving from the WAN uplink). The DNAT flows are
|
||||
reused verbatim from the VXLAN path — they are already table-0 and
|
||||
reg0-independent.
|
||||
|
||||
Returns [] when there is no same-host NSController on this network.
|
||||
"""
|
||||
nsctl = next((p for p in vm_networks if p["is_nscontroller"]), None)
|
||||
if not nsctl or not nsctl["flow_port_name"]:
|
||||
return []
|
||||
|
||||
ns_port = nsctl["flow_port_name"]
|
||||
ns_mac = nsctl["mac"]
|
||||
ns_ip = nsctl["ip"].split("/")[0]
|
||||
|
||||
flows = []
|
||||
for vm in vm_networks:
|
||||
if vm["is_nscontroller"]:
|
||||
continue
|
||||
if not vm["flow_port_name"]:
|
||||
continue
|
||||
flows.append(self._make_flat_metadata_arp(vm["flow_port_name"], ns_mac, bridge))
|
||||
flows.extend(self._make_metadata_dnat(vm, ns_mac, ns_ip, ns_port, bridge))
|
||||
|
||||
return flows
|
||||
|
||||
@staticmethod
|
||||
def _make_flat_metadata_arp(vm_port: str, ns_mac: str, bridge: str) -> dict:
|
||||
"""Table-0 ARP responder for 169.254.169.254 on a flat bridge, scoped to one VM port."""
|
||||
ns_hex = f"{int(ns_mac.replace(':', ''), 16):x}"
|
||||
md_ip_int = int(ipaddress.IPv4Address("169.254.169.254"))
|
||||
|
||||
return {
|
||||
"cookie": COOKIE_MANAGED,
|
||||
"table": 0,
|
||||
"priority": PRIORITY_ARP_METADATA,
|
||||
"match": f'in_port="{vm_port}",arp,arp_tpa=169.254.169.254,arp_op=1',
|
||||
"actions": (
|
||||
"move:NXM_OF_ETH_SRC[]->NXM_OF_ETH_DST[],"
|
||||
f"mod_dl_src:{ns_mac},"
|
||||
"load:0x2->NXM_OF_ARP_OP[],"
|
||||
"move:NXM_NX_ARP_SHA[]->NXM_NX_ARP_THA[],"
|
||||
"move:NXM_OF_ARP_SPA[]->NXM_OF_ARP_TPA[],"
|
||||
f"load:0x{ns_hex}->NXM_NX_ARP_SHA[],"
|
||||
f"load:0x{md_ip_int:x}->NXM_OF_ARP_SPA[],"
|
||||
"IN_PORT"
|
||||
),
|
||||
"bridge": bridge,
|
||||
}
|
||||
|
||||
def _build_flat_passthrough_flow(self, bridge: str) -> list:
|
||||
"""Catch-all NORMAL-switching flow for flat/passthrough networks.
|
||||
|
||||
table=0, priority=0 (PRIORITY_DEFAULT) sits below the only other flows
|
||||
a flat bridge carries — the DHCP/metadata control-plane rules pinned to
|
||||
the NSController (all >=310) — so those are matched first and everything
|
||||
else falls through to plain L2 switching out the physical uplink.
|
||||
Without this, a bridge that has any managed flow at all drops unmatched
|
||||
packets by default; OVS only auto-falls-back to NORMAL when a bridge has
|
||||
zero flows.
|
||||
"""
|
||||
return [
|
||||
{"cookie": COOKIE_MANAGED, "table": 0, "priority": PRIORITY_DEFAULT, "match": "", "actions": "NORMAL", "bridge": bridge},
|
||||
]
|
||||
|
||||
# Helper methods for flow construction
|
||||
@staticmethod
|
||||
def _make_arp_responder(vni: int, ip: str, mac: str, bridge: str) -> dict:
|
||||
|
||||
Reference in New Issue
Block a user