Replaces the template-clone flow with: create empty VM shell, download the cloud image on the node (cached by filename, optional checksum verification), qm importdisk, attach as scsi0. NIC config, snippet upload, ssh keys, disk resize, and start remain unchanged (already generic). New helpers: _run_node_command (strict SSH exec with custom timeout and non-zero-exit detection, unlike the best-effort _exec_ssh_command), _download_cloud_image (idempotent download + checksum check), _find_default_image_storage (content=images discovery, mirrors the existing snippet-storage discovery). 24 tests pass (10 new: _run_node_command x2, _download_cloud_image x4, plus rewrites of the 4 existing create_vm_from_cloud_init tests for the new flow).
561 lines
22 KiB
Python
561 lines
22 KiB
Python
"""VM provisioning mixin for Proxmox — creates, destroys, and monitors VMs via Cloud-Init."""
|
|
|
|
from __future__ import annotations
|
|
|
|
import logging
|
|
import time
|
|
import yaml
|
|
from typing import Any, Dict, List
|
|
from urllib.parse import quote
|
|
|
|
from napalm_device_types.models import NetworkTargetDict, 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. Verifies image_checksum
|
|
(format "<algo>:<hex>", e.g. "sha256:abc123...") if given, re-downloading
|
|
is left to the caller's next attempt if verification fails.
|
|
"""
|
|
filename = image_url.rstrip("/").rsplit("/", 1)[-1]
|
|
local_path = f"{_IMAGE_CACHE_DIR}/{filename}"
|
|
|
|
exists = self._run_node_command(
|
|
f"mkdir -p {_IMAGE_CACHE_DIR} && test -f {local_path} && 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}' && mv {local_path}.tmp {local_path}",
|
|
timeout=timeout,
|
|
)
|
|
|
|
if image_checksum:
|
|
algo, _, expected = image_checksum.partition(":")
|
|
algo = (algo or "sha256").lower()
|
|
actual = self._run_node_command(
|
|
f"{algo}sum {local_path} | awk '{{print $1}}'", timeout=60
|
|
)
|
|
if actual.lower() != expected.lower():
|
|
# Remove the bad file so a retry re-downloads instead of reusing it.
|
|
self._run_node_command(f"rm -f {local_path}", timeout=30)
|
|
raise RuntimeError(
|
|
f"Checksum mismatch for {image_url}: expected {expected}, got {actual}"
|
|
)
|
|
|
|
return local_path
|
|
|
|
def _find_default_image_storage(self) -> str:
|
|
"""Find a storage suitable for VM root disks (content includes 'images')."""
|
|
for storage in self._api.storage.get():
|
|
content = storage.get("content", "")
|
|
if "images" in content and storage.get("enabled"):
|
|
return storage["storage"]
|
|
raise ValueError(
|
|
"No storage with content='images' found. Configure a storage for VM disks."
|
|
)
|
|
|
|
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,
|
|
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 "<algo>:<hex>" 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)
|
|
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",
|
|
)
|
|
|
|
# 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 = 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 + bridge
|
|
net_config = 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
|
|
_logger.info("Checking for snippet storage...")
|
|
storages = self._api.storage.get()
|
|
snippet_storage = None
|
|
for storage in storages:
|
|
content = storage.get("content", "")
|
|
if "snippets" in content and storage.get("enabled"):
|
|
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"Uploading Cloud-Init snippet {filename} to {snippet_storage}")
|
|
|
|
# Upload to snippet storage
|
|
self._node_api().storage(snippet_storage).upload.post(
|
|
content="snippets",
|
|
filename=filename,
|
|
data=user_data_yaml,
|
|
)
|
|
|
|
# Step 7: Configure Cloud-Init references and SSH keys
|
|
_logger.info(f"Setting Cloud-Init config for VM {vmid}")
|
|
|
|
cloud_init_args = {
|
|
"ide2": f"{snippet_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
|
|
_logger.debug(f"Deleting VM {vmid} configuration and disks")
|
|
self._node_api().qemu(vmid_int).delete(
|
|
purge=1,
|
|
destroy_unreferenced_disks=1 if remove_disk else 0,
|
|
)
|
|
|
|
# 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}")
|
|
|
|
# Get VM config to infer net0 MAC (for matching guest-agent results)
|
|
try:
|
|
config = self._node_api().qemu(vmid_int).config.get()
|
|
except Exception:
|
|
# VM may not exist yet or config not readable
|
|
return {"status": "unknown"}
|
|
|
|
# Parse net0 MAC from config (if present)
|
|
net0_line = config.get("net0", "")
|
|
expected_mac = None
|
|
# Example: "virtio,bridge=vmbr0,tag=10" — no explicit MAC
|
|
# Proxmox auto-generates MACs in a deterministic pattern, but we'll
|
|
# match by looking for the first NIC's IP in guest-agent results
|
|
|
|
# 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
|
|
agent_info = self._node_api().qemu(vmid_int).agent.network_get_interfaces.get()
|
|
interfaces = agent_info.get("result", [])
|
|
|
|
# Find net0 (first interface with IP)
|
|
if interfaces:
|
|
net0_iface = interfaces[0] # Assumes net0 is first in list
|
|
net0_mac = net0_iface.get("hardware-address", "")
|
|
ip_addresses = net0_iface.get("ip-addresses", [])
|
|
|
|
if ip_addresses:
|
|
# Found IP
|
|
ip_info = ip_addresses[0]
|
|
ip_addr = ip_info.get("ip-address", "")
|
|
if ip_addr:
|
|
_logger.info(f"VM {vmid} acquired IP {ip_addr}")
|
|
return {
|
|
"status": "running",
|
|
"ip_address": ip_addr,
|
|
"hostname": net0_iface.get("name", ""),
|
|
"mac_address": net0_mac,
|
|
}
|
|
|
|
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
|