"""VM provisioning mixin for Proxmox — creates, destroys, and monitors VMs via Cloud-Init.""" from __future__ import annotations import base64 import hashlib import logging import time import yaml from typing import Any, Dict, List from urllib.parse import quote from napalm_device_types.models import ( NetworkTargetDict, StorageTargetDict, VMProvisionResultDict, VMStatusDict, ) _logger = logging.getLogger(__name__) # Downloaded cloud images are cached here on the hypervisor node, keyed by # filename, so provisioning multiple VMs from the same image only pays the # download cost once. _IMAGE_CACHE_DIR = "/var/lib/vz/template/netork-images" class ProxmoxVMProvisionMixin: """Mixin to add VM provisioning to ProxmoxDriver.""" def _run_node_command(self, command: str, timeout: int) -> str: """ Execute a shell command on the Proxmox node via SSH, raising on failure. Unlike ``_exec_ssh_command`` (best-effort, fixed timeout, swallows errors), this is for critical provisioning steps — image download, disk import — where a non-zero exit or a caller-specific timeout must surface as a hard failure rather than an empty string. """ import paramiko if self._ssh_client is None: ssh_user = self._ssh_username or self.username ssh_pass = self._ssh_password or self.password ssh_pkey = None if self._ssh_key and not ssh_pass: from io import StringIO as _StringIO ssh_pkey = paramiko.RSAKey.from_private_key(_StringIO(self._ssh_key)) self._ssh_client = paramiko.SSHClient() self._ssh_client.set_missing_host_key_policy(paramiko.AutoAddPolicy()) connect_kwargs: Dict[str, Any] = { "hostname": self.hostname, "port": 22, "username": ssh_user, "timeout": self.timeout, } if ssh_pkey: connect_kwargs["pkey"] = ssh_pkey else: connect_kwargs["password"] = ssh_pass self._ssh_client.connect(**connect_kwargs) _, stdout, stderr = self._ssh_client.exec_command(command, timeout=timeout) exit_status = stdout.channel.recv_exit_status() out = stdout.read().decode().strip() err = stderr.read().decode().strip() if exit_status != 0: raise RuntimeError(f"Command failed (exit {exit_status}): {command}\n{err or out}") return out def _download_cloud_image( self, image_url: str, image_checksum: str | None, timeout: int ) -> str: """ Download image_url to the node's image cache dir if not already present. Returns the local path on the hypervisor node. The cache filename is prefixed with a hash of the *full* URL, not just its basename — Ubuntu (and others) publish per-build URLs that change daily under a stable basename (e.g. .../release-20260713/ubuntu-26.04-server- cloudimg-amd64.img), so keying the cache on the basename alone let a stale previous-day build satisfy the "already cached" check and fail checksum verification against today's expected hash. Verifies image_checksum (format ":", e.g. "sha256:abc123...") if given. On mismatch, removes the bad file and retries the download once (covers a corrupted/partial transfer or a stale same-keyed file) before raising. """ filename = image_url.rstrip("/").rsplit("/", 1)[-1] url_hash = hashlib.sha256(image_url.encode()).hexdigest()[:12] local_path = f"{_IMAGE_CACHE_DIR}/{url_hash}-{filename}" algo, _, expected = (image_checksum or "").partition(":") algo = (algo or "sha256").lower() max_attempts = 2 for attempt in range(1, max_attempts + 1): exists = self._run_node_command( f"mkdir -p {_IMAGE_CACHE_DIR} && test -f {local_path} " f"&& echo EXISTS || echo MISSING", timeout=30, ) if "EXISTS" not in exists: _logger.info(f"Downloading cloud image {image_url} -> {local_path}") self._run_node_command( f"wget -q -O {local_path}.tmp '{image_url}' " f"&& mv {local_path}.tmp {local_path}", timeout=timeout, ) if not image_checksum: return local_path actual = self._run_node_command( f"{algo}sum {local_path} | awk '{{print $1}}'", timeout=60 ) if actual.lower() == expected.lower(): return local_path # Remove the bad file so the next attempt re-downloads instead of # reusing it. self._run_node_command(f"rm -f {local_path}", timeout=30) if attempt == max_attempts: raise RuntimeError( f"Checksum mismatch for {image_url}: expected {expected}, got {actual}" ) _logger.warning( f"Checksum mismatch for {image_url} on attempt {attempt}/{max_attempts} " "— retrying download" ) raise AssertionError("unreachable") # loop always returns or raises above def _find_default_image_storage(self) -> str: """Find a storage suitable for VM root disks (content includes 'images'). Proxmox's /storage API omits the "enabled" field entirely for storages that were never explicitly toggled — it is not present-and-falsy, it is just absent, defaulting to enabled. Only an explicit 0 means disabled. Queries the node-scoped /nodes/{node}/storage endpoint, not the cluster-wide /storage one: a storage can be configured with a "nodes" restriction limiting it to other cluster members, and the cluster-wide list doesn't reflect that — it would happily return a storage this node can't actually see, and "qm importdisk" would fail with "storage 'X' is not available on node 'Y'" after the VM shell was already created. """ for storage in self._node_api().storage.get(): content = storage.get("content", "") if "images" in content and storage.get("enabled", 1) != 0: return storage["storage"] raise ValueError( "No storage with content='images' found. Configure a storage for VM disks." ) def _get_storage_path(self, storage: str) -> str: """Resolve a storage's filesystem path on the node. Needed to write Cloud-Init snippets directly: Proxmox's /storage/{s}/upload API only accepts content in {iso, vztmpl, import} — "snippets" is rejected outright, so snippets must be written straight to the filesystem instead. Only dir-backed storages (dir, nfs, cifs, cephfs) expose "path"; those are also the only storage types Proxmox itself allows content='snippets' on. """ config = self._api.storage(storage).get() path = config.get("path") if not path: raise ValueError( f"Storage '{storage}' has no filesystem path (content='snippets' " "requires a dir/nfs/cifs/cephfs-backed storage)" ) return path def get_image_storages(self) -> List[StorageTargetDict]: """List node-available storage pools suitable for a new VM's root disk.""" targets: List[StorageTargetDict] = [] for storage in self._node_api().storage.get(): content = storage.get("content", "") if "images" not in content or storage.get("enabled", 1) == 0: continue if storage.get("active", 1) == 0: continue total = storage.get("total") or 0 avail = storage.get("avail") or 0 targets.append( { "name": storage["storage"], "type": storage.get("type", ""), "total_gb": round(total / (1024**3), 1), "available_gb": round(avail / (1024**3), 1), } ) return targets def _wait_for_task(self, upid: str, timeout: int = 120) -> None: """ Poll a Proxmox task until completion. Polls /nodes/{node}/tasks/{upid}/status until status == 'stopped'. Raises RuntimeError if exitstatus != 'OK' or timeout exceeded. """ start_time = time.time() while True: elapsed = time.time() - start_time if elapsed > timeout: raise RuntimeError(f"Task {upid} timed out after {timeout}s") try: task_status = self._node_api().tasks(upid).status.get() except Exception as e: _logger.debug(f"Error polling task {upid}: {e}") time.sleep(2) continue if task_status.get("status") == "stopped": exitstatus = task_status.get("exitstatus", "UNKNOWN") if exitstatus != "OK": raise RuntimeError( f"Task {upid} failed with exitstatus='{exitstatus}': " f"{task_status.get('exitstatus_text', 'no error message')}" ) return time.sleep(2) def create_vm_from_cloud_init( self, name: str, *, image_url: str, cpu: int, memory: int, nics: List[Dict[str, Any]], cloud_init_config: Dict[str, Any], image_checksum: str | None = None, ssh_public_keys: List[str] | None = None, disk_resize_gb: int | None = None, storage: str | None = None, download_timeout: int = 300, timeout: int = 180, ) -> VMProvisionResultDict: """ Create a new VM from a downloaded cloud image via Proxmox API. Steps: 1. Get next available VMID from cluster 2. Create an empty VM shell (no clone — no pre-existing template needed) 3. Download the cloud image on the node (cached by filename) and import it as the VM's root disk 4. Configure CPU, memory, and network interfaces 5. Verify snippet storage exists 6. Render cloud-init config to YAML and upload 7. Set Cloud-Init config references and SSH keys 8. Optionally resize root disk 9. Start the VM 10. Return VMID, name, node Args: name: new VM display name image_url: URL of the cloud image to download and use as root disk cpu: number of vCPUs memory: RAM in MB nics: list of NIC config dicts (bridge, vlan_tag/trunk_vlan_tags, dhcp flag) cloud_init_config: user-data dict (will be YAML-rendered) image_checksum: expected ":" checksum of the image, verified after download (None = no verification) ssh_public_keys: SSH public keys to inject disk_resize_gb: resize root disk to this size (None = no resize) storage: storage pool for the root disk (None = auto-detect first enabled, node-available storage with content='images') download_timeout: max seconds for the image download (skipped if cached) timeout: max seconds for the remaining provisioning steps Returns: {"vmid": str, "name": str, "node": str} Raises: RuntimeError: provisioning failure (download, import, config, timeout, etc.) ValueError: invalid storage or configuration """ try: _logger.info(f"Creating VM '{name}' from image {image_url}") # Step 1: Get next VMID next_vmid = self._api.cluster.nextid.get() vmid = int(next_vmid) _logger.info(f"Allocated VMID {vmid}") # Step 2: Create empty VM shell (no disks yet) _logger.info(f"Creating VM shell {vmid}") self._node_api().qemu.post( vmid=vmid, name=name, memory=memory, cores=cpu, ostype="l26", scsihw="virtio-scsi-pci", # Without this, Proxmox never attaches the virtio-serial # channel the QEMU guest agent needs — get_vm_status's # agent queries (below) would have nothing to talk to. agent="1", ) # Step 3: Download cloud image (cached) and import as root disk local_path = self._download_cloud_image( image_url, image_checksum, timeout=download_timeout ) image_storage = storage or self._find_default_image_storage() _logger.info(f"Importing {local_path} into VM {vmid} on storage {image_storage}") self._run_node_command( f"qm importdisk {vmid} {local_path} {image_storage} --format qcow2", timeout=timeout, ) # Proxmox leaves the imported disk as an "unusedN" reference — find # it and attach it as the boot disk. imported_config = self._node_api().qemu(vmid).config.get() unused_value = next( (v for k, v in imported_config.items() if k.startswith("unused")), None ) if not unused_value: raise RuntimeError( f"Disk import for VM {vmid} did not produce an unused disk reference" ) self._node_api().qemu(vmid).config.post( scsi0=f"{unused_value},discard=on", boot="order=scsi0", ) # Step 4: Configure network interfaces (CPU/memory already set at shell creation) _logger.info(f"Configuring {len(nics)} NIC(s) for VM {vmid}") config_args: Dict[str, Any] = {} # Build NIC config strings generically for i, nic in enumerate(nics): bridge = nic.get("bridge") if not bridge: raise ValueError(f"NIC {i}: bridge is required") # Build base config: model[=mac] + bridge mac = nic.get("mac") net_config = f"virtio={mac},bridge={bridge}" if mac else f"virtio,bridge={bridge}" # Add VLAN configuration (access vs trunk) if "trunk_vlan_tags" in nic and nic["trunk_vlan_tags"]: vlan_list = ";".join(str(v) for v in nic["trunk_vlan_tags"]) net_config += f",trunks={vlan_list}" elif "vlan_tag" in nic and nic["vlan_tag"] is not None: net_config += f",tag={nic['vlan_tag']}" config_args[f"net{i}"] = net_config self._node_api().qemu(vmid).config.post(**config_args) # Step 5: Verify snippet storage exists # (enabled is absent-not-falsy, and node-scoping matters — see # _find_default_image_storage) _logger.info("Checking for snippet storage...") storages = self._node_api().storage.get() snippet_storage = None for storage in storages: content = storage.get("content", "") if "snippets" in content and storage.get("enabled", 1) != 0: snippet_storage = storage["storage"] break if not snippet_storage: raise ValueError( "No storage with content='snippets' found. " "Configure a snippet-capable storage (e.g. local, nfs dir) " "and enable it." ) _logger.info(f"Using snippet storage: {snippet_storage}") # Step 6: Render and upload Cloud-Init config _logger.info(f"Rendering Cloud-Init config for VMID {vmid}") user_data_yaml = "#cloud-config\n" + yaml.dump( cloud_init_config, default_flow_style=False ) filename = f"{vmid}-user-data.yaml" _logger.debug(f"Writing Cloud-Init snippet {filename} to {snippet_storage}") # Proxmox's /storage/{s}/upload API only accepts content in # {iso, vztmpl, import} — "snippets" is rejected outright # ("does not have a value in the enumeration"). Snippets can only # be written directly to the filesystem, so resolve the storage's # backing path and write the file over SSH instead. storage_path = self._get_storage_path(snippet_storage) encoded = base64.b64encode(user_data_yaml.encode("utf-8")).decode("ascii") self._run_node_command( f"mkdir -p {storage_path}/snippets && " f"echo {encoded} | base64 -d > {storage_path}/snippets/{filename}", timeout=30, ) # Step 7: Configure Cloud-Init references and SSH keys _logger.info(f"Setting Cloud-Init config for VM {vmid}") cloud_init_args = { # The cloud-init drive is a disk image — it needs a storage # with content='images' (same requirement as the root disk), # NOT the snippet storage (content='snippets'). These are # often different storages; Proxmox fails at VM start with # "storage 'X' does not support content-type 'images'" if # this points at a snippets-only storage. "ide2": f"{image_storage}:cloudinit", "citype": "nocloud", "cicustom": f"user={snippet_storage}:snippets/{filename}", } # Configure DHCP for NICs where enabled (default True for index 0, False otherwise) for i, nic in enumerate(nics): dhcp_enabled = nic.get("dhcp", i == 0) # Default DHCP for first NIC only if dhcp_enabled: cloud_init_args[f"ipconfig{i}"] = "ip=dhcp" if ssh_public_keys: # URL-encode SSH keys for Proxmox API sshkeys = ";".join(ssh_public_keys) cloud_init_args["sshkeys"] = quote(sshkeys) self._node_api().qemu(vmid).config.post(**cloud_init_args) # Step 8: Optionally resize root disk if disk_resize_gb is not None: _logger.info(f"Resizing root disk to {disk_resize_gb}GB") # Find root disk (scsi0, virtio0, ide0, sata0 — whichever is first) try: config = self._node_api().qemu(vmid).config.get() root_disk = None for prefix in ("scsi", "virtio", "ide", "sata"): if f"{prefix}0" in config: root_disk = f"{prefix}0" break if root_disk: self._node_api().qemu(vmid).resize.put( disk=root_disk, size=f"{disk_resize_gb}G", ) else: _logger.warning(f"Could not find root disk for VM {vmid}, skipping resize") except Exception as e: _logger.warning(f"Failed to resize disk: {e}, continuing anyway") # Step 9: Start the VM _logger.info(f"Starting VM {vmid}") start_upid = self._node_api().qemu(vmid).status.start.post() self._wait_for_task(start_upid, timeout=timeout) _logger.info(f"VM {vmid} ('{name}') provisioned successfully on {self._node_name}") return { "vmid": str(vmid), "name": name, "node": self._node_name, } except Exception as e: _logger.exception(f"Failed to create VM '{name}': {e}") raise def destroy_vm( self, vmid: str, *, remove_disk: bool = True, timeout: int = 60, ) -> None: """ Destroy a virtual machine and optionally remove its storage. Steps: 1. Stop the VM if running 2. Delete VM configuration and optionally disks 3. Clean up Cloud-Init snippets Args: vmid: hypervisor VMID (string, e.g. "101") remove_disk: if True, also delete disks and storage timeout: max seconds for stop/delete operations Raises: RuntimeError: VM doesn't exist or destruction fails """ try: vmid_int = int(vmid) _logger.info(f"Destroying VM {vmid}") # Step 1: Stop the VM if running try: _logger.debug(f"Stopping VM {vmid}") stop_upid = self._node_api().qemu(vmid_int).status.stop.post() self._wait_for_task(stop_upid, timeout=timeout) except Exception as e: _logger.debug(f"VM {vmid} stop failed (may already be stopped): {e}") # Step 2: Delete VM # Proxmox's API parameter is hyphenated (destroy-unreferenced-disks), # not a valid Python identifier — proxmoxer forwards kwargs to the # request verbatim with no underscore-to-hyphen translation, so this # must be built as a dict and unpacked rather than passed as a kwarg. _logger.debug(f"Deleting VM {vmid} configuration and disks") delete_params = { "purge": 1, "destroy-unreferenced-disks": 1 if remove_disk else 0, } self._node_api().qemu(vmid_int).delete(**delete_params) # Step 3: Clean up Cloud-Init snippets # (This is best-effort; snippet files may be unreachable if storage is unavailable) try: config = self._node_api().qemu(vmid_int).config.get() cicustom = config.get("cicustom", "") if "snippets/" in cicustom: parts = cicustom.split("=") if len(parts) >= 2: snippet_ref = parts[1] # e.g. "snippets:snippets/101-user-data.yaml" storage, filepath = snippet_ref.split(":", 1) _logger.debug(f"Deleting snippet {filepath} from {storage}") try: self._node_api().storage(storage).content(filepath).delete() except Exception as e: _logger.warning(f"Failed to delete snippet {filepath}: {e}") except Exception as e: _logger.debug(f"Could not clean up snippets for VM {vmid}: {e}") _logger.info(f"VM {vmid} destroyed successfully") except Exception as e: _logger.exception(f"Failed to destroy VM {vmid}: {e}") raise def get_vm_status( self, vmid: str, *, wait_for_ip: bool = False, timeout: int = 300, poll_interval: int = 5, ) -> VMStatusDict: """ Get the runtime status of a virtual machine. Optionally waits for the guest-agent to report an IP address on the management NIC (net0), useful after provisioning. Args: vmid: hypervisor VMID (string) wait_for_ip: if True, poll until IP appears on net0 timeout: max seconds to wait for IP (if wait_for_ip=True) poll_interval: seconds between status polls Returns: {"status": str, "ip_address": str, "hostname": str, "mac_address": str} (ip_address, hostname, mac_address only if VM is running and has network info) Raises: RuntimeError: VM doesn't exist or wait_for_ip times out """ try: vmid_int = int(vmid) _logger.debug(f"Getting status for VM {vmid}") # VM must exist / be readable before we start polling. try: self._node_api().qemu(vmid_int).config.get() except Exception: # VM may not exist yet or config not readable return {"status": "unknown"} # Polling loop start_time = time.time() while True: elapsed = time.time() - start_time if wait_for_ip and elapsed > timeout: raise RuntimeError(f"VM {vmid} failed to acquire IP within {timeout}s") try: # Query guest-agent network interfaces. Proxmox's REST path is # "network-get-interfaces" (hyphens) — it must be passed as a # resource id via __call__, not dotted attribute access (which # would silently build a non-existent "network_get_interfaces" # path and 404 on every poll). agent_info = ( self._node_api().qemu(vmid_int).agent("network-get-interfaces").get() ) interfaces = (agent_info or {}).get("result", []) # The guest agent does not report interfaces in a fixed order — # "lo" commonly comes first. Skip it and take the first real # NIC that has an IPv4 address. for iface in interfaces: name = iface.get("name", "") if not name or name == "lo": continue for addr in iface.get("ip-addresses", []): if addr.get("ip-address-type") != "ipv4": continue ip_addr = addr.get("ip-address", "") if ip_addr: _logger.info(f"VM {vmid} acquired IP {ip_addr}") return { "status": "running", "ip_address": ip_addr, "hostname": name, "mac_address": iface.get("hardware-address", ""), } except Exception as e: _logger.debug(f"Error querying guest-agent for VM {vmid}: {e}") if not wait_for_ip: # Return immediate status without IP return {"status": "running"} # Wait before next poll time.sleep(poll_interval) except RuntimeError: raise except Exception as e: _logger.exception(f"Error getting status for VM {vmid}: {e}") raise RuntimeError(f"Failed to get VM {vmid} status: {e}") def get_network_targets(self) -> List[NetworkTargetDict]: """ List selectable network targets (bridges + SDN vnets) for a new VM's NIC. Excludes physical NICs, bonds, and other non-bridge interface types — those are never valid ``NICConfigDict.bridge`` values on Proxmox. """ targets: List[NetworkTargetDict] = [] for iface in self._get_node_network(): iface_type = iface.get("type") name = iface.get("iface", "") if not name: continue if iface_type == "bridge": vlan_aware = bool(int(iface.get("bridge_vlan_aware", 0) or 0)) targets.append({"name": name, "kind": "bridge", "vlan_aware": vlan_aware}) elif iface_type == "OVSBridge": # OVS bridges tag per-port regardless of a dedicated "VLAN aware" setting. targets.append({"name": name, "kind": "bridge", "vlan_aware": True}) for vnet in self._get_sdn_vnets(): name = vnet.get("vnet", "") if not name: continue # A vnet's VLAN is already fixed by its zone/tag — no separate vlan_tag applies. tag = vnet.get("tag") fixed_vlan_tag = int(tag) if tag is not None else None targets.append( { "name": name, "kind": "vnet", "vlan_aware": False, "fixed_vlan_tag": fixed_vlan_tag, } ) return targets