Files
3cloud-backend/app/utils/create_workload_virtual_machine.py
T
JamesBhattarai 38ddcb5d4d Chore: appropos removed constants
Removed redundant code level constatnt with repo based runtime_urls.py
Added run-master, run-worker.sh scripts for convinient clean, builds
2026-05-23 19:20:18 +05:45

521 lines
27 KiB
Python

#app/utils/create_workload_virtual_machine.py
from flask import json, request, jsonify
import requests
from app import db, logger, app
from app.models.models import Volume, Workload, WorkloadHost, VirtualDataCenter, Image, VolumeWorkloadMapping, WorkloadResourceUsage, WorkloadHostFixedResource
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
from app.compute.controller.server_group_affinity import apply_affinity_filters
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.compute.routes.gpu import schedule_gpu_for_vm
from runtime_urls import WEBSOCKET_TASK_URL
websocket_server_url = app.config.get("WEBSOCKET_SERVER_URL", WEBSOCKET_TASK_URL)
def add_VirtualMachine_workload(request_data):
# DONT Validate incoming data
validated_data = request_data.payload
logger.info("Received request to add VirtualMachine workload")
logger.debug(f"Validated payload: {json.dumps(validated_data, indent=2)}")
# TODO - Check if this user has permissions to VIEW and LAUNCH_VirtualMachine in this VDC
logger.info(f"Looking up Virtual Data Center with ID: {validated_data['virtual_data_center']}")
request_vdc = VirtualDataCenter.query.filter_by(id=(validated_data['virtual_data_center'])).first()
if request_vdc is None:
logger.error(f"Virtual Data Center not found with ID: {validated_data['virtual_data_center']}")
return api_response(
success=False,
status=404,
message="Virtual Data Center not found",
error_type="RESOURCE_NOT_FOUND"
)
logger.info(f"Found VDC: {request_vdc.name} (ID: {request_vdc.id}) in region: {request_vdc.region.name}")
logger.info(f"Searching for eligible workload hosts in region {request_vdc.region.name}")
all_workload_hosts = WorkloadHost.query.filter_by(deleted=0, region_id=request_vdc.region.id).all()
logger.info(f"Found {len(all_workload_hosts)} hosts eligible for placement in region {request_vdc.region.name}")
logger.debug(f"Host list: {[host.id for host in all_workload_hosts]}")
filters = [ExcludeAllOfflineHosts(), ExcludeAllDisabledHosts(), LibvirtCapableHosts(), MostAvailableCapacity(),RandomizeHostsOrder()]
logger.info(f"Applying {len(filters)} filters to host list")
filtered_hosts = all_workload_hosts
for filter in filters:
filter_name = filter.__class__.__name__
logger.debug(f"Applying filter: {filter_name}")
pre_filter_count = len(filtered_hosts)
filtered_hosts = filter.apply(filtered_hosts)
post_filter_count = len(filtered_hosts)
logger.info(f"Filter {filter_name} reduced hosts from {pre_filter_count} to {post_filter_count}")
if len(filtered_hosts) == 0:
logger.error(f"No hosts found in region {request_vdc.region.name} after applying all filters")
return api_response(
success=False,
status=400,
message=f"No hosts found in region {request_vdc.region.name} to place this request",
error_type="NO_ELIGIBLE_HOSTS"
)
logger.info(f"Final filtered host count: {len(filtered_hosts)}")
logger.debug(f"Filtered host IDs: {[host.id for host in filtered_hosts]}")
client_response = []
for idx, VirtualMachine in enumerate(validated_data['virtual-machines']):
logger.info(f"Processing VirtualMachine {idx + 1} of {len(validated_data['virtual-machines'])}: {VirtualMachine['name']}")
vm_candidate_hosts, server_group = apply_affinity_filters(list(filtered_hosts), validated_data)
if vm_candidate_hosts is None:
return api_response(
success=False, status=404,
message=f"Server group '{server_group}' not found",
error_type="RESOURCE_NOT_FOUND",
)
if len(vm_candidate_hosts) == 0:
return api_response(
success=False, status=400,
message=f"No eligible hosts for VM '{VirtualMachine['name']}': server group constraint failed",
error_type="NO_ELIGIBLE_HOSTS",
)
# Filter by actual VM resource requirements (RAM + vCPU)
vm_memory = VirtualMachine.get('memory', 0)
vm_vcpu = VirtualMachine.get('vcpu', 0)
vm_gpu = VirtualMachine.get('gpu', 0)
# gpu_requests is the new list-based format; vm_gpu_count is the integer count for filters
gpu_requests = vm_gpu if isinstance(vm_gpu, list) else None
vm_gpu_count = sum(r.get('count', 1) for r in gpu_requests) if gpu_requests else (vm_gpu if isinstance(vm_gpu, int) else 0)
resource_filter = HasSufficientResources(
required_resources={"ram": vm_memory, "cpu": vm_vcpu},
overcommit_ratios={"cpu": 10.0, "ram": 1.0},
)
vm_candidate_hosts = resource_filter.apply(vm_candidate_hosts)
if len(vm_candidate_hosts) == 0:
logger.info(
"Insufficient Resource available for VM '%s' (ram=%s MB, cpu=%s): %d hosts remain",
VirtualMachine['name'], vm_memory, vm_vcpu, len(vm_candidate_hosts),
)
return api_response(
success=False, status=400,
message=(
f"No eligible hosts for VM '{VirtualMachine['name']}': "
f"insufficient capacity (requires {vm_vcpu} vCPU, {vm_memory} MB RAM)"
),
error_type="INSUFFICIENT_CAPACITY",
)
# Filter by GPU requirements if requested
if vm_gpu_count > 0:
logger.info(f"VM '{VirtualMachine['name']}' requires {vm_gpu_count} GPU(s), applying GPU filter")
gpu_filter = GPUCapableHosts(
required_gpu_count=vm_gpu_count,
gpu_model_requests=gpu_requests, # None for legacy int format
)
vm_candidate_hosts = gpu_filter.apply(vm_candidate_hosts)
if len(vm_candidate_hosts) == 0:
logger.error(
f"No eligible hosts for VM '{VirtualMachine['name']}': "
f"insufficient GPU capacity (requires {vm_gpu_count} GPU(s))"
)
return api_response(
success=False, status=400,
message=(
f"No eligible hosts for VM '{VirtualMachine['name']}': "
f"insufficient GPU capacity (requires {vm_gpu_count} GPU(s))"
),
error_type="INSUFFICIENT_GPU_CAPACITY",
)
selected_host = vm_candidate_hosts[0]
logger.info(f"Selected host ID {selected_host.id} ({selected_host.hostname}) for VirtualMachine placement")
# Create the Workload instance in the database
logger.info(f"Creating new Workload record for VirtualMachine {VirtualMachine['name']}")
new_VirtualMachine = Workload(
name=VirtualMachine['name'],
workload_type="VirtualMachine", # Set workload_type to VirtualMachine
vdc_id=request_vdc.id,
# status="pending-allocation",
launch_params=json.dumps(VirtualMachine),
workload_host_id=selected_host.id,
server_group_id=getattr(server_group, 'id', None)
)
db.session.add(new_VirtualMachine)
db.session.commit()
new_VirtualMachine.set_status("pending-allocation")
logger.info(f"Created VirtualMachine workload with ID: {new_VirtualMachine.id}")
logger.debug(f"VirtualMachine DB record: {new_VirtualMachine.to_json()}")
# Record resource usage for scheduler accounting
db.session.add(WorkloadResourceUsage(
resource_type="cpu",
quantity=vm_vcpu,
workload_id=new_VirtualMachine.id,
))
db.session.add(WorkloadResourceUsage(
resource_type="ram",
quantity=vm_memory,
workload_id=new_VirtualMachine.id,
))
if vm_gpu_count > 0:
db.session.add(WorkloadResourceUsage(
resource_type="gpu",
quantity=vm_gpu_count,
workload_id=new_VirtualMachine.id,
))
db.session.commit()
logger.info(
"Recorded resource usage for VM %s: %s vCPU, %s MB RAM%s",
new_VirtualMachine.id, vm_vcpu, vm_memory,
f", {vm_gpu_count} GPU(s)" if vm_gpu_count > 0 else "",
)
# Create volumes in the database
volumes = []
logger.info(f"Processing {len(VirtualMachine.get('volumes', []))} volumes for VirtualMachine")
for vol_idx, volume in enumerate(VirtualMachine.get('volumes', [])):
logger.debug(f"Processing volume {vol_idx + 1}: {json.dumps(volume, indent=2)}")
_volume_source = ""
_volume_image_id = None
if "image_id" in volume:
logger.info(f"[VM={VirtualMachine['name']}] Volume '{volume['name']}' has image_id: {volume['image_id']}")
# Validate the image ID
try:
logger.debug(f"[VM={VirtualMachine['name']}] Validating image ID: {volume['image_id']}")
# Convert the image ID to a UUID object
image_id = (volume['image_id'])
# Query the image from the database
_image = Image.query.filter_by(id=image_id).first()
if _image is None:
logger.error(f"[VM={VirtualMachine['name']}] Image with ID {volume['image_id']} does not exist in database")
# return {"error": f"Image with ID {volume['image_id']} does not exist"}, 404
return api_response(
success=False,
status=400,
message=f"Image with ID {volume['image_id']} does not exist",
error_type="IMAGE_NOT_EXIST"
)
logger.info(f"[VM={VirtualMachine['name']}] Resolved image: '{_image.name}' (id={_image.id}, location={_image.location}, format={_image.format})")
_volume_image_id = _image.id
_volume_source = _image.location
except ValueError as e:
logger.error(f"[VM={VirtualMachine['name']}] Invalid image ID format '{volume['image_id']}': {e}")
return api_response(
success=False,
status=400,
message=f"Invalid image ID format: {volume['image_id']}",
error_type="IMAGE_INVALID_ID"
)
else:
logger.info(f"[VM={VirtualMachine['name']}] No image specified for volume '{volume['name']}'")
_volume_type = volume.get('type') or 'local'
if _volume_type != volume.get('type'):
logger.debug(f"[VM={VirtualMachine['name']}] Volume '{volume['name']}' has no type specified, defaulting to 'local'")
# NOTE: ceph driver exists in worker code but is not connected via the API layer yet.
_supported_volume_types = ['local', 'lvm']
if _volume_type not in _supported_volume_types:
logger.error(
f"[VM={VirtualMachine['name']}] Volume '{volume['name']}' requested unsupported type '{_volume_type}'. "
f"Supported types: {_supported_volume_types}. Rejecting VM creation."
)
return api_response(
success=False,
status=400,
message=f"Volume type '{_volume_type}' is not yet supported. Supported types: {', '.join(_supported_volume_types)}.",
error_type="UNSUPPORTED_VOLUME_TYPE"
)
# 'local' volumes live under LOCAL_VOLUME_PATH; 'lvm' volumes use the VG configured on the worker.
# Both use an empty path here — the worker resolves the actual path using its own settings.
_volume_path = ''
_volume_boot = volume.get('boot', False)
logger.info(
f"[VM={VirtualMachine['name']} id={new_VirtualMachine.id} host={selected_host.hostname}] "
f"Provisioning {_volume_type} volume '{volume['name']}' (boot={_volume_boot}, size={volume['size_gb']}GB)"
)
new_volume = Volume(
name=volume['name'],
image_id=_volume_image_id,
type=_volume_type,
path=_volume_path,
size_gb=volume['size_gb'],
vdc_id=request_vdc.id,
source=_volume_source,
boot=_volume_boot,
status="in-use",
)
logger.debug(f"Volume specification: {new_volume.to_json()}")
db.session.add(new_volume)
db.session.flush()
vol_mapping = VolumeWorkloadMapping(
name=f"Volume mapping for {new_VirtualMachine.name}",
volume_id=new_volume.id,
workload_id=new_VirtualMachine.id,
description=f"VM volume: {volume['name']}"
)
db.session.add(vol_mapping)
db.session.commit()
volumes.append(new_volume)
logger.info(f"Created volume {new_volume.id} (boot={_volume_boot}) with size {new_volume.size_gb}GB and mapping to workload {new_VirtualMachine.id}")
# Create network ports if networks are specified
this_vm_ports = []
this_vm_networks = []
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', [])):
logger.debug(f"Processing network {net_idx + 1}: {json.dumps(network, indent=2)}")
network_obj = Network.query.filter_by(id=(network['id']), deleted=False).first()
if not network_obj:
logger.error(f"Network {network['id']} not found or deleted")
return api_response(
success=False,
status=404,
message=f"Network '{network['id']}' not found or has been deleted",
error_type="NETWORK_NOT_FOUND"
)
this_vm_networks.append(network_obj.id)
if network_obj:
logger.info(f"Found network {network_obj.name} (ID: {network_obj.id})")
try:
port = network_obj.create_port(db.session, workload_id=new_VirtualMachine.id, use_dhcp_range=True)
except Exception as e:
logger.error(
f"Failed to create network port on network '{network_obj.name}' ({network_obj.id}) "
f"for VM '{VirtualMachine['name']}' (id={new_VirtualMachine.id}): {e}"
)
return api_response(
success=False,
status=500,
message=f"Failed to create network port on network '{network_obj.name}' for VM '{VirtualMachine['name']}': {e}",
error_type="PORT_CREATION_FAILED"
)
this_vm_ports.append(port)
logger.info(
f"Created network port for VM '{VirtualMachine['name']}' (workload_id={new_VirtualMachine.id}): "
f"port_id={port.id}, port_name={port.name}, mac={port.mac_address}, "
f"ip={port.ip_address}, network_id={network_obj.id}, network_name={network_obj.name}"
)
logger.debug(f"Network port details: {port.to_json()}")
# Resolve ovs_bridge: prefer network-level, fall back to host's first OVS bridge
bridge = network_obj.ovs_bridge
if not bridge and selected_host.ovs_bridges:
bridge = selected_host.ovs_bridges[0].ovs_bridge_name
logger.info(
f"Network '{network_obj.name}' has no ovs_bridge set; "
f"falling back to host bridge '{bridge}' (host={selected_host.id})"
)
if not bridge:
logger.error(f"No ovs_bridge configured for network {network_obj.name} ({network_obj.id}) on host {selected_host.id}, requested for VM '{VirtualMachine['name']}' (id={new_VirtualMachine.id})")
return api_response(
success=False,
status=400,
message=f"No OVS bridge configured for network '{network_obj.name}'. Configure network.ovs_bridge or add an OVS bridge to the host.",
error_type="NO_OVS_BRIDGE"
)
this_vm_port_ovs_bridges[port.id] = bridge
logger.info(f"Resolved ovs_bridge '{bridge}' for port {port.id} on network {network_obj.name}")
else:
logger.warning(f"Network with ID {network['id']} not found in database")
# === 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
try:
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}")
try:
vm_meta = VmMetadata(
workload_id=new_VirtualMachine.id,
instance_id=new_VirtualMachine.id,
hostname=VirtualMachine['name'],
local_hostname=VirtualMachine['name'],
user_data=VirtualMachine.get('user_data'),
public_keys=VirtualMachine.get('public_keys'),
)
db.session.add(vm_meta)
db.session.commit()
logger.info(f"Created VmMetadata for VM {new_VirtualMachine.id}")
except Exception as meta_exc:
logger.error(f"Failed to create VmMetadata for VM {new_VirtualMachine.id}: {meta_exc}")
# Prepare the payload for websocket server
logger.info("Constructing payload for websocket server")
# Resolve GPU PCI addresses from host's registered fixed resources.
# NVIDIA GPUs expose multiple PCI functions per physical GPU (e.g. c1:00.0 VGA
# and c1:00.1 Audio). All functions of a slot must be passed through together.
# We group by slot (address up to the last dot), pick vm_gpu slots, then include
# every registered function for each chosen slot.
gpu_pci_addresses = []
if gpu_requests:
# New model-based format: allocate via GPU_HOST table, scoped to selected host
allocated_gpu_hosts = schedule_gpu_for_vm(
gpu_requests,
new_VirtualMachine.id,
str(selected_host.id),
)
if allocated_gpu_hosts is None:
db.session.rollback()
return api_response(
success=False, status=400,
message=f"No eligible hosts for VM '{VirtualMachine['name']}': insufficient GPU capacity",
error_type="INSUFFICIENT_GPU_CAPACITY",
)
db.session.commit()
gpu_pci_addresses = [h.pci_address for h in allocated_gpu_hosts]
logger.info(f"Allocated GPU(s) via GPU_HOST: {gpu_pci_addresses}")
elif vm_gpu_count > 0:
# Legacy integer format: resolve from WorkloadHostFixedResource
all_gpu_resources = WorkloadHostFixedResource.query.filter_by(
workload_host_id=selected_host.id,
type="gpu",
deleted=False,
).all()
# Group addresses by slot: "c1:00.0" and "c1:00.1" → slot "c1:00"
from collections import defaultdict
slots: dict = defaultdict(list)
for r in all_gpu_resources:
if r.physical_address:
slot = r.physical_address.rsplit(".", 1)[0] # strip function digit
slots[slot].append(r.physical_address)
# Pick up to vm_gpu_count distinct slots
chosen_slots = list(slots.keys())[:vm_gpu_count]
if len(chosen_slots) < vm_gpu_count:
logger.warning(
f"VM '{VirtualMachine['name']}' requested {vm_gpu_count} GPU slot(s) but host only has "
f"{len(chosen_slots)} slot(s) registered — proceeding with what's available"
)
# Collect all functions for each chosen slot, sorted so .0 comes first
for slot in chosen_slots:
gpu_pci_addresses.extend(sorted(slots[slot]))
logger.info(f"GPU PCI addresses for VM: {gpu_pci_addresses}")
payload = {
"worker_id": selected_host.id,
"task_type": "virtual-machine-create",
"depends_on": this_host_sdn_update_task_id,
"job_details": {
"virtual_machine_name": VirtualMachine['name'],
"virtual_machine_id": new_VirtualMachine.id,
"desired_state": "running",
"virtual_machine_config": {
"memory": VirtualMachine['memory'],
"vcpu": VirtualMachine['vcpu'],
"graphics": {
"type": "vnc"
},
# Build volumes payload, setting source.type to 'http' when the image
# location looks like an HTTP(S) URL or image.location_type signals remote.
"volumes": [
{
"id": vol.id,
"name": vol.name,
"size_gb": vol.size_gb,
"type": vol.type,
"path": vol.path,
**(
(lambda v: {
"source": {
"type": (
"http"
if (getattr(v, 'image', None) and (getattr(v.image, 'location', '') or '').lower().startswith(('http://', 'https://')))
or (getattr(v, 'image', None) and (getattr(v.image, 'location_type', '') or '').lower() in ('http', 'remote', 'url'))
else "image"
),
"image_id": v.image.id,
"path": v.image.location,
"format": v.image.format,
"checksum": v.image.checksum,
}
} if getattr(v, 'image_id', None) else {"source": "None"}
)(vol)
)
}
for vol in volumes
],
"network_ports": [
{
"port_id": port.id,
"port_mac_address": port.mac_address,
"port_name": port.name,
"ovs_bridge": this_vm_port_ovs_bridges[port.id],
}
for port in this_vm_ports
],
**({"gpu_pci_addresses": gpu_pci_addresses} if gpu_pci_addresses else {}),
}
}
}
logger.debug(f"Complete websocket payload: {json.dumps(payload, indent=2)}")
# Send request to websocket server for VM launch
logger.info(f"Sending request to websocket server at {websocket_server_url}")
headers = {"Content-Type": "application/json"}
try:
websocket_server_response = requests.post(websocket_server_url, data=json.dumps(payload), headers=headers)
websocket_server_response_data = websocket_server_response.json()
logger.info(f"Websocket server response status: {websocket_server_response.status_code}")
logger.debug(f"Websocket server response data: {json.dumps(websocket_server_response_data, indent=2)}")
if websocket_server_response.status_code == 201:
logger.info(f"Successfully created task for VirtualMachine {new_VirtualMachine.id}")
new_VirtualMachine.set_status("allocated")
db.session.add(new_VirtualMachine)
db.session.commit()
logger.debug(f"Updated VirtualMachine status to 'allocated'")
else:
logger.error(f"Failed to create task for VirtualMachine {new_VirtualMachine.id}")
new_VirtualMachine.set_status("failed-allocation")
db.session.add(new_VirtualMachine)
db.session.commit()
logger.error(f"Updated VirtualMachine status to 'failed-allocation'")
except Exception as e:
logger.error(f"Exception occurred while contacting websocket server: {str(e)}")
new_VirtualMachine.set_status("failed-allocation")
db.session.add(new_VirtualMachine)
db.session.commit()
logger.error(f"Updated VirtualMachine status to 'failed-allocation' due to exception")
# Prepare the response to the API request
vm_response = {
"virtual_machine_id": new_VirtualMachine.id,
"virtual_machine_name": new_VirtualMachine.name,
"VirtualMachine_status": new_VirtualMachine.status,
}
logger.debug(f"Adding to client response: {vm_response}")
client_response.append(vm_response)
# logger.info(f"Completed processing all VirtualMachines. Returning response for {len(client_response)} VirtualMachines")
# # return client_response, 200
# return api_response(data=client_response, message="VirtualMachines scheduled")