Vm scheduling and deletion works..
This commit is contained in:
@@ -82,7 +82,7 @@ class IaaSClient:
|
||||
|
||||
|
||||
|
||||
def get_virtual_machines(self):
|
||||
def get_containers(self):
|
||||
return self._make_request('GET', 'workloads/containers')
|
||||
|
||||
def get_container(self,container_id):
|
||||
@@ -103,6 +103,7 @@ class IaaSClient:
|
||||
|
||||
|
||||
def get_virtual_machines(self):
|
||||
print("Getting VMs")
|
||||
return self._make_request('GET', 'workloads/virtual_machines')
|
||||
|
||||
def get_virtual_machine(self,virtual_machine_id):
|
||||
|
||||
@@ -276,11 +276,10 @@ def delete_container_workload(workload_id):
|
||||
"worker_id": _container.workload_host_id,
|
||||
"task_type": "container-delete",
|
||||
"job_details": {
|
||||
"tenancyID": 3,
|
||||
"containers": [
|
||||
"container": [
|
||||
{
|
||||
"container_id": _container.id,
|
||||
"deleted": "True"
|
||||
"desired_state": "deleted"
|
||||
}
|
||||
]
|
||||
}
|
||||
|
||||
@@ -40,11 +40,11 @@ def validate_payload(payload):
|
||||
except ValueError:
|
||||
raise ValueError("'vdc' must be a valid UUID.")
|
||||
|
||||
# Check for the 'virtual-machines' key and ensure it is a list with at least one VM
|
||||
# Check for the 'virtual-machines' key and ensure it is a list with at least one VirtualMachine
|
||||
if 'virtual-machines' not in payload:
|
||||
raise ValueError("'virtual-machines' key is missing.")
|
||||
if not isinstance(payload['virtual-machines'], list) or len(payload['virtual-machines']) == 0:
|
||||
raise ValueError("'virtual-machines' must be a list with at least one VM.")
|
||||
raise ValueError("'virtual-machines' must be a list with at least one VirtualMachine.")
|
||||
|
||||
# Initialize the sanitized payload
|
||||
sanitized_payload = {
|
||||
@@ -52,33 +52,33 @@ def validate_payload(payload):
|
||||
'virtual-machines': []
|
||||
}
|
||||
|
||||
# Validate each VM in the 'virtual-machines' list and build the sanitized version
|
||||
for vm in payload['virtual-machines']:
|
||||
if not isinstance(vm, dict):
|
||||
raise ValueError("Each VM must be a dictionary.")
|
||||
# Validate each VirtualMachine in the 'virtual-machines' list and build the sanitized version
|
||||
for VirtualMachine in payload['virtual-machines']:
|
||||
if not isinstance(VirtualMachine , dict):
|
||||
raise ValueError("Each VirtualMachine must be a dictionary.")
|
||||
|
||||
# Check for required keys: 'vm_name' and 'vm_config'
|
||||
if 'vm_name' not in vm:
|
||||
raise ValueError("VM is missing 'vm_name'.")
|
||||
if 'vm_config' not in vm:
|
||||
raise ValueError("VM is missing 'vm_config'.")
|
||||
# Check for required keys: 'virtual_machine_name' and 'virtual_machine_config'
|
||||
if 'virtual_machine_name' not in VirtualMachine :
|
||||
raise ValueError("VirtualMachine is missing 'virtual_machine_name'.")
|
||||
if 'virtual_machine_config' not in VirtualMachine :
|
||||
raise ValueError("VirtualMachine is missing 'virtual_machine_config'.")
|
||||
|
||||
# Validate 'vm_config'
|
||||
vm_config = vm['vm_config']
|
||||
if not isinstance(vm_config, dict):
|
||||
raise ValueError("'vm_config' must be a dictionary.")
|
||||
# Validate 'virtual_machine_config'
|
||||
virtual_machine_config = VirtualMachine['virtual_machine_config']
|
||||
if not isinstance(virtual_machine_config, dict):
|
||||
raise ValueError("'virtual_machine_config' must be a dictionary.")
|
||||
|
||||
# Validate 'memory' and 'vcpu'
|
||||
if 'memory' not in vm_config or not isinstance(vm_config['memory'], int) or vm_config['memory'] <= 0:
|
||||
if 'memory' not in virtual_machine_config or not isinstance(virtual_machine_config['memory'], int) or virtual_machine_config['memory'] <= 0:
|
||||
raise ValueError("'memory' must be a positive integer.")
|
||||
if 'vcpu' not in vm_config or not isinstance(vm_config['vcpu'], int) or vm_config['vcpu'] <= 0:
|
||||
if 'vcpu' not in virtual_machine_config or not isinstance(virtual_machine_config['vcpu'], int) or virtual_machine_config['vcpu'] <= 0:
|
||||
raise ValueError("'vcpu' must be a positive integer.")
|
||||
|
||||
# Validate 'volumes' if present
|
||||
if 'volumes' in vm_config:
|
||||
if not isinstance(vm_config['volumes'], list):
|
||||
if 'volumes' in virtual_machine_config:
|
||||
if not isinstance(virtual_machine_config['volumes'], list):
|
||||
raise ValueError("'volumes' must be a list.")
|
||||
for volume in vm_config['volumes']:
|
||||
for volume in virtual_machine_config['volumes']:
|
||||
if not isinstance(volume, dict):
|
||||
raise ValueError("Each volume must be a dictionary.")
|
||||
if 'name' not in volume or 'size_gb' not in volume:
|
||||
@@ -87,40 +87,40 @@ def validate_payload(payload):
|
||||
raise ValueError("'size_gb' must be a positive number.")
|
||||
|
||||
# Validate 'networks' if present
|
||||
if 'networks' in vm_config:
|
||||
if not isinstance(vm_config['networks'], list):
|
||||
if 'networks' in virtual_machine_config:
|
||||
if not isinstance(virtual_machine_config['networks'], list):
|
||||
raise ValueError("'networks' must be a list.")
|
||||
for network in vm_config['networks']:
|
||||
for network in virtual_machine_config['networks']:
|
||||
if not isinstance(network, dict):
|
||||
raise ValueError("Each network must be a dictionary.")
|
||||
if 'id' not in network:
|
||||
raise ValueError("Network is missing 'id'.")
|
||||
|
||||
# Initialize the sanitized VM
|
||||
sanitized_vm = {
|
||||
'vm_name': vm['vm_name'],
|
||||
'vm_config': {
|
||||
'memory': vm_config['memory'],
|
||||
'vcpu': vm_config['vcpu'],
|
||||
'volumes': vm_config.get('volumes', []),
|
||||
'networks': vm_config.get('networks', [])
|
||||
# Initialize the sanitized VirtualMachine
|
||||
sanitized_VirtualMachine = {
|
||||
'virtual_machine_name': VirtualMachine['virtual_machine_name'],
|
||||
'virtual_machine_config': {
|
||||
'memory': virtual_machine_config['memory'],
|
||||
'vcpu': virtual_machine_config['vcpu'],
|
||||
'volumes': virtual_machine_config.get('volumes', []),
|
||||
'networks': virtual_machine_config.get('networks', [])
|
||||
}
|
||||
}
|
||||
|
||||
# Add the sanitized VM to the sanitized payload
|
||||
sanitized_payload['virtual-machines'].append(sanitized_vm)
|
||||
# Add the sanitized VirtualMachine to the sanitized payload
|
||||
sanitized_payload['virtual-machines'].append(sanitized_VirtualMachine)
|
||||
|
||||
return sanitized_payload
|
||||
|
||||
@api_bp.route('/workloads/virtual_machines', methods=['POST'])
|
||||
def add_vm_workload1():
|
||||
def add_VirtualMachine_workload1():
|
||||
request_data = request.json
|
||||
logger.debug(request_data)
|
||||
|
||||
# Validate incoming data
|
||||
validated_data = validate_payload(request_data)
|
||||
|
||||
# TODO - Check if this user has permissions to VIEW and LAUNCH_VM in this VDC
|
||||
# TODO - Check if this user has permissions to VIEW and LAUNCH_VirtualMachine in this VDC
|
||||
request_vdc = VirtualDataCenter.query.filter_by(id=uuid.UUID(validated_data['vdc'])).first()
|
||||
|
||||
if request_vdc is None:
|
||||
@@ -135,27 +135,27 @@ def add_vm_workload1():
|
||||
return f"No hosts found in region {request_vdc.region.name} to place this request", 400
|
||||
|
||||
client_response = []
|
||||
for vm in validated_data['virtual-machines']:
|
||||
for VirtualMachine in validated_data['virtual-machines']:
|
||||
# TODO - This is where we would slot in the 'placement' routine
|
||||
random_host = random.choice(all_workload_hosts)
|
||||
logger.info(f"Randomly selected host: {random_host}")
|
||||
|
||||
# Create the Workload instance in the database
|
||||
new_vm = Workload(
|
||||
name=vm['vm_name'],
|
||||
workload_type="VM", # Set workload_type to "VM"
|
||||
new_VirtualMachine = Workload(
|
||||
name=VirtualMachine['virtual_machine_name'],
|
||||
workload_type="VirtualMachine", # Set workload_type to VirtualMachine
|
||||
vdc_id=request_vdc.id,
|
||||
status="pending-allocation",
|
||||
launch_params=json.dumps(vm),
|
||||
launch_params=json.dumps(VirtualMachine ),
|
||||
workload_host_id=random_host.id
|
||||
)
|
||||
db.session.add(new_vm)
|
||||
db.session.add(new_VirtualMachine)
|
||||
db.session.commit()
|
||||
logger.debug(f"VM {new_vm.id} workload added to DB")
|
||||
logger.debug(f"VirtualMachine {new_VirtualMachine.id} workload added to DB")
|
||||
|
||||
# Create volumes in the database
|
||||
volumes = []
|
||||
for volume in vm['vm_config'].get('volumes', []):
|
||||
for volume in VirtualMachine['virtual_machine_config'].get('volumes', []):
|
||||
new_volume = Volume(
|
||||
name=volume['name'],
|
||||
size_gb=volume['size_gb'],
|
||||
@@ -168,30 +168,27 @@ def add_vm_workload1():
|
||||
|
||||
# Create network ports if networks are specified
|
||||
networks = []
|
||||
for network in vm['vm_config'].get('networks', []):
|
||||
for network in VirtualMachine['virtual_machine_config'].get('networks', []):
|
||||
network_obj = Network.query.filter_by(id=uuid.UUID(network['id'])).first()
|
||||
if network_obj:
|
||||
port = network_obj.create_port(db, workload_id=new_vm.id)
|
||||
port = network_obj.create_port(db.session, workload_id=new_VirtualMachine.id)
|
||||
networks.append(port)
|
||||
|
||||
# Prepare the VM details for the outbound payload
|
||||
vm_details = {
|
||||
"vm_name": vm['vm_name'],
|
||||
"vm_id": new_vm.id,
|
||||
"desired_state": "running",
|
||||
"vm_config": {
|
||||
"memory": vm['vm_config']['memory'],
|
||||
"vcpu": vm['vm_config']['vcpu'],
|
||||
"volumes": [{"id": vol.id, "size_gb": vol.size_gb, "format": "qcow2"} for vol in volumes],
|
||||
"networks": [{"id": net.id, "mac_address": net.mac_address} for net in networks]
|
||||
}
|
||||
}
|
||||
|
||||
# Send the request off to the websocket server to have this VM launched
|
||||
# Send the request off to the websocket server to have this VirtualMachine launched
|
||||
payload = {
|
||||
"worker_id": random_host.id,
|
||||
"task_type": "vm-create",
|
||||
"job_details": vm_details,
|
||||
"task_type": "virtual-machine-create",
|
||||
"job_details": {
|
||||
"virtual_machine_name": VirtualMachine['virtual_machine_name'],
|
||||
"virtual_machine_id": new_VirtualMachine.id,
|
||||
"desired_state": "running",
|
||||
"virtual_machine_config": {
|
||||
"memory": VirtualMachine['virtual_machine_config']['memory'],
|
||||
"vcpu": VirtualMachine['virtual_machine_config']['vcpu'],
|
||||
"volumes": [{"id": vol.id, "size_gb": vol.size_gb, "format": "qcow2"} for vol in volumes],
|
||||
"networks": [{"id": net.id, "mac_address": net.mac_address} for net in networks]
|
||||
}
|
||||
},
|
||||
}
|
||||
logger.debug(f"Sending payload to websocket server {payload}")
|
||||
headers = {"Content-Type": "application/json"}
|
||||
@@ -200,65 +197,98 @@ def add_vm_workload1():
|
||||
websocket_server_response_data = websocket_server_response.json()
|
||||
|
||||
if websocket_server_response.status_code == 201:
|
||||
logger.info(f"Task created successfully for VM {new_vm.id}")
|
||||
new_vm.status = "pending-allocated"
|
||||
db.session.add(new_vm)
|
||||
logger.info(f"Task created successfully for VirtualMachine {new_VirtualMachine.id}")
|
||||
new_VirtualMachine.status = "pending-allocated"
|
||||
db.session.add(new_VirtualMachine)
|
||||
db.session.commit()
|
||||
else:
|
||||
logger.error(f"Creation request failed for VM {new_vm.id}")
|
||||
new_vm.status = "failed-allocation"
|
||||
db.session.add(new_vm)
|
||||
logger.error(f"Creation request failed for VirtualMachine {new_VirtualMachine.id}")
|
||||
new_VirtualMachine.status = "failed-allocation"
|
||||
db.session.add(new_VirtualMachine)
|
||||
db.session.commit()
|
||||
|
||||
# Prepare the response to the API request
|
||||
client_response.append({
|
||||
"vm_id": new_vm.id,
|
||||
"vm_name": new_vm.name,
|
||||
"vm_status": new_vm.status,
|
||||
"virtual_machine_id": new_VirtualMachine.id,
|
||||
"virtual_machine_name": new_VirtualMachine.name,
|
||||
"VirtualMachine_status": new_VirtualMachine.status,
|
||||
})
|
||||
|
||||
return client_response, 200
|
||||
|
||||
|
||||
@api_bp.route('/workloads/virtual_machines/status_update/<virtual_machine_id>', methods=['PUT'])
|
||||
def edit_VirtualMachine_workload(virtual_machine_id):
|
||||
data = request.json
|
||||
new_status = data.get('new_status')
|
||||
|
||||
logger.info(f"Updating VirtualMachine ID:{virtual_machine_id} to new status {new_status}")
|
||||
workload = Workload.query.filter_by(id=uuid.UUID(virtual_machine_id), workload_type="VirtualMachine").first_or_404()
|
||||
workload.status = new_status
|
||||
db.session.commit()
|
||||
return jsonify(success=True)
|
||||
|
||||
@api_bp.route('/workloads/vm_test', methods=['GET'])
|
||||
def add_vm_workload():
|
||||
@api_bp.route('/workloads/virtual_machines/<workload_id>', methods=['GET'])
|
||||
def get_VirtualMachine_workload(workload_id):
|
||||
try:
|
||||
workload_uuid = uuid.UUID(workload_id)
|
||||
except ValueError:
|
||||
return jsonify({"error": "Invalid workload ID"}), 404
|
||||
|
||||
workload = Workload.query.filter(
|
||||
Workload.id == workload_uuid,
|
||||
Workload.workload_type == "VirtualMachine",
|
||||
Workload.deleted == False
|
||||
).first_or_404()
|
||||
logger.info(f"{workload.workload_type} {workload.deleted}")
|
||||
return jsonify(workload.to_json())
|
||||
|
||||
test_vm_name="cory_vm"
|
||||
|
||||
vm_details = {
|
||||
"vm_name": test_vm_name,
|
||||
"vm_id": uuid.uuid4(),
|
||||
"desired_state": "running",
|
||||
"vm_config": {
|
||||
"memory": 2048,
|
||||
"vcpu": 2,
|
||||
"volumes": [
|
||||
{"id": uuid.uuid4(), "size_gb": 20, "format": "qcow2"},
|
||||
],
|
||||
"networks": [
|
||||
{"name": "default", "mac_address": "52:54:00:12:34:58"}
|
||||
],
|
||||
@api_bp.route('/workloads/virtual_machines/<workload_id>', methods=['DELETE'])
|
||||
def delete_VirtualMachine_workload(workload_id):
|
||||
try:
|
||||
workload_uuid = uuid.UUID(workload_id)
|
||||
except ValueError:
|
||||
return jsonify({"error": "Invalid workload ID"}), 404
|
||||
|
||||
_VirtualMachine = Workload.query.filter(
|
||||
Workload.id == workload_uuid,
|
||||
Workload.workload_type == "VirtualMachine",
|
||||
Workload.deleted == False
|
||||
).first_or_404()
|
||||
|
||||
payload = {
|
||||
"worker_id": _VirtualMachine.workload_host_id,
|
||||
"task_type": "virtual-machine-delete",
|
||||
"job_details": {
|
||||
"virtual_machine_id": _VirtualMachine.id,
|
||||
"desired_state": "deleted"
|
||||
}
|
||||
}
|
||||
payload = {
|
||||
"worker_id": "4426077a-9973-4370-b1c4-98362fef9828",
|
||||
"task_type": "vm-create",
|
||||
"job_details": vm_details,
|
||||
}
|
||||
|
||||
logger.debug(f"Sending payload to websocket server {payload}")
|
||||
headers = {"Content-Type": "application/json"}
|
||||
|
||||
websocket_server_response = requests.post(websocket_server_url, data=json.dumps(payload), headers=headers)
|
||||
websocket_server_response_data = websocket_server_response.json()
|
||||
|
||||
|
||||
if websocket_server_response.status_code == 201:
|
||||
logger.info(f"Task created successfully for vm {test_vm_name}")
|
||||
|
||||
logger.info(f"Task created successfully for VirtualMachine {_VirtualMachine.id}")
|
||||
_VirtualMachine.status = "pending-deleted"
|
||||
db.session.add(_VirtualMachine)
|
||||
db.session.commit()
|
||||
else:
|
||||
logger.error(f"Creation request failed for vm {test_vm_name}")
|
||||
logger.error(f"Creation request failed for VirtualMachine {_VirtualMachine.id}")
|
||||
_VirtualMachine.status = "failed-deleted"
|
||||
db.session.add(_VirtualMachine)
|
||||
db.session.commit()
|
||||
|
||||
_VirtualMachine.soft_delete()
|
||||
logger.info(f"Deleted VirtualMachine {_VirtualMachine.id}")
|
||||
return jsonify({'message': 'VirtualMachine workload deleted successfully'}), 200
|
||||
|
||||
|
||||
return "ok",200
|
||||
|
||||
@api_bp.route('/workloads/virtual_machines', methods=['GET'])
|
||||
def get_VirtualMachine_workloads():
|
||||
workloads = Workload.query.filter(
|
||||
Workload.workload_type == "VirtualMachine",
|
||||
Workload.deleted == False
|
||||
).all()
|
||||
return jsonify([workload.to_json() for workload in workloads])
|
||||
|
||||
-400
@@ -1,400 +0,0 @@
|
||||
import requests
|
||||
from datetime import datetime
|
||||
import os
|
||||
from colorlog import ColoredFormatter
|
||||
import socketio
|
||||
import argparse
|
||||
import logging
|
||||
import json
|
||||
import asyncio
|
||||
import docker
|
||||
|
||||
from worker_tasks.report import ReportTask
|
||||
from worker_tasks.ping import PingTask
|
||||
from worker_tasks.file_presence import FilePresenceTask
|
||||
from worker_tasks.container import ContainerTask
|
||||
|
||||
from dotenv import load_dotenv
|
||||
|
||||
# Configure logging
|
||||
DEBUG_SOCKETIO = False
|
||||
|
||||
# Custom filter to include the function name
|
||||
class FunctionNameFilter(logging.Filter):
|
||||
def filter(self, record):
|
||||
record.funcName = record.funcName if hasattr(record, 'funcName') else '<unknown>'
|
||||
return True
|
||||
|
||||
# Define the colorized log format
|
||||
log_format = (
|
||||
"%(log_color)s%(asctime)s - %(levelname)s - %(funcName)s - %(message)s"
|
||||
)
|
||||
date_format = "%Y-%m-%d %H:%M:%S"
|
||||
|
||||
# Configure the formatter with colors
|
||||
formatter = ColoredFormatter(
|
||||
log_format,
|
||||
datefmt=date_format,
|
||||
log_colors={
|
||||
"DEBUG": "cyan",
|
||||
"INFO": "green",
|
||||
"WARNING": "yellow",
|
||||
"ERROR": "red",
|
||||
"CRITICAL": "bold_red",
|
||||
},
|
||||
)
|
||||
|
||||
# Configure the handler
|
||||
handler = logging.StreamHandler()
|
||||
handler.setFormatter(formatter)
|
||||
|
||||
# Configure the logger
|
||||
logger = logging.getLogger(__name__)
|
||||
logger.setLevel(logging.DEBUG)
|
||||
logger.addFilter(FunctionNameFilter())
|
||||
logger.addHandler(handler)
|
||||
|
||||
class WorkerClient:
|
||||
def __init__(self, worker_ID, worker_secret, server_URL):
|
||||
self.worker_id = worker_ID
|
||||
self.worker_secret = worker_secret
|
||||
self.server_url = server_URL
|
||||
if DEBUG_SOCKETIO:
|
||||
self.sio = socketio.AsyncClient(logger=logger, engineio_logger=logger)
|
||||
else:
|
||||
self.sio = socketio.AsyncClient(logger=False, engineio_logger=False)
|
||||
self.joined_server = False
|
||||
|
||||
# Bind events
|
||||
self.sio.on("connect", self.on_connect)
|
||||
self.sio.on("connect_error", self.on_connect_error)
|
||||
self.sio.on("disconnect", self.on_disconnect)
|
||||
self.sio.on("task", self.handle_task)
|
||||
self.sio.on("join_accept", self.on_join_accept)
|
||||
self.sio.on("join_reject", self.on_join_reject)
|
||||
self.sio.on("message", self.on_message)
|
||||
self.sio.on("*", self.debug_all_events)
|
||||
|
||||
async def on_connect(self):
|
||||
logger.info("Connected to the server, requesting to join.")
|
||||
await self.send_join_request()
|
||||
|
||||
async def on_connect_error(self, data):
|
||||
logger.error(f"Connection failed: {data}")
|
||||
self.joined_server = False
|
||||
|
||||
async def on_disconnect(self):
|
||||
logger.info("Disconnected from the server.")
|
||||
self.joined_server = False
|
||||
|
||||
async def on_message(self, data):
|
||||
logger.info(f"Message received: {data}")
|
||||
|
||||
async def on_join_accept(self, data):
|
||||
logger.info("Join accepted.")
|
||||
self.joined_server = True
|
||||
|
||||
async def on_join_reject(self, data):
|
||||
logger.error("Join rejected.")
|
||||
self.joined_server = False
|
||||
exit()
|
||||
|
||||
async def send_join_request(self):
|
||||
"""Notify the server about this worker (Initial join)."""
|
||||
logger.info("Sending join request")
|
||||
if not self.joined_server:
|
||||
await self.sio.emit("join_request", {"worker_id": self.worker_id, "worker_secret": self.worker_secret})
|
||||
logger.info(f"Worker {self.worker_id} asked to join the server.")
|
||||
|
||||
def handle_task(self, data):
|
||||
"""Handle incoming tasks from the server."""
|
||||
if not self.joined_server:
|
||||
logger.error("Not joined to server. ignoring task")
|
||||
return
|
||||
task_type = data["type"]
|
||||
task_id = data["task_id"]
|
||||
job_details = json.loads(data["job_details"])
|
||||
task_worker_id = data["worker_id"]
|
||||
|
||||
logger.info(f"Received task {task_id} of type '{task_type}' with job_details: {job_details}")
|
||||
|
||||
|
||||
if task_type == "report":
|
||||
result = ReportTask("", logger).Execute()
|
||||
elif task_type == "ping":
|
||||
result = PingTask(job_details, logger).Execute()
|
||||
elif task_type == "file_presence":
|
||||
result = FilePresenceTask(job_details, logger).Execute()
|
||||
elif task_type == "container-create":
|
||||
result = ContainerTask(job_details, logger).Create()
|
||||
elif task_type == "container-delete":
|
||||
result = ContainerTask(job_details, logger).Delete()
|
||||
else:
|
||||
raise ValueError(f"Unknown task type: {task_type}")
|
||||
try:
|
||||
self.send_task_result(task_id, result, task_worker_id)
|
||||
except Exception as e:
|
||||
logger.error(f"Error processing task {task_id}: {e}")
|
||||
self.send_task_result(task_id, {"success": False, "response": str(e)}, task_worker_id)
|
||||
|
||||
def send_task_result(self, task_id, result, worker_id):
|
||||
"""Send task result back to the server."""
|
||||
self.sio.emit("ack", {"task_id": task_id, "worker_id": worker_id, "result": result})
|
||||
logger.info(f"Sent result for task {task_id}: {result}")
|
||||
|
||||
def debug_all_events(self, event, data):
|
||||
"""Debug all incoming data."""
|
||||
logger.debug(f"Event: {event} | Data: {json.dumps(data, indent=2)}")
|
||||
|
||||
async def start(self):
|
||||
"""Worker connects to the API server and processes tasks."""
|
||||
try:
|
||||
logger.info(f"Worker {self.worker_id} connecting to {self.server_url}...")
|
||||
await self.sio.connect(self.server_url)
|
||||
await self.sio.wait()
|
||||
except Exception as e:
|
||||
logger.error(f"Error: {e}")
|
||||
await self.sio.disconnect()
|
||||
|
||||
async def watch_docker_events(worker_id, server_url):
|
||||
"""Watch Docker events and send them to the server for containers managed by this worker."""
|
||||
logger.info("Starting Docker event watcher...")
|
||||
docker_client = docker.from_env()
|
||||
sio = socketio.AsyncClient()
|
||||
|
||||
try:
|
||||
for event in await asyncio.to_thread(docker_client.events, decode=True):
|
||||
logger.debug(f"Caught docker event {event}")
|
||||
if event['Type'] == 'container' and event['Action'] in ['start', 'stop']:
|
||||
container_id = event['id']
|
||||
try:
|
||||
# Fetch the container details to check its labels
|
||||
container = await asyncio.to_thread(docker_client.containers.get, container_id)
|
||||
labels = container.attrs['Config']['Labels']
|
||||
|
||||
# Check if the container is managed by this worker
|
||||
if labels and labels.get("managed_by") == "worker_agent":
|
||||
logger.info(f"Detected Docker event for managed container: {event}")
|
||||
await send_event_to_server(sio, worker_id, server_url, event)
|
||||
except docker.errors.NotFound:
|
||||
logger.warning(f"Container {container_id} not found. Skipping event.")
|
||||
except Exception as e:
|
||||
logger.error(f"Error fetching container details: {e}")
|
||||
except Exception as e:
|
||||
logger.error(f"Error watching Docker events: {e}")
|
||||
|
||||
async def send_event_to_server(sio, worker_id, server_url, event):
|
||||
"""Send Docker event to the server over WebSocket."""
|
||||
try:
|
||||
await sio.connect(server_url)
|
||||
await sio.emit("docker_event", {"worker_id": worker_id, "event": event})
|
||||
logger.info(f"Sent Docker event to server: {event}")
|
||||
except Exception as e:
|
||||
logger.error(f"Failed to send Docker event to server: {e}")
|
||||
|
||||
def validate_json_input(text_input):
|
||||
try:
|
||||
# Step 1: Validate that it's valid JSON
|
||||
data = json.loads(text_input)
|
||||
except json.JSONDecodeError as e:
|
||||
raise ValueError(f"Invalid JSON format: {e}")
|
||||
|
||||
# Step 2: Ensure 'tenancyID' is present
|
||||
if 'tenancyID' not in data:
|
||||
raise ValueError("Key 'tenancyID' is missing.")
|
||||
|
||||
# Step 3: Ensure 1 or more containers are present
|
||||
if 'containers' not in data or not isinstance(data['containers'], list) or len(data['containers']) == 0:
|
||||
raise ValueError("Key 'containers' is missing, not a list, or empty.")
|
||||
|
||||
# Step 4: For each container, ensure 'container_name' and 'docker_image' are present
|
||||
for container in data['containers']:
|
||||
if not isinstance(container, dict):
|
||||
raise ValueError("Container is not a dictionary.")
|
||||
if 'container_name' not in container or 'docker_image' not in container:
|
||||
raise ValueError("Container is missing 'container_name' or 'docker_image'.")
|
||||
|
||||
# If all validations pass, return the parsed JSON object
|
||||
return data
|
||||
|
||||
def container_test():
|
||||
logger.info("Running container test")
|
||||
payload_text = '''{
|
||||
"tenancyID": "tenant123a",
|
||||
"containers": [
|
||||
{
|
||||
"container_id": "af47632e-43b7-4474-82fa-475a24fe88d5",
|
||||
"docker_image": "nginx",
|
||||
"cpu_shares": 1,
|
||||
"mem_limit": 128,
|
||||
"container_name": "web-server1",
|
||||
"command": ["nginx", "-g", "daemon off;"],
|
||||
"working_dir": "/usr/share/nginx/html"
|
||||
},{
|
||||
"container_id": "82c1c6ae-9dd4-41c0-ba05-45a086d58cf5",
|
||||
"docker_image": "python:3.9-alpine",
|
||||
"container_name": "python-webserver",
|
||||
"cpu_shares": 1,
|
||||
"mem_limit": 128,
|
||||
"command": ["python", "-m", "http.server","82"],
|
||||
"environment": {
|
||||
"DELAY_START_MSEC": "2000"
|
||||
},
|
||||
"restart_policy": {
|
||||
"Name": "always"
|
||||
}
|
||||
}
|
||||
]
|
||||
|
||||
}'''
|
||||
|
||||
# Parse payload
|
||||
payload_json = validate_json_input(payload_text)
|
||||
logger.debug(f"Container test payload {payload_json}")
|
||||
# Initialize and execute ContainerTask
|
||||
container_task = ContainerTask(payload_json, logger)
|
||||
result = container_task.Execute()
|
||||
logger.info(result)
|
||||
|
||||
def enroll_worker(region_id):
|
||||
# Configuration
|
||||
API_URL = "http://127.0.0.1:5000/api/workload_hosts/enroll"
|
||||
|
||||
# Gather system information
|
||||
def get_system_info():
|
||||
try:
|
||||
with open('/sys/class/dmi/id/sys_vendor', 'r') as f:
|
||||
system_manufacturer = f.read().strip()
|
||||
except:
|
||||
system_manufacturer = "Unknown"
|
||||
|
||||
try:
|
||||
with open('/sys/class/dmi/id/product_name', 'r') as f:
|
||||
system_model = f.read().strip()
|
||||
except:
|
||||
system_model = "Unknown"
|
||||
|
||||
try:
|
||||
with open('/sys/class/dmi/id/product_serial', 'r') as f:
|
||||
physical_identifier = f.read().strip()
|
||||
except:
|
||||
physical_identifier = "Unknown"
|
||||
|
||||
try:
|
||||
dcim_identifier = os.popen('dmidecode -s system-uuid').read().strip()
|
||||
except:
|
||||
dcim_identifier = "Unknown"
|
||||
|
||||
installed_date = datetime.utcnow().strftime("%Y-%m-%dT%H:%M:%SZ")
|
||||
hostname = os.uname().nodename
|
||||
|
||||
return {
|
||||
"system_manufacturer": system_manufacturer,
|
||||
"system_model": system_model,
|
||||
"physical_identifier": physical_identifier,
|
||||
"dcim_identifier": dcim_identifier,
|
||||
"installed_date": installed_date,
|
||||
"hostname": hostname
|
||||
}
|
||||
|
||||
# Construct JSON payload
|
||||
system_info = get_system_info()
|
||||
payload = {
|
||||
"region_id": region_id,
|
||||
"system_manufacturer": system_info["system_manufacturer"],
|
||||
"system_model": system_info["system_model"],
|
||||
"physical_identifier": system_info["physical_identifier"],
|
||||
"dcim_identifier": system_info["dcim_identifier"],
|
||||
"installed_date": system_info["installed_date"],
|
||||
"hostname": system_info["hostname"]
|
||||
}
|
||||
|
||||
# Send enrollment request
|
||||
response = requests.post(API_URL, json=payload, headers={"Content-Type": "application/json"})
|
||||
|
||||
# Handle response
|
||||
if response.status_code != 201:
|
||||
print(f"Enrollment failed: {response.status_code} - {response.text}")
|
||||
exit(1)
|
||||
|
||||
# Parse response JSON
|
||||
try:
|
||||
response_data = response.json()
|
||||
worker_id = response_data.get("worker_id")
|
||||
worker_secret = response_data.get("worker_secret")
|
||||
|
||||
if not worker_id or not worker_secret:
|
||||
print("Enrollment failed: Worker_ID or Worker_Secret not found in response")
|
||||
exit(1)
|
||||
|
||||
# Ensure .env file exists before appending
|
||||
if not os.path.exists(".env"):
|
||||
open(".env", "w").close() # Create an empty .env file if it doesn't exist
|
||||
|
||||
# Append worker credentials to .env
|
||||
with open(".env", "a") as env_file:
|
||||
env_file.write(f"WORKER_ID={worker_id}\n")
|
||||
env_file.write(f"WORKER_SECRET={worker_secret}\n")
|
||||
|
||||
logger.info("Enrollment successful. Worker_ID and Worker_Secret written to .env file.")
|
||||
return response_data
|
||||
|
||||
except json.JSONDecodeError:
|
||||
print("Enrollment failed: Invalid JSON response from server")
|
||||
exit(1)
|
||||
|
||||
def main():
|
||||
load_dotenv()
|
||||
|
||||
parser = argparse.ArgumentParser(description="Worker CLI")
|
||||
parser.add_argument("--server-url", help="WebSocket Server URL (e.g., http://localhost:5000)", default=os.getenv("SERVER_URL"))
|
||||
parser.add_argument("--region-id", help="Region ID (Only required during enrollment)", default=os.getenv("REGION_ID"))
|
||||
parser.add_argument("--container-test", action="store_true", help="Run container test")
|
||||
args = parser.parse_args()
|
||||
|
||||
# Verbose logging of environment variables
|
||||
logger.info("Starting Worker CLI")
|
||||
logger.info(f"SERVER_URL: {args.server_url if args.server_url else 'Not Provided'}")
|
||||
logger.info(f"REGION_ID: {args.region_id if args.region_id else 'Not Provided'}")
|
||||
|
||||
if args.container_test:
|
||||
logger.info("Running container test...")
|
||||
container_test()
|
||||
exit()
|
||||
|
||||
# Ensure required arguments are provided
|
||||
if not args.server_url:
|
||||
parser.error("--server-url is required when --container-test is not used and not defined in .env.")
|
||||
|
||||
# Enroll worker only if WORKER_ID and WORKER_SECRET are not in .env
|
||||
worker_secret = os.getenv("WORKER_SECRET")
|
||||
worker_id = os.getenv("WORKER_ID")
|
||||
if not worker_id or not worker_secret:
|
||||
logger.info("Worker ID or Worker Secret missing. Enrolling worker...")
|
||||
# Region ID is only required for enrollment
|
||||
if not args.region_id:
|
||||
parser.error("--region-id must be provided either in the .env file or via CLI.")
|
||||
|
||||
enroll_worker(args.region_id)
|
||||
load_dotenv()
|
||||
else:
|
||||
logger.info("Worker ID and Worker Secret found. Skipping enrollment.")
|
||||
|
||||
worker_ID = os.getenv("WORKER_ID")
|
||||
worker_secret = os.getenv("WORKER_SECRET")
|
||||
server_URL = os.getenv("SERVER_URL") if os.getenv("SERVER_URL") else args.server_url
|
||||
|
||||
# Create and start the worker
|
||||
logger.info("Starting Worker with the following params...")
|
||||
logger.info(f"WORKER_ID: {worker_ID}")
|
||||
logger.info(f"WORKER_SECRET: {worker_secret}")
|
||||
logger.info(f"SERVER_URL: {server_URL}")
|
||||
|
||||
worker = WorkerClient(worker_ID, worker_secret, server_URL)
|
||||
|
||||
asyncio.run(worker.start())
|
||||
logger.debug("Done with starting working")
|
||||
asyncio.run(watch_docker_events(worker_ID, server_URL))
|
||||
logger.debug("Done with starting docker")
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
@@ -24,6 +24,8 @@ def render():
|
||||
render_detail_view_network(resource)
|
||||
elif resource_type == 'container':
|
||||
render_detail_view_container(resource)
|
||||
elif resource_type == 'VirtualMachine':
|
||||
render_detail_view_virtual_machine(resource)
|
||||
|
||||
def render_detail_view_container(container):
|
||||
|
||||
@@ -207,39 +209,39 @@ def render_detail_view_region(region):
|
||||
|
||||
import streamlit as st
|
||||
|
||||
def render_detail_view_vm(vm):
|
||||
def render_detail_view_virtual_machine(vm):
|
||||
"""
|
||||
Renders the details of a specific VM workload.
|
||||
"""
|
||||
vm = st.session_state.client.get_vm(vm['id'])
|
||||
VirtualMachine = st.session_state.client.get_virtual_machine(vm['id'])
|
||||
if not vm:
|
||||
st.error("VM workload not found.")
|
||||
return
|
||||
|
||||
st.title(f"VM Workload: {vm['vm_name']}")
|
||||
print(vm)
|
||||
st.title(f"VM Workload: {VirtualMachine['name']}")
|
||||
|
||||
st.header("VM Workload Details")
|
||||
col1, col2 = st.columns(2)
|
||||
with col1:
|
||||
st.markdown(f"**ID:** `{vm['id']}`")
|
||||
st.markdown(f"**Created:** {format_timestamp(vm.get('created_at', 'Unknown'))}")
|
||||
st.markdown(f"**VDC ID:** {vm['vdc_id']}")
|
||||
st.markdown(f"**Workload Host ID:** {vm.get('workload_host_id', 'N/A')}")
|
||||
st.markdown(f"**ID:** `{VirtualMachine['id']}`")
|
||||
st.markdown(f"**Created:** {format_timestamp(VirtualMachine.get('created_at', 'Unknown'))}")
|
||||
st.markdown(f"**VDC ID:** {VirtualMachine['vdc_id']}")
|
||||
st.markdown(f"**Workload Host ID:** {VirtualMachine.get('workload_host_id', 'N/A')}")
|
||||
with col2:
|
||||
if vm.get('updated_at'):
|
||||
st.markdown(f"**Last Updated:** {format_timestamp(vm['updated_at'])}")
|
||||
st.markdown(f"**Status:** {vm['status']}")
|
||||
st.markdown(f"**Launch Parameters:** {vm.get('launch_params', 'N/A')}")
|
||||
if VirtualMachine.get('updated_at'):
|
||||
st.markdown(f"**Last Updated:** {format_timestamp(VirtualMachine['updated_at'])}")
|
||||
st.markdown(f"**Status:** {VirtualMachine['status']}")
|
||||
st.markdown(f"**Launch Parameters:** {VirtualMachine.get('launch_params', 'N/A')}")
|
||||
|
||||
# Edit VM Workload Form
|
||||
with st.expander("Edit VM Workload"):
|
||||
with st.form("edit_vm"):
|
||||
new_status = st.selectbox("Status", ["running", "stopped", "pending", "failed-deleted", "pending-deleted", "pending-allocation", "failed-allocation", "pending-allocated"], index=["running", "stopped", "pending", "failed-deleted", "pending-deleted", "pending-allocation", "failed-allocation", "pending-allocated"].index(vm['status']))
|
||||
new_launch_params = st.text_input("Launch Parameters", value=vm.get('launch_params', ''))
|
||||
new_status = st.selectbox("Status", ["running", "stopped", "pending", "failed-deleted", "pending-deleted", "pending-allocation", "failed-allocation", "pending-allocated"], index=["running", "stopped", "pending", "failed-deleted", "pending-deleted", "pending-allocation", "failed-allocation", "pending-allocated"].index(VirtualMachine['status']))
|
||||
new_launch_params = st.text_input("Launch Parameters", value=VirtualMachine.get('launch_params', ''))
|
||||
submit = st.form_submit_button("Update VM Workload")
|
||||
|
||||
if submit:
|
||||
result = st.session_state.client.edit_vm(vm['id'], {
|
||||
result = st.session_state.client.edit_virtual_machine(VirtualMachine['id'], {
|
||||
"status": new_status,
|
||||
"launch_params": new_launch_params
|
||||
})
|
||||
@@ -247,15 +249,15 @@ def render_detail_view_vm(vm):
|
||||
st.success("VM workload updated successfully!")
|
||||
st.session_state.selected_resource = vm
|
||||
st.session_state.view_type = 'detail'
|
||||
st.session_state.selected_resource_type = 'vm'
|
||||
st.session_state.selected_resource_type = 'VirtualMachine'
|
||||
st.rerun()
|
||||
|
||||
# Delete VM Workload Button
|
||||
if st.button("Delete VM Workload"):
|
||||
result = st.session_state.client.delete_vm(vm['id'])
|
||||
result = st.session_state.client.delete_virtual_machine(VirtualMachine['id'])
|
||||
if result:
|
||||
st.success("VM workload deleted successfully!")
|
||||
st.session_state.selected_resource = None
|
||||
st.session_state.view_type = 'list'
|
||||
st.session_state.selected_resource_type = 'vm'
|
||||
st.session_state.selected_resource_type = 'VirtualMachine'
|
||||
st.rerun()
|
||||
|
||||
@@ -74,11 +74,57 @@ def render():
|
||||
result = st.session_state.client.create_virtual_machine({
|
||||
"vdc": vdc_id,
|
||||
"virtual-machines": [{
|
||||
"vm_name": vm_name,
|
||||
"vm_config": vm_config
|
||||
"virtual_machine_name": vm_name,
|
||||
"virtual_machine_config": vm_config
|
||||
}]
|
||||
})
|
||||
if result:
|
||||
st.success("VM workload created successfully!")
|
||||
st.session_state.volumes = [] # Reset volumes after creation
|
||||
st.rerun()
|
||||
|
||||
# List Container Workloads
|
||||
st.subheader("Existing Container Workloads")
|
||||
VirtualMachines = st.session_state.client.get_virtual_machines()
|
||||
|
||||
if VirtualMachines:
|
||||
# Filter VirtualMachines based on the selected VDC
|
||||
if selected_vdc_name:
|
||||
selected_vdc_id = vdc_options[selected_vdc_name]
|
||||
filtered_VirtualMachines = [VirtualMachine for VirtualMachine in VirtualMachines if VirtualMachine['vdc_id'] == selected_vdc_id]
|
||||
else:
|
||||
filtered_VirtualMachines = VirtualMachines
|
||||
|
||||
for VirtualMachine in filtered_VirtualMachines:
|
||||
with st.container():
|
||||
col1, col2, col3 = st.columns([3, 1, 1])
|
||||
with col1:
|
||||
# Safely handle empty or None launch_params
|
||||
launch_params = VirtualMachine.get('launch_params') # Use .get() to avoid KeyError if 'launch_params' doesn't exist
|
||||
if launch_params: # Check if launch_params is not empty or None
|
||||
try:
|
||||
VirtualMachine_name = json.loads(launch_params).get('docker_image') # Use .get() to avoid KeyError if 'docker_image' doesn't exist
|
||||
except json.JSONDecodeError:
|
||||
print("LP not vald json")
|
||||
# Handle the case where launch_params is not valid JSON
|
||||
VirtualMachine_name = None
|
||||
else:
|
||||
print("lp empty")
|
||||
# Handle the case where launch_params is empty or None
|
||||
VirtualMachine_name = None
|
||||
st.markdown(f"### {VirtualMachine['name']}")
|
||||
st.text(f"Status: {VirtualMachine['status']}")
|
||||
st.text(f"VDC ID: {VirtualMachine['vdc_id']}")
|
||||
st.text(f"Image: {VirtualMachine_name}")
|
||||
with col2:
|
||||
st.text(f"ID: {VirtualMachine['id']}")
|
||||
with col3:
|
||||
if st.button("View Details", key=f"view_{VirtualMachine['id']}"):
|
||||
print("Button is pressed")
|
||||
st.session_state.selected_resource = VirtualMachine
|
||||
st.session_state.view_type = 'detail'
|
||||
st.session_state.selected_resource_type = 'VirtualMachine'
|
||||
st.rerun()
|
||||
|
||||
else:
|
||||
st.info("No VirtualMachine workloads found.")
|
||||
|
||||
+6
-3
@@ -4,7 +4,7 @@ import json
|
||||
import asyncio
|
||||
from worker_tasks.container import ContainerTask
|
||||
from worker_tasks.file_presence import FilePresenceTask
|
||||
from worker_tasks.libvirt import LibvirtVMTask
|
||||
from worker_tasks.libvirt import LibvirtVirtualMachineTask
|
||||
from worker_tasks.ping import PingTask
|
||||
from worker_tasks.report import ReportTask
|
||||
|
||||
@@ -89,11 +89,14 @@ class WorkerClient:
|
||||
result = ContainerTask(job_details, logger).Create()
|
||||
elif task_type == "container-delete":
|
||||
result = ContainerTask(job_details, logger).Delete()
|
||||
elif task_type == "vm-create":
|
||||
result = LibvirtVMTask(job_details,self.libvirt_config, logger).execute()
|
||||
elif task_type == "virtual-machine-create":
|
||||
result = LibvirtVirtualMachineTask(job_details,self.libvirt_config, logger).execute()
|
||||
elif task_type == "virtual-machine-delete":
|
||||
result = LibvirtVirtualMachineTask(job_details,self.libvirt_config, logger).execute()
|
||||
else:
|
||||
raise ValueError(f"Unknown task type: {task_type}")
|
||||
|
||||
|
||||
await self.send_task_result(task_id, result, task_worker_id)
|
||||
except Exception as e:
|
||||
logger.error(f"Error processing task {task_id}: {e}")
|
||||
|
||||
@@ -101,7 +101,7 @@ class ContainerTask:
|
||||
container_name = container['container_id']
|
||||
|
||||
# Check if the container is marked for deletion
|
||||
if container.get('deleted', False):
|
||||
if container.get('desired_state', False):
|
||||
self.logger.info(f"Container '{container_name}' is marked for deletion. Ensuring it is not present...")
|
||||
self.delete_container(client, container_name)
|
||||
return # Exit the function as no further action is needed for deleted containers
|
||||
@@ -221,7 +221,7 @@ class ContainerTask:
|
||||
|
||||
self.logger.info(_container)
|
||||
# Is this container a NS controller or regular container?
|
||||
# container_workload_type=self.params['containers']
|
||||
# container_workload_type=self.params['container']
|
||||
|
||||
create_result=self.launch_container(client, _container)
|
||||
#TODO - Track the resultsof the tasks so we can rollback\delete if any subsequent part of the task fails
|
||||
@@ -239,7 +239,7 @@ class ContainerTask:
|
||||
|
||||
def Delete(self):
|
||||
"""Delete the specified container and its namespace controller if it's the last one."""
|
||||
container_name=self.params['containers'][0]['container_id']
|
||||
container_name=self.params['container'][0]['container_id']
|
||||
self.logger.info(f"Executing ContainerTask - Delete for container '{container_name}'...")
|
||||
try:
|
||||
# Initialize Docker client
|
||||
|
||||
+98
-97
@@ -7,15 +7,15 @@ from logging import Logger
|
||||
import xml.etree.ElementTree as ET
|
||||
from xml.dom import minidom
|
||||
|
||||
class LibvirtVMTask:
|
||||
class LibvirtVirtualMachineTask:
|
||||
def __init__(self, params, config, logger):
|
||||
"""
|
||||
Initialize the LibvirtVMTask.
|
||||
Initialize the LibvirtVirtualMachineTask.
|
||||
|
||||
Args:
|
||||
params (dict): Contains 'vm_id', 'desired_state', optional 'action' and 'vm_config'/'xml_config'.
|
||||
vm_config (dict): VM configuration details including CPU, RAM, networks, volumes.
|
||||
xml_config (str): Optional direct XML configuration (used if vm_config not provided).
|
||||
params (dict): Contains 'virtual_machine_id', 'desired_state', optional 'action' and 'virtual_machine_config'/'xml_config'.
|
||||
virtual_machine_config (dict): VirtualMachine configuration details including CPU, RAM, networks, volumes.
|
||||
xml_config (str): Optional direct XML configuration (used if virtual_machine_config not provided).
|
||||
config (dict): Configuration including default volume path.
|
||||
logger (Logger): Logger instance for logging.
|
||||
"""
|
||||
@@ -23,17 +23,17 @@ class LibvirtVMTask:
|
||||
raise TypeError("Logger must be an instance of the logging.Logger class.")
|
||||
self.logger = logger
|
||||
|
||||
self.logger.info("Initializing LibvirtVMTask")
|
||||
self.logger.info("Initializing LibvirtVirtualMachineTask")
|
||||
|
||||
# Validate params
|
||||
required_params = {
|
||||
"vm_id": str,
|
||||
"virtual_machine_id": str,
|
||||
"desired_state": str, # e.g., 'running', 'stopped', 'deleted'
|
||||
}
|
||||
optional_params = {
|
||||
"action": str, # 'reboot', 'start', 'stop'
|
||||
"xml_config": str, # Direct XML configuration
|
||||
"vm_config": dict, # JSON configuration for VM details
|
||||
"virtual_machine_config": dict, # JSON configuration for VirtualMachine details
|
||||
}
|
||||
|
||||
if not isinstance(params, dict):
|
||||
@@ -53,12 +53,12 @@ class LibvirtVMTask:
|
||||
self.logger.error(f"Parameter '{key}' must be of type {expected_type.__name__}.")
|
||||
raise TypeError(f"Parameter '{key}' must be of type {expected_type.__name__}.")
|
||||
|
||||
self.vm_id = params["vm_id"]
|
||||
self.virtual_machine_id = params["virtual_machine_id"]
|
||||
self.desired_state = params["desired_state"]
|
||||
self.action = params.get("action")
|
||||
|
||||
# Handle VM configuration
|
||||
self.vm_config = params.get("vm_config")
|
||||
# Handle VirtualMachine configuration
|
||||
self.virtual_machine_config = params.get("virtual_machine_config")
|
||||
self.xml_config = params.get("xml_config")
|
||||
self.volume_paths = [] # Will store volume paths for cleaning up
|
||||
|
||||
@@ -68,27 +68,27 @@ class LibvirtVMTask:
|
||||
self.logger.error("Config must be a dictionary containing 'default_volume_path'.")
|
||||
raise ValueError("Config must be a dictionary containing 'default_volume_path'.")
|
||||
|
||||
# If vm_config is provided, ensure volumes exist and generate XML from it
|
||||
if self.vm_config:
|
||||
self.logger.info(f"VM configuration provided as JSON. Processing volumes and generating XML.")
|
||||
# If virtual_machine_config is provided, ensure volumes exist and generate XML from it
|
||||
if self.virtual_machine_config:
|
||||
self.logger.info(f"VirtualMachine configuration provided as JSON. Processing volumes and generating XML.")
|
||||
# Handle storage volumes if needed
|
||||
if "volumes" in self.vm_config:
|
||||
if "volumes" in self.virtual_machine_config:
|
||||
self._process_storage_volumes()
|
||||
self.xml_config = self._generate_xml_from_config()
|
||||
elif not self.xml_config and self.desired_state != "deleted":
|
||||
self.logger.error("Either vm_config or xml_config must be provided for non-delete operations.")
|
||||
raise ValueError("Either vm_config or xml_config must be provided for non-delete operations.")
|
||||
self.logger.error("Either virtual_machine_config or xml_config must be provided for non-delete operations.")
|
||||
raise ValueError("Either virtual_machine_config or xml_config must be provided for non-delete operations.")
|
||||
|
||||
self.logger.info(f"LibvirtVMTask initialized with vm_id={self.vm_id}, desired_state={self.desired_state}")
|
||||
self.logger.info(f"LibvirtVirtualMachineTask initialized with virtual_machine_id={self.virtual_machine_id}, desired_state={self.desired_state}")
|
||||
|
||||
def _process_storage_volumes(self):
|
||||
"""
|
||||
Process storage volumes defined in vm_config.
|
||||
Process storage volumes defined in virtual_machine_config.
|
||||
Create any volumes that don't exist.
|
||||
Keep track of volume paths for potential cleanup during delete.
|
||||
"""
|
||||
self.logger.info("Processing storage volumes")
|
||||
for volume in self.vm_config["volumes"]:
|
||||
for volume in self.virtual_machine_config["volumes"]:
|
||||
if "id" not in volume:
|
||||
self.logger.error("Volume configuration missing required 'id' attribute.")
|
||||
raise ValueError("Volume configuration missing required 'id' attribute.")
|
||||
@@ -209,44 +209,45 @@ class LibvirtVMTask:
|
||||
|
||||
def _generate_xml_from_config(self):
|
||||
"""
|
||||
Generate libvirt XML configuration from vm_config dictionary.
|
||||
Generate libvirt XML configuration from virtual_machine_config dictionary.
|
||||
|
||||
Args:
|
||||
vm_id (str): A unique identifier for the VM to be embedded in the XML.
|
||||
virtual_machine_id (str): A unique identifier for the VirtualMachine to be embedded in the XML.
|
||||
|
||||
Returns:
|
||||
str: XML configuration for the VM.
|
||||
str: XML configuration for the VirtualMachine.
|
||||
"""
|
||||
self.logger.debug(f"Inside _generate_xml_from_config for {self.virtual_machine_id}")
|
||||
try:
|
||||
# Validate vm_config structure
|
||||
# Validate virtual_machine_config structure
|
||||
required_config = ["memory", "vcpu"]
|
||||
for key in required_config:
|
||||
if key not in self.vm_config:
|
||||
raise ValueError(f"Missing required VM configuration parameter: {key}")
|
||||
if key not in self.virtual_machine_config:
|
||||
raise ValueError(f"Missing required VirtualMachine configuration parameter: {key}")
|
||||
|
||||
# Create root domain element
|
||||
domain = ET.Element("domain")
|
||||
domain.set("type", "kvm") # Default to KVM, could be made configurable
|
||||
domain.set("type", "kvm") # Default to KVirtualMachine, could be made configurable
|
||||
domain.set("xmlns:custom", "http://example.com/xmlns/libvirt/custom")
|
||||
|
||||
# Basic VM information
|
||||
ET.SubElement(domain, "name").text = self.vm_id
|
||||
|
||||
# Basic VirtualMachine information
|
||||
ET.SubElement(domain, "name").text = self.virtual_machine_id
|
||||
self.logger.debug(f"Set name")
|
||||
# Memory configuration (in KiB)
|
||||
memory = int(self.vm_config["memory"]) * 1024 # Convert MB to KiB
|
||||
memory = int(self.virtual_machine_config["memory"]) * 1024 # Convert MB to KiB
|
||||
ET.SubElement(domain, "memory", unit="KiB").text = str(memory)
|
||||
ET.SubElement(domain, "currentMemory", unit="KiB").text = str(memory)
|
||||
|
||||
# CPU configuration
|
||||
ET.SubElement(domain, "vcpu", placement="static").text = str(self.vm_config["vcpu"])
|
||||
ET.SubElement(domain, "vcpu", placement="static").text = str(self.virtual_machine_config["vcpu"])
|
||||
|
||||
# OS configuration
|
||||
os = ET.SubElement(domain, "os")
|
||||
ET.SubElement(os, "type", arch="x86_64", machine="pc-q35-6.0").text = "hvm"
|
||||
|
||||
# Boot options if specified
|
||||
if "boot_devices" in self.vm_config:
|
||||
for device in self.vm_config["boot_devices"]:
|
||||
if "boot_devices" in self.virtual_machine_config:
|
||||
for device in self.virtual_machine_config["boot_devices"]:
|
||||
ET.SubElement(os, "boot", dev=device)
|
||||
|
||||
# Features
|
||||
@@ -262,8 +263,8 @@ class LibvirtVMTask:
|
||||
ET.SubElement(devices, "emulator").text = "/usr/bin/qemu-system-x86_64"
|
||||
|
||||
# Add disks/volumes
|
||||
if "volumes" in self.vm_config:
|
||||
for idx, volume in enumerate(self.vm_config["volumes"]):
|
||||
if "volumes" in self.virtual_machine_config:
|
||||
for idx, volume in enumerate(self.virtual_machine_config["volumes"]):
|
||||
disk = ET.SubElement(devices, "disk", type="file", device="disk")
|
||||
ET.SubElement(disk, "driver", name="qemu", type="qcow2")
|
||||
|
||||
@@ -278,10 +279,10 @@ class LibvirtVMTask:
|
||||
ET.SubElement(disk, "target", dev=device_name, bus="virtio")
|
||||
|
||||
# Add network interfaces
|
||||
if "networks" in self.vm_config:
|
||||
for network in self.vm_config["networks"]:
|
||||
if "networks" in self.virtual_machine_config:
|
||||
for network in self.virtual_machine_config["networks"]:
|
||||
interface = ET.SubElement(devices, "interface", type="network")
|
||||
ET.SubElement(interface, "source", network=network["name"])
|
||||
ET.SubElement(interface, "source", network=network["id"])
|
||||
|
||||
if "mac_address" in network:
|
||||
ET.SubElement(interface, "mac", address=network["mac_address"])
|
||||
@@ -289,25 +290,25 @@ class LibvirtVMTask:
|
||||
ET.SubElement(interface, "model", type="virtio")
|
||||
|
||||
# Add graphics if specified
|
||||
if "graphics" in self.vm_config:
|
||||
graphics_config = self.vm_config["graphics"]
|
||||
if "graphics" in self.virtual_machine_config:
|
||||
graphics_config = self.virtual_machine_config["graphics"]
|
||||
graphics = ET.SubElement(devices, "graphics", type=graphics_config.get("type", "vnc"))
|
||||
|
||||
for attr, value in graphics_config.items():
|
||||
if attr != "type":
|
||||
graphics.set(attr, str(value))
|
||||
|
||||
# Add metadata section to store VM ID and VM name
|
||||
# Add metadata section to store VirtualMachine ID and VirtualMachine name
|
||||
metadata = ET.SubElement(domain, "metadata")
|
||||
custom_metadata = ET.SubElement(metadata, "custom:metadata", xmlns_custom="http://example.com/xmlns/libvirt/custom")
|
||||
ET.SubElement(custom_metadata, "custom:vm_id").text = self.vm_id
|
||||
ET.SubElement(custom_metadata, "custom:vm_name").text = self.vm_id
|
||||
ET.SubElement(custom_metadata, "custom:virtual_machine_id").text = self.virtual_machine_id
|
||||
ET.SubElement(custom_metadata, "custom:VirtualMachine_name").text = self.virtual_machine_id
|
||||
|
||||
# Convert to pretty XML string
|
||||
rough_string = ET.tostring(domain, 'utf-8')
|
||||
reparsed = minidom.parseString(rough_string)
|
||||
xml_config = reparsed.toprettyxml(indent=" ")
|
||||
self.logger.info(f"Generated XML configuration for VM '{self.vm_id}' with ID '{self.vm_id}'.")
|
||||
self.logger.info(f"Generated XML configuration for VirtualMachine '{self.virtual_machine_id}' with ID '{self.virtual_machine_id}'.")
|
||||
return xml_config
|
||||
|
||||
except Exception as e:
|
||||
@@ -317,7 +318,7 @@ class LibvirtVMTask:
|
||||
|
||||
def execute(self):
|
||||
"""
|
||||
Execute the task: Manage VM state and optionally perform actions.
|
||||
Execute the task: Manage VirtualMachine state and optionally perform actions.
|
||||
|
||||
Returns:
|
||||
dict: Response payload indicating success or failure.
|
||||
@@ -328,33 +329,33 @@ class LibvirtVMTask:
|
||||
if conn is None:
|
||||
raise RuntimeError("Failed to open connection to libvirt.")
|
||||
|
||||
# Extract volume paths from VM if it already exists
|
||||
# This is needed for deletion when the task is called without vm_config
|
||||
# Extract volume paths from VirtualMachine if it already exists
|
||||
# This is needed for deletion when the task is called without virtual_machine_config
|
||||
if not self.volume_paths:
|
||||
try:
|
||||
domain = conn.lookupByName(self.vm_id)
|
||||
domain = conn.lookupByName(self.virtual_machine_id)
|
||||
self._extract_volume_paths_from_domain(domain)
|
||||
except libvirt.libvirtError:
|
||||
# VM doesn't exist, so no volumes to extract
|
||||
# VirtualMachine doesn't exist, so no volumes to extract
|
||||
pass
|
||||
|
||||
# Lookup VM by name
|
||||
self.logger.info("Looking up VM Name")
|
||||
# Lookup VirtualMachine by name
|
||||
self.logger.info("Looking up VirtualMachine Name")
|
||||
try:
|
||||
# Get all domains (VMs) managed by the libvirt connection
|
||||
# Get all domains (VirtualMachines) managed by the libvirt connection
|
||||
all_domains = conn.listAllDomains(0) # 0 means no flags, return all domains
|
||||
self.logger.info(f"Retrieved {len(all_domains)} VMs from libvirt.")
|
||||
self.logger.info(f"Retrieved {len(all_domains)} VirtualMachines from libvirt.")
|
||||
|
||||
# Search for the VM with the given name
|
||||
# Search for the VirtualMachine with the given name
|
||||
domain = None
|
||||
for dom in all_domains:
|
||||
if dom.name() == self.vm_id:
|
||||
if dom.name() == self.virtual_machine_id:
|
||||
domain = dom
|
||||
self.logger.info(f"VM '{self.vm_id}' exists. Managing state.")
|
||||
self.logger.info(f"VirtualMachine '{self.virtual_machine_id}' exists. Managing state.")
|
||||
break
|
||||
|
||||
if domain is None:
|
||||
self.logger.info(f"VM '{self.vm_id}' does not exist.")
|
||||
self.logger.info(f"VirtualMachine '{self.virtual_machine_id}' does not exist.")
|
||||
except libvirt.libvirtError as e:
|
||||
self.logger.error(f"Error retrieving domains from libvirt: {e}")
|
||||
raise e
|
||||
@@ -362,104 +363,104 @@ class LibvirtVMTask:
|
||||
# Handle desired state
|
||||
if self.desired_state == "deleted":
|
||||
if domain:
|
||||
self.logger.info(f"Deleting VM '{self.vm_id}'.")
|
||||
self.logger.info(f"Deleting VirtualMachine '{self.virtual_machine_id}'.")
|
||||
if domain.isActive():
|
||||
domain.destroy() # Stop the VM if running
|
||||
domain.destroy() # Stop the VirtualMachine if running
|
||||
domain.undefine() # Remove its definition
|
||||
|
||||
# Delete associated storage volumes
|
||||
self._delete_storage_volumes()
|
||||
|
||||
response_payload["response"]["status"] = "VM deleted"
|
||||
response_payload["response"]["status"] = "VirtualMachine deleted"
|
||||
else:
|
||||
self.logger.info(f"VM '{self.vm_id}' is already deleted.")
|
||||
response_payload["response"]["status"] = "VM already deleted"
|
||||
self.logger.info(f"VirtualMachine '{self.virtual_machine_id}' is already deleted.")
|
||||
response_payload["response"]["status"] = "VirtualMachine already deleted"
|
||||
response_payload["success"] = True
|
||||
|
||||
elif self.desired_state in ["running", "stopped"]:
|
||||
if domain:
|
||||
# Check if VM config has changed
|
||||
# Check if VirtualMachine config has changed
|
||||
if self.xml_config:
|
||||
current_xml = domain.XMLDesc()
|
||||
if self._is_xml_config_different(current_xml, self.xml_config):
|
||||
self.logger.info(f"VM configuration has changed. Updating VM '{self.vm_id}'.")
|
||||
# Need to recreate the VM with new configuration
|
||||
self.logger.info(f"VirtualMachine configuration has changed. Updating VirtualMachine '{self.virtual_machine_id}'.")
|
||||
# Need to recreate the VirtualMachine with new configuration
|
||||
if domain.isActive():
|
||||
domain.destroy() # Stop the VM if running
|
||||
domain.destroy() # Stop the VirtualMachine if running
|
||||
domain.undefine() # Remove its definition
|
||||
domain = conn.defineXML(self.xml_config) # Create with new config
|
||||
|
||||
self.logger.info(f"VM '{self.vm_id}' exists. Checking state.")
|
||||
self.logger.info(f"VirtualMachine '{self.virtual_machine_id}' exists. Checking state.")
|
||||
if self.desired_state == "running":
|
||||
self.logger.info(f"VM desired state is running, current state is {domain.isActive()}")
|
||||
self.logger.info(f"VirtualMachine desired state is running, current state is {domain.isActive()}")
|
||||
if not domain.isActive():
|
||||
self.logger.info(f"Starting VM '{self.vm_id}'.")
|
||||
domain.create() # Start the VM
|
||||
response_payload["response"]["status"] = "VM running"
|
||||
self.logger.info(f"Starting VirtualMachine '{self.virtual_machine_id}'.")
|
||||
domain.create() # Start the VirtualMachine
|
||||
response_payload["response"]["status"] = "VirtualMachine running"
|
||||
else: # desired_state == "stopped"
|
||||
self.logger.info(f"VM desired state is stopped, current state is {domain.isActive()}")
|
||||
self.logger.info(f"VirtualMachine desired state is stopped, current state is {domain.isActive()}")
|
||||
if domain.isActive():
|
||||
self.logger.info(f"Stopping VM '{self.vm_id}'.")
|
||||
domain.destroy() # Stop the VM
|
||||
response_payload["response"]["status"] = "VM stopped"
|
||||
self.logger.info(f"Stopping VirtualMachine '{self.virtual_machine_id}'.")
|
||||
domain.destroy() # Stop the VirtualMachine
|
||||
response_payload["response"]["status"] = "VirtualMachine stopped"
|
||||
else:
|
||||
if not self.xml_config:
|
||||
raise ValueError(f"XML configuration is required to create VM '{self.vm_id}'.")
|
||||
self.logger.info(f"Creating VM '{self.vm_id}' with the provided configuration.")
|
||||
raise ValueError(f"XML configuration is required to create VirtualMachine '{self.virtual_machine_id}'.")
|
||||
self.logger.info(f"Creating VirtualMachine '{self.virtual_machine_id}' with the provided configuration.")
|
||||
domain = conn.defineXML(self.xml_config)
|
||||
if self.desired_state == "running":
|
||||
domain.create() # Start the VM
|
||||
response_payload["response"]["status"] = "VM created and running"
|
||||
domain.create() # Start the VirtualMachine
|
||||
response_payload["response"]["status"] = "VirtualMachine created and running"
|
||||
else:
|
||||
response_payload["response"]["status"] = "VM created and stopped"
|
||||
response_payload["response"]["status"] = "VirtualMachine created and stopped"
|
||||
response_payload["success"] = True
|
||||
|
||||
# Handle optional actions
|
||||
if self.action:
|
||||
if not domain:
|
||||
raise RuntimeError(f"Cannot perform action '{self.action}' on a non-existent VM.")
|
||||
raise RuntimeError(f"Cannot perform action '{self.action}' on a non-existent VirtualMachine.")
|
||||
|
||||
if self.action == "start":
|
||||
if not domain.isActive():
|
||||
self.logger.info(f"Starting VM '{self.vm_id}'.")
|
||||
self.logger.info(f"Starting VirtualMachine '{self.virtual_machine_id}'.")
|
||||
domain.create()
|
||||
response_payload["response"]["action_status"] = "VM started"
|
||||
response_payload["response"]["action_status"] = "VirtualMachine started"
|
||||
elif self.action == "stop":
|
||||
if domain.isActive():
|
||||
self.logger.info(f"Stopping VM '{self.vm_id}'.")
|
||||
self.logger.info(f"Stopping VirtualMachine '{self.virtual_machine_id}'.")
|
||||
domain.destroy()
|
||||
response_payload["response"]["action_status"] = "VM stopped"
|
||||
response_payload["response"]["action_status"] = "VirtualMachine stopped"
|
||||
elif self.action == "reboot":
|
||||
if domain.isActive():
|
||||
self.logger.info(f"Rebooting VM '{self.vm_id}'.")
|
||||
self.logger.info(f"Rebooting VirtualMachine '{self.virtual_machine_id}'.")
|
||||
domain.reboot()
|
||||
response_payload["response"]["action_status"] = "VM rebooted"
|
||||
response_payload["response"]["action_status"] = "VirtualMachine rebooted"
|
||||
else:
|
||||
raise RuntimeError("Cannot reboot a stopped VM.")
|
||||
raise RuntimeError("Cannot reboot a stopped VirtualMachine.")
|
||||
else:
|
||||
raise ValueError(f"Unsupported action: {self.action}")
|
||||
self.logger.info(f"Action '{self.action}' completed successfully for VM '{self.vm_id}'.")
|
||||
self.logger.info(f"Action '{self.action}' completed successfully for VirtualMachine '{self.virtual_machine_id}'.")
|
||||
|
||||
except Exception as e:
|
||||
self.logger.error(f"Error executing LibvirtVMTask: {e}")
|
||||
self.logger.error(f"Error executing LibvirtVirtualMachineTask: {e}")
|
||||
response_payload["response"]["error"] = str(e)
|
||||
|
||||
finally:
|
||||
if 'conn' in locals() and conn:
|
||||
conn.close()
|
||||
|
||||
self.logger.info(f"LibvirtVMTask response: {response_payload}")
|
||||
self.logger.info(f"LibvirtVirtualMachineTask response: {response_payload}")
|
||||
return response_payload
|
||||
|
||||
def _extract_volume_paths_from_domain(self, domain):
|
||||
"""
|
||||
Extract volume paths from an existing domain.
|
||||
This is useful when deleting a VM without having its original configuration.
|
||||
This is useful when deleting a VirtualMachine without having its original configuration.
|
||||
|
||||
Args:
|
||||
domain: libvirt domain object
|
||||
"""
|
||||
self.logger.info(f"Extracting volume paths from existing domain '{self.vm_id}'")
|
||||
self.logger.info(f"Extracting volume paths from existing domain '{self.virtual_machine_id}'")
|
||||
try:
|
||||
xml_desc = domain.XMLDesc()
|
||||
root = ET.fromstring(xml_desc)
|
||||
@@ -479,16 +480,16 @@ class LibvirtVMTask:
|
||||
"""
|
||||
Delete all associated storage volumes.
|
||||
"""
|
||||
self.logger.info(f"Deleting storage volumes for VM '{self.vm_id}'")
|
||||
self.logger.info(f"Deleting storage volumes for VirtualMachine '{self.virtual_machine_id}'")
|
||||
for path in self.volume_paths:
|
||||
self._delete_volume(path)
|
||||
|
||||
def _is_xml_config_different(self, current_xml, new_xml):
|
||||
"""
|
||||
Compare current VM XML with the new XML configuration to determine if update is needed.
|
||||
Compare current VirtualMachine XML with the new XML configuration to determine if update is needed.
|
||||
|
||||
Args:
|
||||
current_xml (str): Current XML configuration of the VM.
|
||||
current_xml (str): Current XML configuration of the VirtualMachine.
|
||||
new_xml (str): New XML configuration to be applied.
|
||||
|
||||
Returns:
|
||||
|
||||
Reference in New Issue
Block a user