Feat: VM reconsilation
This commit is contained in:
@@ -76,11 +76,19 @@ def register_socketio_handlers(socketio):
|
||||
|
||||
# Check if this is a reconcile_and_delete task completion
|
||||
if task.task_type == "reconcile_and_delete" and result.get("success"):
|
||||
logger.info(f"[{worker_id}] Reconcile and delete task completed successfully, initiating host online transition")
|
||||
logger.info(f"[{worker_id}] Reconcile and delete task completed successfully, initiating VM reconciliation")
|
||||
# Import the function here to avoid circular imports
|
||||
from websocket_server.worker_manager import handle_reconcile_and_delete_completion
|
||||
handle_reconcile_and_delete_completion(worker_id)
|
||||
|
||||
# VM reconciliation is the last step before the host goes online. Handled even
|
||||
# on failure so a libvirt error cannot leave the host stuck in 'reconciling'.
|
||||
if task.task_type == "virtual-machine-reconcile":
|
||||
if not result.get("success"):
|
||||
logger.warning(f"[{worker_id}] VM reconciliation failed, bringing host online without a VM sweep")
|
||||
from websocket_server.worker_manager import handle_vm_reconcile_completion
|
||||
handle_vm_reconcile_completion(worker_id, result)
|
||||
|
||||
# Find workers that have tasks waiting on this task so we can wake them up.
|
||||
# Needed for cross-worker depends_on (e.g. VM create waiting on a DNS update
|
||||
# that ran on a different host).
|
||||
|
||||
@@ -215,10 +215,118 @@ def send_reconcile_and_delete_task(worker_id, expected_containers):
|
||||
|
||||
def handle_reconcile_and_delete_completion(worker_id):
|
||||
"""
|
||||
Handle the completion of the reconcile_and_delete task by moving host to online status.
|
||||
Handle the completion of the reconcile_and_delete task by reconciling VMs.
|
||||
|
||||
The host stays in 'reconciling' until the VM sweep finishes; it is moved online
|
||||
from handle_vm_reconcile_completion.
|
||||
"""
|
||||
logger.info(f"[{worker_id}] Reconcile and delete completed, moving host to online status")
|
||||
update_host_status(worker_id, "online")
|
||||
logger.info(f"[{worker_id}] Reconcile and delete completed, starting VM reconciliation")
|
||||
send_vm_reconcile_task(worker_id)
|
||||
|
||||
|
||||
def send_vm_reconcile_task(worker_id):
|
||||
"""
|
||||
Ask the worker to report the libvirt domains it currently has defined.
|
||||
"""
|
||||
try:
|
||||
task = Task(
|
||||
worker_id=worker_id,
|
||||
task_type="virtual-machine-reconcile",
|
||||
job_details=json.dumps({}),
|
||||
status="pending"
|
||||
)
|
||||
|
||||
session = Session()
|
||||
session.add(task)
|
||||
session.commit()
|
||||
session.close()
|
||||
|
||||
logger.debug(f"[{worker_id}] Created virtual-machine-reconcile task")
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"[{worker_id}] Error creating virtual-machine-reconcile task: {e}")
|
||||
# Never strand the host in 'reconciling' just because we could not ask for the report
|
||||
update_host_status(worker_id, "online")
|
||||
|
||||
|
||||
def send_vm_delete_task(worker_id, virtual_machine_id):
|
||||
"""
|
||||
Queue a virtual-machine-delete task for a single domain on this worker.
|
||||
"""
|
||||
try:
|
||||
task = Task(
|
||||
worker_id=worker_id,
|
||||
task_type="virtual-machine-delete",
|
||||
job_details=json.dumps({
|
||||
"virtual_machine_id": virtual_machine_id,
|
||||
"desired_state": "deleted",
|
||||
}),
|
||||
status="pending"
|
||||
)
|
||||
|
||||
session = Session()
|
||||
session.add(task)
|
||||
session.commit()
|
||||
session.close()
|
||||
|
||||
logger.debug(f"[{worker_id}] Created virtual-machine-delete task for {virtual_machine_id}")
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"[{worker_id}] Error creating virtual-machine-delete task for {virtual_machine_id}: {e}")
|
||||
|
||||
|
||||
def handle_vm_reconcile_completion(worker_id, result):
|
||||
"""
|
||||
Decide which of the worker's reported domains are orphans and queue deletes.
|
||||
|
||||
A domain is an orphan when the workload table says it should no longer be here —
|
||||
typically because failover moved the VM to another host while this one was
|
||||
unreachable, leaving the original still running and writing to shared storage.
|
||||
|
||||
Domains with no matching workload row are left strictly alone: this may be a BYO
|
||||
host whose own VMs predate enrolment, and destroying those would be unrecoverable.
|
||||
"""
|
||||
from websocket_server.config import get_db_session
|
||||
from app.models.models import Workload
|
||||
|
||||
domains = (result or {}).get("domains", [])
|
||||
orphans = []
|
||||
|
||||
try:
|
||||
with get_db_session() as session:
|
||||
for domain in domains:
|
||||
name = domain.get("name")
|
||||
if not name:
|
||||
continue
|
||||
|
||||
workload = session.query(Workload).filter(
|
||||
Workload.id == name,
|
||||
Workload.workload_type == "VirtualMachine",
|
||||
).first()
|
||||
|
||||
if not workload:
|
||||
logger.debug(f"[{worker_id}] Domain {name} is not a known workload, leaving it alone")
|
||||
continue
|
||||
|
||||
if workload.deleted or str(workload.workload_host_id) != str(worker_id):
|
||||
orphans.append(name)
|
||||
|
||||
for virtual_machine_id in orphans:
|
||||
send_vm_delete_task(worker_id, virtual_machine_id)
|
||||
|
||||
if orphans:
|
||||
logger.info(
|
||||
f"[{worker_id}] VM reconciliation queued deletes for {len(orphans)} "
|
||||
f"orphaned domain(s): {orphans}"
|
||||
)
|
||||
else:
|
||||
logger.info(f"[{worker_id}] VM reconciliation found no orphans ({len(domains)} domains reported)")
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"[{worker_id}] Error during VM reconciliation: {e}")
|
||||
|
||||
finally:
|
||||
update_host_status(worker_id, "online")
|
||||
|
||||
|
||||
def worker_dispatch_flag_check():
|
||||
|
||||
@@ -353,6 +353,9 @@ class WorkerClient:
|
||||
elif task_type == "virtual-machine-delete":
|
||||
from worker_tasks.libvirt import LibvirtVirtualMachineTask
|
||||
result = await loop.run_in_executor(None, functools.partial(LibvirtVirtualMachineTask(job_details, logger).execute))
|
||||
elif task_type == "virtual-machine-reconcile":
|
||||
from worker_tasks.libvirt import reconcile_virtual_machines
|
||||
result = await loop.run_in_executor(None, functools.partial(reconcile_virtual_machines, logger))
|
||||
elif task_type == "container-reconcile":
|
||||
result = await loop.run_in_executor(None, functools.partial(ContainerTask(logger,self.docker_monitor).reconcile_all_containers))
|
||||
elif task_type == "reconcile_and_delete":
|
||||
|
||||
@@ -2064,17 +2064,36 @@ class ContainerTask:
|
||||
3. Deletes any container that is NOT in the expected list
|
||||
|
||||
Args:
|
||||
job_details (dict): Contains expected_container_ids list
|
||||
|
||||
job_details (dict): Contains expected_container_ids list and the
|
||||
authoritative flag saying whether that list can be trusted
|
||||
|
||||
Returns:
|
||||
dict: Result of the reconciliation operation
|
||||
"""
|
||||
expected_container_ids = set(job_details.get("expected_container_ids", []))
|
||||
deleted_containers = []
|
||||
failed_deletions = []
|
||||
|
||||
|
||||
# An empty expected set means "delete everything", so only act on it when the
|
||||
# server confirms it actually read the set from the API. Without this, an API
|
||||
# blip during reconciliation would wipe every container on the host.
|
||||
if not job_details.get("authoritative", False):
|
||||
self.logger.warning(
|
||||
"Skipping reconcile_and_delete: expected container set is not authoritative"
|
||||
)
|
||||
return {
|
||||
"success": True,
|
||||
"skipped": True,
|
||||
"reason": "non_authoritative_expected_set",
|
||||
"deleted_containers": [],
|
||||
"failed_deletions": [],
|
||||
"expected_count": len(expected_container_ids),
|
||||
"deleted_count": 0,
|
||||
"failed_count": 0,
|
||||
}
|
||||
|
||||
self.logger.info(f"Starting reconcile_and_delete: expected {len(expected_container_ids)} containers")
|
||||
|
||||
|
||||
try:
|
||||
# Get all containers managed by this worker
|
||||
worker_id = settings.get_value("WORKER_ID")
|
||||
|
||||
@@ -531,3 +531,39 @@ class LibvirtVirtualMachineTask:
|
||||
on_error(str(exc))
|
||||
|
||||
self.logger.info(f"[vm-log-stream] Stream ended for VM '{vm_id}'")
|
||||
|
||||
|
||||
def reconcile_virtual_machines(logger):
|
||||
"""
|
||||
Report every libvirt domain defined on this host.
|
||||
|
||||
Deliberately report-only. Unlike containers, domains carry no marker saying we
|
||||
created them — an enrolled BYO host may already run VMs that predate it joining
|
||||
the region — so the worker cannot tell an orphan from a stranger. The server
|
||||
decides what to delete by checking the reported names against the workload
|
||||
table, and issues virtual-machine-delete for the ones it owns.
|
||||
|
||||
Returns:
|
||||
dict: {"success": bool, "domains": [{"name": str, "active": bool}, ...]}
|
||||
"""
|
||||
conn = None
|
||||
try:
|
||||
conn = libvirt.open("qemu:///system")
|
||||
if conn is None:
|
||||
raise SchedulingError(ErrorType.LIBVIRT_CONNECTION_FAILED, "Failed to open connection to libvirt")
|
||||
|
||||
domains = [
|
||||
{"name": dom.name(), "active": bool(dom.isActive())}
|
||||
for dom in conn.listAllDomains(0)
|
||||
]
|
||||
|
||||
logger.info(f"virtual-machine-reconcile: reporting {len(domains)} domains")
|
||||
return {"success": True, "domains": domains}
|
||||
|
||||
except libvirt.libvirtError as e:
|
||||
logger.error(f"Error listing domains during virtual-machine-reconcile: {e}")
|
||||
return {"success": False, "response": str(e), "domains": []}
|
||||
|
||||
finally:
|
||||
if conn:
|
||||
conn.close()
|
||||
|
||||
Reference in New Issue
Block a user