Merge pull request 'Fix(Feat): VxLAN and Stray ports' (#56) from JamesBhattarai/theapi:v1.01/network into master

Reviewed-on: coryHawkvelt/theapi#56
This commit is contained in:
2026-06-03 03:50:59 +00:00
14 changed files with 469 additions and 159 deletions
+2 -71
View File
@@ -14,6 +14,7 @@ from app import app, db, logger
from app.models.models import VirtualDataCenter, Workload, WorkloadRequest, WorkloadHost, WorkloadResourceUsage
from app.models.network import Network, NetworkPort
from app.utils.nscontroller_config import build_nscontroller_config, get_nscontroller_resource_requirements
from app.utils.nscontroller_helpers import dispatch_nscontroller_to_host
from app.utils.standard_responses import api_response
from app.utils.sdn_helpers import send_sdn_updates_for_networks
from app.utils.dns_helpers import send_dns_updates_for_vdc
@@ -44,76 +45,6 @@ def _try_silent_operation(operation, operation_name, is_fatal=False):
raise
def _dispatch_nscontroller_to_host(nsc_workload: Workload, network: Network, host: WorkloadHost) -> dict:
"""
Build a minimal standalone pod payload for an NSController and dispatch it
to the given host via the WebSocket/task server.
Returns the response dict from the task server.
"""
vdc_id = nsc_workload.vdc_id
network_id = network.id
# Rebuild config with up-to-date DB state (ports exist now)
nsc_config = build_nscontroller_config(
vdc_id=vdc_id,
container_name=nsc_workload.name,
pod_id=None,
include_ports=False,
network_id=network_id,
)
# Attach the network port (if any) so the worker can wire up OVS
port = NetworkPort.query.filter_by(workload_id=nsc_workload.id, deleted=False).first()
if port:
bridge = network.ovs_bridge or (host.ovs_bridges[0].ovs_bridge_name if host.ovs_bridges else None)
nsc_config["network_ports"] = [{
"network_id": str(network.id),
"name": port.name,
"ip_address": port.ip_address,
"mac_address": port.mac_address,
"subnet_mask": port.subnet_mask,
"dns_servers": port.dns_servers.split(",") if port.dns_servers else [],
"vni": network.vni,
"ovs_bridge": bridge,
"gateway": network.ipv4_gateway,
}]
# Standalone pod-update style payload understood by the worker's container task
nsc_container = {
"container_id": str(nsc_workload.id),
"docker_image": nsc_config.get("docker_image", "xcloudify-nscontroller:latest"),
"container_name": nsc_workload.name,
"workload_type": "NSController",
"desired_state": "running",
"use_dns": True,
"vdc_id": vdc_id,
"dns_config": nsc_config.get("dns_config", {}),
"env": nsc_config.get("env", {}),
"cpu": nsc_config.get("cpu", 4),
"mem_limit": nsc_config.get("mem_limit", 128),
}
if "network_ports" in nsc_config:
nsc_container["network_ports"] = nsc_config["network_ports"]
payload = {
"worker_id": str(host.id),
"task_type": "pod-update",
"job_details": {
"pod_id": f"standalone-nsc-{nsc_workload.id}",
"containers": [nsc_container],
},
}
resp = http_requests.post(
app.config["WEBSOCKET_SERVER_URL"],
data=json.dumps(payload),
headers={"Content-Type": "application/json"},
timeout=10,
)
return resp
@api_bp.route('/networks/<network_id>/nscontroller', methods=['POST'])
def create_nscontroller_for_network(network_id):
"""
@@ -247,7 +178,7 @@ def create_nscontroller_for_network(network_id):
# Dispatch to worker
try:
resp = _dispatch_nscontroller_to_host(nsc_workload, network, host)
resp = dispatch_nscontroller_to_host(nsc_workload, network, host)
if resp.status_code in [200, 201, 202]:
nsc_workload.set_status("allocated")
db.session.commit()
+2 -1
View File
@@ -354,7 +354,8 @@ class WorkloadHost(BaseModel):
sdn_payload = {
'peer-ports': [],
'local-ports': []
'local-ports': [],
'local_ip': self.ip_address_eastwest
}
# Gather all workloads (excluding deleted) on this host + optional new one
+22 -6
View File
@@ -14,6 +14,7 @@ from werkzeug.exceptions import abort
from app.utils.sdn_helpers import send_sdn_updates_for_networks
from app.utils.dns_helpers import send_dns_updates_for_vdc
from app.utils.standard_responses import api_response
from app.utils.nscontroller_helpers import ensure_nscontroller_on_host
from app.compute.routes.gpu import schedule_gpu_for_vm
from runtime_urls import WEBSOCKET_SERVER_URL
@@ -284,6 +285,7 @@ def add_VirtualMachine_workload(request_data):
# Create network ports if networks are specified
this_vm_ports = []
this_vm_networks = []
this_vm_network_objs = []
this_vm_port_ovs_bridges = {} # port.id -> resolved ovs_bridge name
logger.info(f"Processing {len(VirtualMachine.get('networks', []))} networks for VirtualMachine")
for net_idx, network in enumerate(VirtualMachine.get('networks', [])):
@@ -298,6 +300,7 @@ def add_VirtualMachine_workload(request_data):
error_type="NETWORK_NOT_FOUND"
)
this_vm_networks.append(network_obj.id)
this_vm_network_objs.append(network_obj)
if network_obj:
logger.info(f"Found network {network_obj.name} (ID: {network_obj.id})")
try:
@@ -341,15 +344,23 @@ def add_VirtualMachine_workload(request_data):
else:
logger.warning(f"Network with ID {network['id']} not found in database")
# Ensure an NSController is running on the selected host for each network.
# If one doesn't exist yet, auto-create and dispatch it so DHCP/DNS is available
# before the VM boots.
for net_obj in this_vm_network_objs:
ensure_nscontroller_on_host(net_obj, selected_host, request_vdc.id)
# === Calculate all peers across selected_host for all its VMs ===
# logger.info("Calculating peer VMs for network connectivity")
# sdn_update_payload = selected_host.generate_sdn_payload(db.session)
this_host_sdn_update_task_id=send_sdn_updates_for_networks(this_vm_networks,selected_host.id)
# Push updated dnsmasq config to the NSController so the new VM's dhcp-host entry is present before the VM boots and sends a DHCP
# sdn_update_payload = selected_host.generate_sdn_payload(db.session)
# Push updated dnsmasq config to the NSController so the new VM's dhcp-host entry is
# present before the VM boots and sends a DHCP request. We capture the last DNS task
# ID and pass it as depends_on on the VM create task so the worker won't start the
# VM domain until the DNS update has been acknowledged (cross-host safe via ack.py).
dns_task_id = None
try:
send_dns_updates_for_vdc(request_vdc.id, logger)
dns_task_id = send_dns_updates_for_vdc(request_vdc.id, logger)
except Exception as dns_exc:
logger.error(f"DNS update after VM port creation failed: {dns_exc}")
@@ -423,7 +434,7 @@ def add_VirtualMachine_workload(request_data):
payload = {
"worker_id": selected_host.id,
"task_type": "virtual-machine-create",
"depends_on": this_host_sdn_update_task_id,
**({"depends_on": dns_task_id} if dns_task_id else {}),
"job_details": {
"virtual_machine_name": VirtualMachine['name'],
"virtual_machine_id": new_VirtualMachine.id,
@@ -493,6 +504,11 @@ def add_VirtualMachine_workload(request_data):
db.session.add(new_VirtualMachine)
db.session.commit()
logger.debug(f"Updated VirtualMachine status to 'allocated'")
vm_create_task_id = websocket_server_response_data.get("task_id")
send_sdn_updates_for_networks(
this_vm_networks,
depends_on_for_host=(selected_host.id, vm_create_task_id)
)
else:
logger.error(f"Failed to create task for VirtualMachine {new_VirtualMachine.id}")
new_VirtualMachine.set_status("failed-allocation")
+11 -7
View File
@@ -590,7 +590,7 @@ 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) -> bool:
def send_dns_updates_for_vdc(vdc_id: str, logger_instance=None) -> Optional[int]:
"""
Send DNS zone updates to all NSControllers in a VDC.
@@ -656,10 +656,11 @@ def send_dns_updates_for_vdc(vdc_id: str, logger_instance=None) -> bool:
# WHY: WebSocket server pattern matches existing SDN update mechanism
# for consistency and reliability
success_count = 0
last_task_id: Optional[int] = None
# Use environment variable for WebSocket server URL, with fallback to default
websocket_server_url = app.config["WEBSOCKET_SERVER_URL"]
headers = {"Content-Type": "application/json"}
for host_id in host_ids:
host = WorkloadHost.query.filter_by(id=host_id, deleted=False).first()
if not host:
@@ -681,12 +682,15 @@ def send_dns_updates_for_vdc(vdc_id: str, logger_instance=None) -> bool:
"nscontroller_container_ids": container_ids
}
}
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}"
)
try:
response = requests.post(
websocket_server_url,
@@ -696,9 +700,9 @@ def send_dns_updates_for_vdc(vdc_id: str, logger_instance=None) -> bool:
)
if response.status_code == 201:
task_id = response.json().get("task_id")
last_task_id = response.json().get("task_id")
success_count += 1
log.info(f"DNS update task {task_id} queued for host {host.id}")
log.info(f"DNS update task {last_task_id} queued for host {host.id}")
else:
log.error(
f"Failed to queue DNS task for host {host.id}: "
@@ -708,4 +712,4 @@ def send_dns_updates_for_vdc(vdc_id: str, logger_instance=None) -> bool:
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 success_count > 0
return last_task_id
+180
View File
@@ -0,0 +1,180 @@
"""
Shared helpers for creating and dispatching NSController workloads.
Used by both the NSController API routes and the VM scheduling path.
"""
import uuid
import json
import requests as http_requests
from app import app, db, logger
from app.models.models import Workload, WorkloadResourceUsage
from app.models.network import Network, NetworkPort
from app.utils.nscontroller_config import build_nscontroller_config, get_nscontroller_resource_requirements
from app.utils.sdn_helpers import send_sdn_updates_for_networks
def dispatch_nscontroller_to_host(nsc_workload: Workload, network: Network, host) -> http_requests.Response:
"""
Build a standalone pod payload for an NSController and dispatch it to the given host.
Returns the HTTP response from the task server.
"""
vdc_id = nsc_workload.vdc_id
network_id = network.id
nsc_config = build_nscontroller_config(
vdc_id=vdc_id,
container_name=nsc_workload.name,
pod_id=None,
include_ports=False,
network_id=network_id,
)
port = NetworkPort.query.filter_by(workload_id=nsc_workload.id, deleted=False).first()
if port:
bridge = network.ovs_bridge or (host.ovs_bridges[0].ovs_bridge_name if host.ovs_bridges else None)
nsc_config["network_ports"] = [{
"network_id": str(network.id),
"name": port.name,
"ip_address": port.ip_address,
"mac_address": port.mac_address,
"subnet_mask": port.subnet_mask,
"dns_servers": port.dns_servers.split(",") if port.dns_servers else [],
"vni": network.vni,
"ovs_bridge": bridge,
"gateway": network.ipv4_gateway,
}]
nsc_container = {
"container_id": str(nsc_workload.id),
"docker_image": nsc_config.get("docker_image", "xcloudify-nscontroller:latest"),
"container_name": nsc_workload.name,
"workload_type": "NSController",
"desired_state": "running",
"use_dns": True,
"vdc_id": vdc_id,
"dns_config": nsc_config.get("dns_config", {}),
"env": nsc_config.get("env", {}),
"cpu": nsc_config.get("cpu", 4),
"mem_limit": nsc_config.get("mem_limit", 128),
}
if "network_ports" in nsc_config:
nsc_container["network_ports"] = nsc_config["network_ports"]
payload = {
"worker_id": str(host.id),
"task_type": "pod-update",
"job_details": {
"pod_id": f"standalone-nsc-{nsc_workload.id}",
"containers": [nsc_container],
},
}
return http_requests.post(
app.config["WEBSOCKET_SERVER_URL"],
data=json.dumps(payload),
headers={"Content-Type": "application/json"},
timeout=10,
)
def ensure_nscontroller_on_host(network: Network, host, vdc_id: str) -> "Workload | None":
"""
Ensure an NSController exists on `host` for `network`.
If one is already present (active port), returns it immediately.
Otherwise creates, provisions, and dispatches a new one.
Returns the workload on success, or None if creation failed.
"""
existing = Workload.query.filter_by(
vdc_id=vdc_id,
workload_type="NSController",
workload_host_id=host.id,
deleted=False,
).first()
if existing:
existing_port = NetworkPort.query.filter_by(
workload_id=existing.id, network_id=network.id, deleted=False
).first()
if existing_port:
logger.info(
f"NSController {existing.id} already present on host {host.id} for network {network.id}"
)
return existing
if not network.ipv4_cidr:
logger.warning(
f"Cannot auto-create NSController: network {network.id} ({network.name}) has no CIDR configured"
)
return None
nsc_name = f"NSC-{network.name[:12]}-{uuid.uuid4().hex[:6]}"
logger.info(f"Auto-creating NSController '{nsc_name}' on host {host.id} for network {network.id}")
try:
nsc_config = build_nscontroller_config(
vdc_id=vdc_id,
container_name=nsc_name,
pod_id=None,
include_ports=False,
network_id=network.id,
)
except Exception as e:
logger.error(f"Failed to build NSController config for network {network.id}: {e}")
return None
nsc_workload = Workload(
id=str(uuid.uuid4()),
vdc_id=vdc_id,
workload_type="NSController",
name=nsc_name,
workload_host_id=host.id,
launch_params=json.dumps(nsc_config),
deleted=False,
)
db.session.add(nsc_workload)
db.session.flush()
res = get_nscontroller_resource_requirements()
db.session.add(WorkloadResourceUsage(resource_type="cpu", quantity=res["cpu"], workload_id=nsc_workload.id))
db.session.add(WorkloadResourceUsage(resource_type="ram", quantity=res["ram"], workload_id=nsc_workload.id))
try:
port = network.create_port(
db.session,
workload_id=nsc_workload.id,
port_type="NSController",
ip_address=network.nscontroller_ip if network.nscontroller_ip else None,
)
network.nscontroller_ip = port.ip_address
except Exception as e:
db.session.rollback()
logger.error(f"Failed to allocate NSController port for network {network.id}: {e}")
return None
nsc_workload.set_status("pending-allocation")
db.session.commit()
try:
send_sdn_updates_for_networks([network.id], host.id)
except Exception as e:
logger.warning(f"SDN update failed for auto-NSController on network {network.id}: {e}")
try:
resp = dispatch_nscontroller_to_host(nsc_workload, network, host)
if resp.status_code in [200, 201, 202]:
nsc_workload.set_status("allocated")
db.session.commit()
logger.info(
f"Auto-created and dispatched NSController {nsc_workload.id} "
f"on host {host.id} for network {network.id}"
)
else:
nsc_workload.set_status("failed-allocation")
db.session.commit()
logger.error(
f"Worker rejected NSController dispatch for network {network.id}: {resp.text}"
)
except Exception as e:
nsc_workload.set_status("failed-allocation")
db.session.commit()
logger.error(f"Failed to dispatch auto-NSController for network {network.id}: {e}")
return nsc_workload
+12 -1
View File
@@ -14,13 +14,20 @@ from app.models.network import NetworkPort
from flask import json
import requests
def send_sdn_updates_for_networks(network_ids: List[str], return_task_for_host: Optional[str] = None) -> Optional[str]:
def send_sdn_updates_for_networks(
network_ids: List[str],
return_task_for_host: Optional[str] = None,
depends_on_for_host: Optional[tuple] = None,
) -> Optional[str]:
"""
Send SDN update tasks to all hosts that have workloads connected to the specified networks.
Args:
network_ids (List[str]): List of network UUIDs whose connectivity changed.
return_task_for_host (Optional[str]): If specified, the task_id for this host's SDN update will be returned.
depends_on_for_host (Optional[tuple]): (host_id, task_id) — if specified, the sdn-update task
for that host will include depends_on so it runs after the given task completes. Used to
ensure the VM's TAP exists before flows are installed on the host creating the VM.
Returns:
Optional[str]: The task_id for the specified host, if requested and successful. Otherwise None.
@@ -61,6 +68,10 @@ def send_sdn_updates_for_networks(network_ids: List[str], return_task_for_host:
"job_details": sdn_payload
}
if depends_on_for_host and str(host.id) == str(depends_on_for_host[0]):
task_payload["depends_on"] = depends_on_for_host[1]
logger.debug("sdn-update for host %s will depend on task %s", host.id, depends_on_for_host[1])
try:
response = requests.post(websocket_server_url, data=json.dumps(task_payload), headers=headers)
response_data = response.json()
+15
View File
@@ -48,8 +48,23 @@ def register_socketio_handlers(socketio):
from websocket_server.worker_manager import handle_reconcile_and_delete_completion
handle_reconcile_and_delete_completion(worker_id)
# 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).
dependent_worker_ids = [
row[0] for row in session.query(Task.worker_id).filter(
Task.depends_on == task_id,
Task.status == "pending",
Task.worker_id != worker_id,
).distinct().all()
]
redis_client = get_redis_client()
redis_client.set(f"worker_status_{worker_id}", "idle")
for dep_worker_id in dependent_worker_ids:
logger.info(f"[{worker_id}] Waking worker {dep_worker_id} — its task depends on completed task {task_id}")
redis_client.set(f"worker_queue_{dep_worker_id}", "True")
except Exception as e:
logger.exception(f"[{worker_id}] Failed to acknowledge task {task_id}: {e}")
@@ -22,12 +22,11 @@
<interface type='bridge'>
<mac address='{{ port.port_mac_address }}'/>
<source bridge='{{ port.bridge | default("br-int") }}'/>
<target dev='{{ port.port_name | default("vnet" + loop.index0 | string) }}'/>
{% if port.ovs_interface_id %}
<virtualport type='openvswitch'>
<parameters interfaceid='{{ port.ovs_interface_id }}'/>
</virtualport>
{% else %}
<target dev='{{ port.port_name | default("vnet" + loop.index0 | string) }}'/>
{% endif %}
<model type='virtio'/>
<alias name='{{ port.port_name | default("net" + loop.index0 | string) }}'/>
-14
View File
@@ -14,7 +14,6 @@ from xml.dom import minidom
from worker_tasks.volumes import VolumeProcessor
from settings import settings # Import global settings
from worker_tasks.libvirt_xml_builder import _generate_xml_from_config
from worker_tasks.ovs_sdn import OVS_SDN
# Import VFIO passthrough controller (app/ is two levels up from worker/)
_REPO_ROOT = os.path.abspath(os.path.join(os.path.dirname(__file__), "..", ".."))
@@ -111,19 +110,6 @@ class LibvirtVirtualMachineTask:
self.logger.info("No volumes found in virtual_machine_config or volumes list is empty.")
network_ports = self.virtual_machine_config.get("network_ports", [])
if network_ports:
ovs = OVS_SDN({}, self.logger)
for np in network_ports:
port_name = np.get("port_name")
mac = np.get("port_mac_address")
bridge = np.get("ovs_bridge")
if port_name and mac and bridge:
self.logger.info(f"Ensuring OVS port '{port_name}' on bridge '{bridge}' before XML generation")
ovs.port_mgr.ensure_vm_port(port_name, mac, bridge)
else:
self.logger.warning(f"Network port entry missing required fields, skipping: {np}")
# Generate XML configuration
self.xml_config = _generate_xml_from_config(self)
+4 -24
View File
@@ -22,27 +22,11 @@ templates always have the data they need.
"""
import copy
import subprocess
from libvirt_template.parser.parser import VMPayload
from libvirt_template.templating import render_vm
def _get_ovs_interface_id(interface_name: str, logger) -> str | None:
"""Return the OVS Interface UUID for *interface_name*, or None on failure."""
try:
result = subprocess.run(
["ovs-vsctl", "get", "Interface", interface_name, "_uuid"],
capture_output=True, text=True, check=True,
)
return result.stdout.strip().strip('"')
except subprocess.CalledProcessError as e:
logger.error(
f"Failed to get OVS interface ID for '{interface_name}': "
f"{e.stderr.strip()}"
)
return None
def _pci_addr_to_feature_gpu(addr: str) -> dict:
"""Convert a short/long-form BDF string to a VMFeatures.gpu entry dict.
@@ -95,15 +79,11 @@ def _generate_xml_from_config(self) -> str:
vol["path"] = resolved
# ── 2. Inject OVS interface IDs into network ports ─────────────────
# The template emits <virtualport type='openvswitch'> when
# port.ovs_interface_id is set. Also normalise ovs_bridge → bridge
# so the template key matches.
# Use the DB port_id as the OVS iface-id so libvirt emits
# <virtualport type='openvswitch'> without needing a pre-existing
# OVS port. Also normalise ovs_bridge → bridge for the template.
for port in cfg.get("network_ports", []):
port_name = port.get("port_name")
if port_name:
iface_id = _get_ovs_interface_id(port_name, self.logger)
port["ovs_interface_id"] = iface_id or ""
# cloudify uses "ovs_bridge"; template uses "bridge"
port["ovs_interface_id"] = port.get("port_id", "")
if "ovs_bridge" in port and "bridge" not in port:
port["bridge"] = port["ovs_bridge"]
+88 -18
View File
@@ -58,7 +58,9 @@ class FlowBuilder:
for vni, vm_networks in by_vni.items():
bridge = vm_networks[0]["bridge"]
flows += self._build_broadcast_flows(vm_networks, vni, bridge, vxlan_ports[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:
flows += self._build_default_flows(bridge)
@@ -86,7 +88,6 @@ class FlowBuilder:
if bridge not in vxlan_ports:
vxlan_ports[bridge] = self.port_mgr.get_existing_vxlan_ports()
tap = self.tap_resolver.get_tap_for_mac(network["mac"], bridge)
flow_port = None
if is_nsc:
@@ -100,15 +101,26 @@ class FlowBuilder:
)
flow_port = port_name
else:
# Regular VMs should use the libvirt TAP that already exists on the host.
# Do not create an internal OVS port for them.
flow_port = tap
flow_port = None
try:
if self.port_mgr.client.port_exists(port_name, bridge):
iface_type = self.port_mgr.client.get_interface_attr(
port_name, "type"
).strip().strip('"')
if iface_type != "internal":
flow_port = port_name
except Exception as e:
self.logger.debug("Could not check port %s type: %s", port_name, e)
if not tap and not is_nsc:
self.logger.warning(
"No TAP found for VM MAC %s on %s; flow_port will be None until TAP appears",
network["mac"], bridge
)
if not flow_port:
flow_port = self.tap_resolver.get_tap_for_mac(
network["mac"], bridge, attempts=10, delay=0.5
)
if not flow_port:
self.logger.warning(
"No TAP found for VM MAC %s on %s; flow_port will be None",
network["mac"], bridge
)
port_data = {
"workload_id": workload_id,
@@ -151,7 +163,7 @@ class FlowBuilder:
"""Build table 0-10 inbound classification flows."""
flows = []
for peer in vni_peers:
vxlan_port = self.port_mgr.ensure_vxlan_port(peer["workload_host_ip"], bridge, vxlan_map)
vxlan_port = self.port_mgr.ensure_vxlan_port(remote_ip=peer["workload_host_ip"], bridge=bridge, existing_ports=vxlan_map)
flows.append({
"cookie": COOKIE_MANAGED,
"table": 0,
@@ -242,7 +254,7 @@ class FlowBuilder:
created_meters.add(meter_id)
action_flow["actions"] = f"meter:{meter_id},resubmit(,40)"
else:
action_flow["actions"] = "resubmit(,40)" if action_flow["table"] == 30 else "resubmit(,40)"
action_flow["actions"] = "resubmit(,31)" if action_flow["table"] == 30 else "resubmit(,40)"
flows.append(action_flow)
@@ -352,6 +364,8 @@ class FlowBuilder:
for vm in vm_networks:
if vm["is_nscontroller"]:
continue
if not vm["flow_port_name"]:
continue
flows.append({
"cookie": COOKIE_MANAGED,
"table": 0,
@@ -384,7 +398,58 @@ class FlowBuilder:
return flows
def _build_broadcast_flows(self, vm_networks, vni, bridge, vxlan_map) -> list:
def _build_unicast_flows(self, vm_networks, vni_peers, vni, bridge, vxlan_map) -> list:
"""Build per-MAC unicast delivery flows.
A known destination MAC lives in exactly one place:
- local VM/NSC on this host → output to its TAP/port (table 40,
above the broadcast flood at PRIORITY_BROADCAST so unicast wins).
- remote VM on a peer host → fall through to table 41, write the
VNI into the tunnel id and send out the single VXLAN port to the
host that owns that MAC.
Without these flows, generic VM↔VM traffic (same host or cross-host)
has no delivery rule and is dropped at the end of the pipeline.
"""
flows = []
# Local destinations → their own port.
for port in vm_networks:
if not port.get("flow_port_name"):
continue
flows.append({
"cookie": COOKIE_MANAGED,
"table": 40,
"priority": PRIORITY_LOCAL_DELIVERY,
"match": f'reg0=0x{vni:x},dl_dst={port["mac"]}',
"actions": f'output:"{port["flow_port_name"]}"',
"bridge": bridge,
})
# Remote destinations → tunnel to the host that owns the MAC.
for peer in vni_peers:
vxlan_port = vxlan_map.get(peer["workload_host_ip"])
if not vxlan_port:
self.logger.warning(
"No VXLAN port for peer host %s (vni 0x%x); skipping unicast flow for %s",
peer["workload_host_ip"], vni, peer["mac"]
)
continue
flows.append({
"cookie": COOKIE_MANAGED,
"table": 41,
"priority": PRIORITY_LOCAL_DELIVERY,
"match": f'reg0=0x{vni:x},dl_dst={peer["mac"]}',
"actions": (
"move:NXM_NX_REG0[]->NXM_NX_TUN_ID[0..31],"
f'output:"{vxlan_port}"'
),
"bridge": bridge,
})
return flows
def _build_broadcast_flows(self, vm_networks, vni_peers, vni, bridge, vxlan_map) -> list:
"""Build broadcast delivery flows."""
flows = []
local_outputs = [f'output:"{p["flow_port_name"]}"' for p in vm_networks if p.get("flow_port_name")]
@@ -395,18 +460,23 @@ class FlowBuilder:
"table": 40,
"priority": PRIORITY_BROADCAST,
"match": f"reg0=0x{vni:x},dl_dst=ff:ff:ff:ff:ff:ff",
"actions": ",".join(local_outputs),
"actions": ",".join(local_outputs) + ",resubmit(,41)",
"bridge": bridge
})
# VXLAN broadcast (table 41)
vxlan_outputs = [f'output:"{v}"' for v in vxlan_map.values()] if vxlan_map else []
# VXLAN broadcast (table 41) — only to hosts that have peers on this VNI
vni_peer_ips = {p["workload_host_ip"] for p in vni_peers}
vxlan_outputs = [
f'output:"{vxlan_map[ip]}"'
for ip in vni_peer_ips
if ip in vxlan_map
]
if vxlan_outputs:
flows.append({
"cookie": COOKIE_MANAGED,
"table": 41,
"priority": PRIORITY_DEFAULT,
"match": "dl_dst=ff:ff:ff:ff:ff:ff",
"match": f"reg0=0x{vni:x},dl_dst=ff:ff:ff:ff:ff:ff",
"actions": "move:NXM_NX_REG0[]->NXM_NX_TUN_ID[0..31]," + ",".join(sorted(vxlan_outputs)),
"bridge": bridge
})
@@ -549,7 +619,7 @@ class FlowBuilder:
return []
flows = []
ns_port = nsctl["port_name"]
ns_port = nsctl["flow_port_name"]
ns_mac = nsctl["mac"]
ns_ip = nsctl["ip"].split("/")[0]
+92 -11
View File
@@ -1,18 +1,21 @@
"""OVS port management (VM ports, VXLAN ports)."""
import os
import re
import time
import uuid
from logging import Logger
from .base_client import OVSClient
from .constants import VXLAN_PORT_PREFIX
class PortManager:
"""Manage VM and VXLAN ports on OVS bridges."""
def __init__(self, logger: Logger, client: OVSClient):
def __init__(self, logger: Logger, client: OVSClient, local_ip: str = None):
self.logger = logger
self.client = client
self.local_ip = local_ip
def ensure_vm_port(self, port_name: str, mac_address: str, bridge: str,
allow_netns_existing: bool = False) -> str:
@@ -75,17 +78,24 @@ class PortManager:
if not bridge:
raise ValueError(f"Missing 'ovs_bridge' for deleting port {port_name}")
# Protect NSController infrastructure ports
if port_name and port_name.startswith("port-"):
raise ValueError(
f"SAFETY BLOCK: Cannot delete NSController port '{port_name}'. "
f"These are control-plane infrastructure."
)
if not self.client.port_exists(port_name, bridge):
self.logger.warning(f"Port {port_name} doesn't exist on {bridge}")
return False
# Protect OVS internal ports — these are NSController control-plane interfaces.
# Regular VM TAPs are type "" or "system", never "internal".
try:
iface_type = self.client.get_interface_attr(port_name, "type").strip().strip('"')
if iface_type == "internal":
raise ValueError(
f"SAFETY BLOCK: Cannot delete OVS internal port '{port_name}'. "
f"This is a control-plane infrastructure port."
)
except ValueError:
raise
except Exception as e:
self.logger.debug("Could not check type for port %s: %s", port_name, e)
self.logger.info(f"Deleting port {port_name} from {bridge}")
self.client.del_port(bridge, port_name)
return True
@@ -105,20 +115,24 @@ class PortManager:
if not bridge:
raise ValueError(f"Missing 'ovs_bridge' for VXLAN to {remote_ip}")
if existing_ports and remote_ip in existing_ports:
if existing_ports is not None and remote_ip in existing_ports:
return existing_ports[remote_ip]
port_name = f"vxlan-{uuid.uuid4().hex[:6]}"
self.logger.info(f"Creating VXLAN {port_name} on {bridge} → {remote_ip}")
self.client.vsctl([
"add-port", bridge, port_name,
"--", "set", "interface", port_name, "type=vxlan",
f"options:remote_ip={remote_ip}",
*([f"options:local_ip={self.local_ip}"] if self.local_ip else []),
"options:key=flow",
"options:csum=true"
])
if existing_ports is not None:
existing_ports[remote_ip] = port_name
return port_name
def get_existing_vxlan_ports(self) -> dict:
@@ -141,3 +155,70 @@ class PortManager:
current_port = None
return port_map
def _list_vxlan_ports(self, bridge: str) -> list:
"""List managed VXLAN ports on a bridge.
Returns a list of {"name", "remote_ip", "healthy"} dicts. A port is
"healthy" when it has a valid ofport (>0); broken tunnels (e.g. a
duplicate that failed with "File exists") report ofport -1/empty.
"""
ports = []
for name in self.client.list_ports(bridge):
if not name.startswith(f"{VXLAN_PORT_PREFIX}-"):
continue
remote_ip = None
try:
opts = self.client.get_interface_attr(name, "options")
match = re.search(r'remote_ip="?([^",}\s]+)"?', opts)
if match:
remote_ip = match.group(1)
except Exception as e:
self.logger.debug(f"Could not read options for {name}: {e}")
healthy = False
try:
healthy = int(self.client.get_interface_attr(name, "ofport")) > 0
except (ValueError, TypeError):
healthy = False
ports.append({"name": name, "remote_ip": remote_ip, "healthy": healthy})
return ports
def reconcile_vxlan_ports(self, bridge: str, desired_remote_ips: set) -> None:
"""Garbage-collect VXLAN tunnel ports on a bridge.
Keeps at most one healthy tunnel per remote_ip that is still needed,
and removes:
- stale tunnels — remote_ip no longer present in the peer set,
- duplicate tunnels — more than one port to the same remote_ip,
- broken tunnels — failed to attach to ofproto (no valid ofport).
When a desired remote_ip has only broken ports, all of them are removed
so the flow builder recreates a fresh tunnel on this run.
"""
if not bridge:
raise ValueError("Bridge required for VXLAN reconciliation")
by_ip = {}
for port in self._list_vxlan_ports(bridge):
by_ip.setdefault(port["remote_ip"], []).append(port)
for remote_ip, plist in by_ip.items():
# Stale: no peer needs this host anymore -> remove every tunnel to it.
if remote_ip not in desired_remote_ips:
for port in plist:
self.logger.info(f"🧹 GC stale VXLAN {port['name']} → {remote_ip}")
self.client.del_port(bridge, port["name"])
continue
# Desired: keep a single healthy tunnel, drop duplicates/broken.
healthy = [p for p in plist if p["healthy"]]
keep = healthy[0]["name"] if healthy else None
for port in plist:
if port["name"] == keep:
continue
reason = "duplicate" if keep else "broken"
self.logger.info(f"🧹 GC {reason} VXLAN {port['name']} → {remote_ip}")
self.client.del_port(bridge, port["name"])
+22 -1
View File
@@ -22,7 +22,7 @@ class OVS_SDN:
# Initialize modules
self.client = OVSClient(logger)
self.port_mgr = PortManager(logger, self.client)
self.port_mgr = PortManager(logger, self.client, local_ip=params.get("local_ip"))
self.tap_resolver = TAPResolver(logger)
self.flow_parser = FlowParser()
self.flow_enforcer = FlowEnforcer(logger, self.client)
@@ -34,6 +34,8 @@ class OVS_SDN:
def Execute(self) -> dict:
"""Execute the SDN configuration."""
self.logger.debug(f"Input params: {self.params}")
self._reconcile_vxlan_ports()
# Build intended flows
flows, response_payload, meter_ids, bridges = self.flow_builder.build_flows(self.params)
@@ -55,3 +57,22 @@ class OVS_SDN:
"success": True,
"response": response_payload
}
def _reconcile_vxlan_ports(self) -> None:
"""Remove stale/duplicate/broken VXLAN ports prior to flow construction."""
desired_remote_ips = {
peer["workload_host_ip"]
for peer in self.params.get("peer-ports", [])
if peer.get("workload_host_ip")
}
bridges = {
net.get("ovs_bridge")
for local in self.params.get("local-ports", [])
for net in local.get("networks", [])
if net.get("ovs_bridge")
}
for bridge in bridges:
try:
self.port_mgr.reconcile_vxlan_ports(bridge, desired_remote_ips)
except Exception as e:
self.logger.warning(f"VXLAN reconciliation failed for {bridge}: {e}")
+18 -3
View File
@@ -1,6 +1,7 @@
"""TAP/MAC resolution utilities."""
import subprocess
import time
from logging import Logger
from .constants import TAP_MAC_PREFIX, VM_MAC_PREFIX
@@ -11,13 +12,28 @@ class TAPResolver:
def __init__(self, logger: Logger):
self.logger = logger
def get_tap_for_mac(self, mac: str, bridge: str) -> str:
def get_tap_for_mac(self, mac: str, bridge: str, attempts: int = 1, delay: float = 0.5) -> str:
"""
Resolve VM MAC to actual kernel TAP interface on bridge.
Retries up to `attempts` times with `delay` seconds between tries to
handle races where the TAP hasn't appeared on OVS yet.
Returns:
TAP port name or None if not found
"""
for attempt in range(attempts):
result = self._try_get_tap_for_mac(mac, bridge)
if result:
return result
if attempt < attempts - 1:
self.logger.debug(f"TAP for MAC {mac} not found, retrying ({attempt + 1}/{attempts})")
time.sleep(delay)
self.logger.warning(f"No live TAP found for MAC {mac} on {bridge}")
return None
def _try_get_tap_for_mac(self, mac: str, bridge: str) -> str:
tap_mac = TAP_MAC_PREFIX + mac[2:]
# Strategy A: Sysfs scan over bridge ports
@@ -57,7 +73,6 @@ class TAPResolver:
except Exception as e:
self.logger.debug(f"Strategy A (sysfs) failed: {e}")
self.logger.warning(f"No live TAP found for MAC {mac} on {bridge}")
return None
def _parse_ovs_record(self, output: str) -> str: