feat: VM provisioning, and reboot_host on ESXi

create_vm_from_cloud_init takes the same qcow2/raw cloud images netOrk
offers for Proxmox. ESXi can neither boot nor download them, so the
driver does both: download with checksum check and one retry, convert
with qemu-img to a streamOptimized VMDK (cached by URL), import through
a minimal OVF descriptor over NFC, pin requested MACs, grow the disk,
and attach a NoCloud seed ISO placed next to the VM's files. NoCloud
rather than guestinfo because it needs nothing in the guest; the
user-data installs open-vm-tools, which the driver declares as its
guest agent. A failure after the import removes the VM again.

Placement is a pure decision over inventory rows: a connected host
outside maintenance mode that sees the datastore and every port group,
with the resource pool and VM folder found by walking up to the
datacenter -- one path for a standalone host and for a vCenter.

Two faults vcsim surfaced and the tests now pin: a chunked upload body
next to a Content-Length is refused with 500, so the disk goes up as a
sized file object that also reports lease progress; and a device edit
replaces the device as sent, so disk and NIC edits start from the live
objects, backing included (vcsim panicked on a disk without one).

destroy_vm, get_vm_status, get_network_targets (port groups with their
fixed VLAN) and get_image_storages complete the contract.

reboot_host on ESXi uses RebootHost_Task and refuses outside
maintenance mode: force=True would cut power to running VMs.

Tested against vcsim in ESXi and vCenter mode, end to end.
This commit is contained in:
Christian Manivong
2026-09-24 10:00:33 +02:00
parent c8e472b4c6
commit 2b84ceae3a
25 changed files with 1851 additions and 14 deletions
+13 -7
View File
@@ -56,9 +56,13 @@ class Inventory:
finally:
view.Destroy()
def properties(self, obj: Any, paths: Sequence[str]) -> dict[str, Any]:
"""The given properties of one managed object; ``{}`` if it is gone."""
rows = self._retrieve(_PC.ObjectSpec(obj=obj), type(obj), paths)
def properties(self, obj: Any, paths: Sequence[str], *, raw: bool = False) -> dict[str, Any]:
"""The given properties of one managed object; ``{}`` if it is gone.
``raw=True`` keeps the pyVmomi objects -- for a caller that has to hand
a device back to vSphere whole, backing and all, in a reconfigure.
"""
rows = self._retrieve(_PC.ObjectSpec(obj=obj), type(obj), paths, raw=raw)
return rows[0] if rows else {}
def licenses(self) -> list[dict[str, Any]]:
@@ -68,7 +72,9 @@ class Inventory:
return []
return self.properties(manager, ["licenses"]).get("licenses", [])
def _retrieve(self, obj_spec: Any, vim_type: Any, paths: Sequence[str]) -> list[dict[str, Any]]:
def _retrieve(
self, obj_spec: Any, vim_type: Any, paths: Sequence[str], *, raw: bool = False
) -> list[dict[str, Any]]:
spec = _PC.FilterSpec(
objectSet=[obj_spec], propSet=[_PC.PropertySpec(type=vim_type, pathSet=list(paths))]
)
@@ -76,15 +82,15 @@ class Inventory:
result = pc.RetrievePropertiesEx([spec], _PC.RetrieveOptions(maxObjects=_PAGE_SIZE))
rows: list[dict[str, Any]] = []
while result is not None:
rows.extend(_row(obj) for obj in result.objects)
rows.extend(_row(obj, raw) for obj in result.objects)
if not result.token:
break
result = pc.ContinueRetrievePropertiesEx(token=result.token)
return rows
def _row(obj_content: Any) -> dict[str, Any]:
def _row(obj_content: Any, raw: bool = False) -> dict[str, Any]:
row: dict[str, Any] = {"_moref": obj_content.obj._moId}
for prop in obj_content.propSet or []:
row[prop.name] = to_plain(prop.val)
row[prop.name] = prop.val if raw else to_plain(prop.val)
return row
+4
View File
@@ -2,6 +2,7 @@
from __future__ import annotations
from pathlib import Path
from typing import Any, ClassVar
from napalm.base.exceptions import ConnectionException
@@ -22,6 +23,7 @@ from napalm_vmware.parse.storage import storage_pools
from napalm_vmware.parse.vm_config import vm_config
from napalm_vmware.parse.vms import vm_list
from napalm_vmware.parse.warnings import host_warnings
from napalm_vmware.provision.image import DEFAULT_CACHE_DIR
_DEFAULT_PORT = 443
#: ``about.apiType`` -> the driver that handles it, for a helpful refusal.
@@ -55,6 +57,8 @@ class VmwareBaseDriver(HypervisorDriver):
args = optional_args or {}
self._port = int(args.get("port") or _DEFAULT_PORT)
self._verify_ssl = bool(args.get("verify_ssl", args.get("ssl_verify", True)))
# Where converted cloud images are kept between provisioning jobs.
self._image_cache_dir = Path(args.get("image_cache_dir") or DEFAULT_CACHE_DIR)
self._si: Any = None
self._inventory: Any = None
+18 -1
View File
@@ -14,9 +14,10 @@ from napalm_vmware.base import VmwareBaseDriver
from napalm_vmware.parse.facts import esxi_facts
from napalm_vmware.parse.interfaces import host_interfaces, host_interfaces_ip
from napalm_vmware.parse.lldp import lldp_neighbors
from napalm_vmware.provisioning import VmwareProvisioningMixin
class VmwareEsxiDriver(VmwareActionsMixin, VmwareBaseDriver):
class VmwareEsxiDriver(VmwareProvisioningMixin, VmwareActionsMixin, VmwareBaseDriver):
"""One ESXi host. Its VMs, vmnics, vmkernel NICs, datastores and port groups."""
DRIVER_NAME = "vmware_esxi"
@@ -46,6 +47,22 @@ class VmwareEsxiDriver(VmwareActionsMixin, VmwareBaseDriver):
def get_interfaces_ip(self) -> dict[str, dict[str, Any]]:
return host_interfaces_ip(self._host())
def reboot_host(self) -> None:
"""Restart the host, which has to be in maintenance mode.
``force=True`` would restart a host with running VMs by cutting their
power. Evacuating them -- or deciding not to -- is the operator's call,
made by entering maintenance mode in vSphere first.
"""
host = self._host()
if not host.get("runtime.inMaintenanceMode"):
raise RuntimeError(
f"{host.get('name', self.hostname)} is not in maintenance mode; "
"enter it in vSphere first so running VMs are not powered off"
)
# The task finishes as the host goes down; there is nothing to wait for.
invoke(self._mo(vim.HostSystem, host["_moref"]).RebootHost_Task, force=False)
def get_lldp_neighbors(self) -> dict[str, list[dict[str, str]]]:
network_system = self._host().get("configManager.networkSystem")
if not network_system:
+16
View File
@@ -26,6 +26,9 @@ HOST = (
"config.autoStart",
"configIssue",
"configManager.networkSystem",
"parent",
"datastore",
"network",
)
VM = (
@@ -47,6 +50,8 @@ VM = (
"summary.quickStats",
"guest.net",
"guest.ipAddress",
"guest.hostName",
"config.files.vmPathName",
"guest.toolsStatus",
"guest.toolsRunningStatus",
"snapshot",
@@ -66,6 +71,13 @@ DV_PORTGROUP = (
DV_SWITCH = ("name",)
# Placement of new VMs: which network objects exist, and the parent chain from
# a host's compute resource up to its datacenter.
NETWORK = ("name",)
COMPUTE_RESOURCE = ("parent", "resourcePool")
FOLDER = ("parent",)
DATACENTER = ("name", "vmFolder")
#: Managed object type name -> paths, in the order harvest.py dumps them.
ALL = {
"HostSystem": HOST,
@@ -73,4 +85,8 @@ ALL = {
"Datastore": DATASTORE,
"DistributedVirtualPortgroup": DV_PORTGROUP,
"DistributedVirtualSwitch": DV_SWITCH,
"Network": NETWORK,
"ComputeResource": COMPUTE_RESOURCE,
"Folder": FOLDER,
"Datacenter": DATACENTER,
}
View File
+135
View File
@@ -0,0 +1,135 @@
"""Cloud images (qcow2/raw) converted to streamOptimized VMDKs, cached by URL.
ESXi cannot boot a qcow2 and cannot download one itself, so the conversion
runs where the driver runs: download, verify, ``qemu-img convert``. The
result is kept, keyed by URL, so the next VM from the same image skips both
steps. Only the converted disk is kept; the download is deleted.
"""
from __future__ import annotations
import hashlib
import json
import os
import shutil
import subprocess
import tempfile
from collections.abc import Callable
from pathlib import Path
import requests
_CHUNK = 1024 * 1024
DEFAULT_CACHE_DIR = Path(tempfile.gettempdir()) / "napalm-vmware-images"
Fetch = Callable[[str, Path, float], None]
Convert = Callable[[Path, Path], None]
VirtualSize = Callable[[Path], int]
def _qemu_img() -> str:
path = shutil.which("qemu-img")
if path is None:
raise RuntimeError(
"qemu-img is not installed where netOrk runs the provisioning job; "
"it is needed to convert cloud images for VMware (package qemu-utils)"
)
return path
def fetch(url: str, dest: Path, timeout: float) -> None: # pragma: no cover - network
with requests.get(url, stream=True, timeout=timeout) as response:
response.raise_for_status()
with dest.open("wb") as out:
for chunk in response.iter_content(_CHUNK):
out.write(chunk)
def convert_to_vmdk(src: Path, dst: Path) -> None:
subprocess.run(
[
_qemu_img(),
"convert",
"-O",
"vmdk",
"-o",
"subformat=streamOptimized",
str(src),
str(dst),
],
check=True,
capture_output=True,
)
def virtual_size(path: Path) -> int:
"""The disk size the image describes, in bytes (not the file size)."""
out = subprocess.run(
[_qemu_img(), "info", "--output", "json", str(path)],
check=True,
capture_output=True,
text=True,
).stdout
return int(json.loads(out)["virtual-size"])
def _digest(path: Path, algorithm: str) -> str:
h = hashlib.new(algorithm)
with path.open("rb") as fh:
for chunk in iter(lambda: fh.read(_CHUNK), b""):
h.update(chunk)
return h.hexdigest()
class ImageCache:
"""Converted images in ``directory``, one ``<key>.vmdk`` plus ``<key>.json`` each."""
def __init__(
self,
directory: Path = DEFAULT_CACHE_DIR,
*,
fetch: Fetch = fetch,
convert: Convert = convert_to_vmdk,
virtual_size: VirtualSize = virtual_size,
) -> None:
self._dir = Path(directory)
self._fetch = fetch
self._convert = convert
self._virtual_size = virtual_size
def vmdk(self, url: str, checksum: str | None, timeout: float) -> tuple[Path, int]:
"""``(path to the VMDK, virtual disk size in bytes)`` for ``url``."""
self._dir.mkdir(parents=True, exist_ok=True)
key = hashlib.sha256(url.encode()).hexdigest()[:16]
vmdk, meta = self._dir / f"{key}.vmdk", self._dir / f"{key}.json"
if vmdk.exists() and meta.exists():
return vmdk, int(json.loads(meta.read_text())["virtual_size"])
download = self._dir / f"{key}.download"
try:
self._download_verified(url, checksum, download, timeout)
size = self._virtual_size(download)
partial = self._dir / f"{key}.vmdk.partial"
self._convert(download, partial)
# Rename last: a VMDK that exists is a VMDK that is complete, even
# when two jobs convert the same image at once.
os.replace(partial, vmdk)
meta.write_text(json.dumps({"url": url, "virtual_size": size}))
finally:
download.unlink(missing_ok=True)
return vmdk, size
def _download_verified(
self, url: str, checksum: str | None, dest: Path, timeout: float
) -> None:
algorithm, _, expected = (checksum or "").rpartition(":")
algorithm = (algorithm or "sha256").lower()
for _attempt in (1, 2):
self._fetch(url, dest, timeout)
if not checksum:
return
actual = _digest(dest, algorithm)
if actual.lower() == expected.lower():
return
dest.unlink(missing_ok=True)
raise RuntimeError(f"Checksum mismatch for {url}: expected {expected}, got {actual}")
+120
View File
@@ -0,0 +1,120 @@
"""A minimal OVF 1.0 descriptor for importing one converted cloud image.
vSphere builds the VM from this (``OvfManager.CreateImportSpec``) and then
takes the disk contents over NFC. The descriptor only has to describe what
the image does not: CPU, memory, one disk and the NICs.
Devices are the ones VMware recommends for a 64-bit Linux guest -- PVSCSI and
VMXNET3. Both drivers are in the mainline kernel; whether every distribution's
*cloud* kernel carries them is one of the things #305 checks on real hardware.
"""
from __future__ import annotations
from xml.sax.saxutils import escape, quoteattr
_HW_VERSION = "vmx-13" # ESXi 6.5 and later
_GUEST_OS = "otherLinux64Guest"
_STREAM_OPTIMIZED = "http://www.vmware.com/interfaces/specifications/vmdk.html#streamOptimized"
def _item(**elements: str) -> str:
# CIM requires the rasd elements in alphabetical order.
body = "".join(f"<rasd:{k}>{escape(v)}</rasd:{k}>" for k, v in sorted(elements.items()))
return f"<Item>{body}</Item>"
def _nic_items(nic_count: int, first_instance: int) -> str:
return "".join(
_item(
AutomaticAllocation="true",
Connection=f"net{i}",
ElementName=f"Network adapter {i + 1}",
InstanceID=str(first_instance + i),
ResourceSubType="VmxNet3",
ResourceType="10",
)
for i in range(nic_count)
)
def ovf_descriptor(
name: str,
*,
cpu: int,
memory_mb: int,
capacity_bytes: int,
nic_count: int,
firmware: str = "bios",
) -> str:
networks = "".join(
f'<Network ovf:name="net{i}"><Description>net{i}</Description></Network>'
for i in range(nic_count)
)
hardware = "".join(
[
_item(
AllocationUnits="hertz * 10^6",
ElementName=f"{cpu} virtual CPU(s)",
InstanceID="1",
ResourceType="3",
VirtualQuantity=str(cpu),
),
_item(
AllocationUnits="byte * 2^20",
ElementName=f"{memory_mb} MB of memory",
InstanceID="2",
ResourceType="4",
VirtualQuantity=str(memory_mb),
),
_item(
Address="0",
ElementName="SCSI controller 0",
InstanceID="3",
ResourceSubType="VirtualSCSI",
ResourceType="6",
),
_item(
AddressOnParent="0",
ElementName="Hard disk 1",
HostResource="ovf:/disk/vmdisk1",
InstanceID="4",
Parent="3",
ResourceType="17",
),
_nic_items(nic_count, first_instance=5),
]
)
firmware_config = (
'<vmw:Config ovf:required="false" vmw:key="firmware" vmw:value="efi"/>'
if firmware == "efi"
else ""
)
return (
'<?xml version="1.0" encoding="UTF-8"?>'
'<Envelope xmlns="http://schemas.dmtf.org/ovf/envelope/1"'
' xmlns:ovf="http://schemas.dmtf.org/ovf/envelope/1"'
' xmlns:rasd="http://schemas.dmtf.org/wbem/wscim/1/cim-schema/2/CIM_ResourceAllocationSettingData"'
' xmlns:vssd="http://schemas.dmtf.org/wbem/wscim/1/cim-schema/2/CIM_VirtualSystemSettingData"'
' xmlns:vmw="http://www.vmware.com/schema/ovf">'
'<References><File ovf:href="disk.vmdk" ovf:id="file1"/></References>'
"<DiskSection><Info>Virtual disks</Info>"
f'<Disk ovf:capacity="{capacity_bytes}" ovf:capacityAllocationUnits="byte"'
f' ovf:diskId="vmdisk1" ovf:fileRef="file1" ovf:format="{_STREAM_OPTIMIZED}"/>'
"</DiskSection>"
f"<NetworkSection><Info>Networks</Info>{networks}</NetworkSection>"
f"<VirtualSystem ovf:id={quoteattr(name)}>"
"<Info>A virtual machine</Info>"
f"<Name>{escape(name)}</Name>"
f'<OperatingSystemSection ovf:id="101" vmw:osType="{_GUEST_OS}"><Info>Guest OS</Info>'
"</OperatingSystemSection>"
"<VirtualHardwareSection><Info>Virtual hardware</Info>"
"<System><vssd:ElementName>Virtual Hardware Family</vssd:ElementName>"
"<vssd:InstanceID>0</vssd:InstanceID>"
f"<vssd:VirtualSystemIdentifier>{escape(name)}</vssd:VirtualSystemIdentifier>"
f"<vssd:VirtualSystemType>{_HW_VERSION}</vssd:VirtualSystemType></System>"
f"{hardware}{firmware_config}"
"</VirtualHardwareSection>"
"</VirtualSystem>"
"</Envelope>"
)
+107
View File
@@ -0,0 +1,107 @@
"""Where a new VM goes: a pure decision over plain inventory rows.
A host qualifies when it is connected, not in maintenance mode, and sees both
a usable datastore and every requested network. Among those, the pair with
the most free space wins -- or the named datastore, when one is given. The
resource pool comes from the host's compute resource and the VM folder from
the datacenter above it; a standalone ESXi host has the same shape
(``ha-compute-res`` under ``ha-datacenter``), so one path serves both.
"""
from __future__ import annotations
from dataclasses import dataclass, field
from typing import Any
Row = dict[str, Any]
@dataclass
class Inventory:
hosts: list[Row]
datastores: list[Row]
networks: list[Row]
compute_resources: list[Row]
folders: list[Row] = field(default_factory=list)
datacenters: list[Row] = field(default_factory=list)
@dataclass
class Placement:
host: str
host_name: str
datastore: str
datastore_name: str
resource_pool: str
folder: str
datacenter_name: str
networks: list[str]
def _usable_datastores(inv: Inventory, storage: str | None) -> dict[str, Row]:
usable = {
ds["_moref"]: ds
for ds in inv.datastores
if (ds.get("summary") or {}).get("accessible")
and (ds.get("summary") or {}).get("maintenanceMode", "normal") == "normal"
and (storage is None or ds.get("name") == storage)
}
if storage is not None and not usable:
raise ValueError(f"There is no usable datastore named {storage!r}")
return usable
def _host_networks(
host: Row, names_by_moref: dict[str, str], wanted: list[str]
) -> list[str] | None:
"""The host's network moref for each wanted name, or None if one is missing."""
by_name = {names_by_moref.get(m): m for m in host.get("network") or []}
found = [by_name.get(name) for name in wanted]
return None if None in found else [m for m in found if m]
def _datacenter(inv: Inventory, compute_resource: str) -> Row:
parents = {row["_moref"]: row.get("parent") for row in inv.compute_resources + inv.folders}
datacenters = {dc["_moref"]: dc for dc in inv.datacenters}
node: str | None = compute_resource
while node is not None and node not in datacenters:
node = parents.get(node)
if node is None:
raise RuntimeError(f"No datacenter above {compute_resource!r}")
return datacenters[node]
def choose_placement(inv: Inventory, storage: str | None, networks: list[str]) -> Placement:
usable = _usable_datastores(inv, storage)
names = {n["_moref"]: n.get("name", "") for n in inv.networks}
best: tuple[int, Row, Row, list[str]] | None = None
for host in sorted(inv.hosts, key=lambda h: h.get("name", "")):
if host.get("runtime.connectionState") != "connected" or host.get(
"runtime.inMaintenanceMode"
):
continue
host_nets = _host_networks(host, names, networks)
if host_nets is None:
continue
for ds_ref in host.get("datastore") or []:
ds = usable.get(ds_ref)
free = int((ds or {}).get("summary", {}).get("freeSpace", 0))
if ds is not None and (best is None or free > best[0]):
best = (free, host, ds, host_nets)
if best is None:
raise ValueError(
f"There is no host that sees datastore {storage or '(any)'} and networks {networks}"
)
_, host, ds, host_nets = best
compute = next(c for c in inv.compute_resources if c["_moref"] == host.get("parent"))
dc = _datacenter(inv, compute["_moref"])
return Placement(
host=host["_moref"],
host_name=host.get("name", ""),
datastore=ds["_moref"],
datastore_name=ds.get("name", ""),
resource_pool=compute["resourcePool"],
folder=dc["vmFolder"],
datacenter_name=dc.get("name", ""),
networks=host_nets,
)
+67
View File
@@ -0,0 +1,67 @@
"""cloud-init's NoCloud seed: user-data, meta-data, network-config on an ISO.
NoCloud rather than VMware's guestinfo datasource because it needs nothing in
the guest: guestinfo is read through open-vm-tools, which a generic cloud
image does not ship -- it is one of the things the user-data installs.
"""
from __future__ import annotations
import io
from typing import Any
import pycdlib
import yaml
#: The volume label cloud-init's NoCloud datasource looks for.
_LABEL = "cidata"
def user_data(cloud_init_config: dict[str, Any], ssh_public_keys: list[str] | None) -> str:
"""The ``#cloud-config`` document, with ``ssh_public_keys`` merged in."""
config = dict(cloud_init_config)
if ssh_public_keys:
keys = list(config.get("ssh_authorized_keys") or [])
keys += [k for k in ssh_public_keys if k not in keys]
config["ssh_authorized_keys"] = keys
return "#cloud-config\n" + yaml.safe_dump(config, sort_keys=False)
def meta_data(hostname: str, instance_id: str) -> str:
return yaml.safe_dump({"instance-id": instance_id, "local-hostname": hostname}, sort_keys=False)
def network_config(nics: list[tuple[str, bool]]) -> dict[str, Any] | None:
"""Netplan v2 config: DHCP on every ``(mac, dhcp)`` NIC that asks for it.
``None`` when no NIC does -- cloud-init then leaves networking alone
rather than being told to configure nothing.
"""
ethernets = {
f"nic{i}": {"match": {"macaddress": mac}, "dhcp4": True}
for i, (mac, dhcp) in enumerate(nics)
if dhcp
}
return {"version": 2, "ethernets": ethernets} if ethernets else None
def nocloud_iso(user: str, meta: str, network: dict[str, Any] | None) -> bytes:
"""An ISO 9660 image (Rock Ridge + Joliet) holding the seed files."""
files = {"user-data": user, "meta-data": meta}
if network is not None:
files["network-config"] = yaml.safe_dump(network, sort_keys=False)
iso = pycdlib.PyCdlib()
iso.new(interchange_level=3, joliet=3, rock_ridge="1.09", vol_ident=_LABEL)
for index, (name, content) in enumerate(files.items()):
data = content.encode()
iso.add_fp(
io.BytesIO(data),
len(data),
f"/SEED{index}.;1",
rr_name=name,
joliet_path=f"/{name}",
)
out = io.BytesIO()
iso.write_fp(out)
iso.close()
return out.getvalue()
+113
View File
@@ -0,0 +1,113 @@
"""The two HTTP transfers provisioning needs: a disk over NFC, a file to a datastore.
Both authenticate with the vSphere session's own cookie, so no second login
and no credentials beyond the ones the driver already holds.
"""
from __future__ import annotations
import time
from collections.abc import Callable
from pathlib import Path
from typing import Any
from urllib.parse import quote
import requests
#: vSphere drops an NFC lease that reports no progress for five minutes.
_PROGRESS_EVERY_SECONDS = 20.0
class ProgressFile:
"""A file body with a length, calling ``report(percent)`` as it is read.
Sized so ``requests`` sends it with a Content-Length instead of chunked:
vSphere's NFC endpoint refuses a chunked disk upload.
"""
def __init__(
self,
path: Path,
report: Callable[[int], None],
clock: Callable[[], float] = time.monotonic,
) -> None:
self._fh = path.open("rb")
self._size = path.stat().st_size
self._sent = 0
self._report = report
self._clock = clock
self._last = clock()
def __len__(self) -> int:
return self._size
def read(self, size: int = -1) -> bytes:
chunk = self._fh.read(size)
self._sent += len(chunk)
if chunk and self._clock() - self._last >= _PROGRESS_EVERY_SECONDS:
self._report(min(99, self._sent * 100 // (self._size or 1)))
self._last = self._clock()
return chunk
def close(self) -> None:
self._fh.close()
def upload_disk(
url: str, vmdk: Path, cookie: str, verify: bool, report: Callable[[int], None]
) -> None:
"""Stream a streamOptimized VMDK to an NFC lease's device URL."""
body = ProgressFile(vmdk, report)
try:
response = requests.post(
url,
data=body,
headers={"Content-Type": "application/x-vnd.vmware-streamVmdk", "Cookie": cookie},
verify=verify,
timeout=(30, 600),
)
finally:
body.close()
response.raise_for_status()
def lease_url(url: str, hostname: str, port: int) -> str:
"""An NFC device URL with its ``*`` host filled in.
ESXi answers ``https://*/nfc/...``: "the host you are talking to". A port
other than 443 has to be carried over, or the upload misses a NAT'd host.
Through a vCenter the URL already names the ESXi host that takes the disk.
"""
host = hostname if port == 443 else f"{hostname}:{port}"
return url.replace("://*/", f"://{host}/", 1)
def datastore_url(base: str, path: str) -> str:
"""``https://host/folder/<path>`` -- the datastore file browser endpoint."""
return f"{base}/folder/{quote(path)}"
def upload_file(
base: str,
datacenter: str,
datastore: str,
path: str,
data: bytes,
cookie: str,
verify: bool,
) -> None:
"""Put ``data`` at ``[datastore] path``."""
response = requests.put(
datastore_url(base, path),
params={"dcPath": datacenter, "dsName": datastore},
data=data,
headers={"Content-Type": "application/octet-stream", "Cookie": cookie},
verify=verify,
timeout=(30, 120),
)
response.raise_for_status()
def session_cookie(si: Any) -> str:
"""The ``vmware_soap_session`` cookie of a pyVmomi ServiceInstance."""
return si._stub.cookie
+396
View File
@@ -0,0 +1,396 @@
"""HypervisorDriver provisioning: new VMs from netOrk's cloud-image catalog.
The image is downloaded and converted where the driver runs
(``provision/image.py``), imported over OVF/NFC, and configured by cloud-init
from a NoCloud ISO placed next to the VM's files. See the README for why this
path and not OVA templates or a Content Library.
"""
from __future__ import annotations
import logging
import posixpath
import time
import uuid
from pathlib import Path
from typing import Any
from napalm_device_types.models import (
NetworkTargetDict,
NICConfigDict,
StorageTargetDict,
VMProvisionResultDict,
VMStatusDict,
)
from pyVmomi import vim
from napalm_vmware import paths
from napalm_vmware._plain import to_plain
from napalm_vmware._tasks import fault_message, invoke, wait_for_task, wait_until
from napalm_vmware.parse import vm_devices as dev
from napalm_vmware.parse.networks import virtual_networks
from napalm_vmware.parse.vms import vm_list
from napalm_vmware.provision import transfer
from napalm_vmware.provision.image import ImageCache
from napalm_vmware.provision.ovf import ovf_descriptor
from napalm_vmware.provision.placement import Inventory as PlacementInventory
from napalm_vmware.provision.placement import Placement, choose_placement
from napalm_vmware.provision.seed import meta_data, network_config, nocloud_iso, user_data
logger = logging.getLogger(__name__)
_GB = 1024**3
_SEED_ISO = "cidata.iso"
_TASK_TIMEOUT = 300
_STATUS = {"poweredOn": "running", "poweredOff": "stopped", "suspended": "suspended"}
class VmwareProvisioningMixin:
"""Mixed into both drivers ahead of :class:`VmwareBaseDriver`."""
GUEST_AGENT_PACKAGES: tuple[str, ...] = ("open-vm-tools",)
# The unit is open-vm-tools on Debian/Ubuntu and vmtoolsd on EL/Fedora.
GUEST_AGENT_RUNCMD: tuple[str, ...] = (
"systemctl enable --now open-vm-tools || systemctl enable --now vmtoolsd",
)
hostname: str
_port: int
_verify_ssl: bool
_image_cache_dir: Path
_si: Any
_inventory: Any
_mo: Any
_find_vm: Any
_hosts: Any
_index: Any
# -- choosing targets --------------------------------------------------------
def get_image_storages(self) -> list[StorageTargetDict]:
result: list[StorageTargetDict] = []
for ds in self._inventory.collect("Datastore", paths.DATASTORE):
summary = ds.get("summary") or {}
if (
not summary.get("accessible")
or summary.get("maintenanceMode", "normal") != "normal"
):
continue
result.append(
{
"name": ds.get("name", ""),
"type": str(summary.get("type", "")).lower(),
"total_gb": round(int(summary.get("capacity") or 0) / _GB, 1),
"available_gb": round(int(summary.get("freeSpace") or 0) / _GB, 1),
}
)
return sorted(result, key=lambda s: s["name"])
def get_network_targets(self) -> list[NetworkTargetDict]:
networks = virtual_networks(
self._hosts(),
self._inventory.collect("DistributedVirtualPortgroup", paths.DV_PORTGROUP),
self._inventory.collect("DistributedVirtualSwitch", paths.DV_SWITCH),
)
# The port group fixes the VLAN; a NIC cannot add a tag of its own.
return [
{
"name": net["name"],
"kind": "portgroup",
"vlan_aware": False,
"fixed_vlan_tag": net["vlan_id"] or None,
}
for net in sorted(networks.values(), key=lambda n: n["name"])
]
def _check_nics(self, nics: list[NICConfigDict]) -> None:
if not nics:
raise RuntimeError("A VM needs at least one network adapter")
targets = {t["name"]: t for t in self.get_network_targets()}
for nic in nics:
target = targets.get(nic["bridge"])
if target is None:
raise RuntimeError(f"There is no port group named {nic['bridge']!r}")
if nic.get("trunk_vlan_tags"):
raise RuntimeError("VMware port groups cannot trunk per adapter")
tag = nic.get("vlan_tag")
if tag and tag != target.get("fixed_vlan_tag"):
raise RuntimeError(
f"Port group {nic['bridge']!r} carries VLAN "
f"{target.get('fixed_vlan_tag') or 'untagged'}, not {tag}; "
"choose a port group on the VLAN instead"
)
def _placement_inventory(self) -> PlacementInventory:
collect = self._inventory.collect
return PlacementInventory(
hosts=self._hosts(),
datastores=collect("Datastore", paths.DATASTORE),
networks=collect("Network", paths.NETWORK),
compute_resources=collect("ComputeResource", paths.COMPUTE_RESOURCE),
folders=collect("Folder", paths.FOLDER),
datacenters=collect("Datacenter", paths.DATACENTER),
)
# -- creating ----------------------------------------------------------------
def create_vm_from_cloud_init(
self,
name: str,
*,
image_url: str,
cpu: int,
memory: int,
nics: list[NICConfigDict],
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:
self._check_nics(nics)
try:
vmdk, image_size = ImageCache(self._image_cache_dir).vmdk(
image_url, image_checksum, download_timeout
)
place = choose_placement(
self._placement_inventory(), storage, [n["bridge"] for n in nics]
)
except ValueError as exc:
raise RuntimeError(str(exc)) from exc
descriptor = ovf_descriptor(
name, cpu=cpu, memory_mb=memory, capacity_bytes=image_size, nic_count=len(nics)
)
vm_ref = self._import_ovf(name, descriptor, vmdk, place, timeout)
try:
self._finish_vm(
vm_ref,
name,
place,
nics,
user_data(cloud_init_config, ssh_public_keys),
disk_gb=disk_resize_gb,
image_size=image_size,
)
self._run_task(self._mo(vim.VirtualMachine, vm_ref).PowerOnVM_Task)
except Exception as exc:
logger.warning("Provisioning %s failed after import, removing it: %s", name, exc)
self._discard(vm_ref)
raise RuntimeError(f"Provisioning {name!r} failed: {exc}") from exc
props = self._inventory.properties(
self._mo(vim.VirtualMachine, vm_ref), ["config.instanceUuid"]
)
return {"vmid": props["config.instanceUuid"], "name": name, "node": place.host_name}
def _run_task(self, call: Any, *args: Any, **kwargs: Any) -> None:
wait_for_task(self._inventory, invoke(call, *args, **kwargs), _TASK_TIMEOUT)
def _import_ovf( # pragma: no cover - live NFC session; covered by tests/test_vcsim.py
self, name: str, descriptor: str, vmdk: Path, place: Placement, timeout: int
) -> str:
"""Create the VM from ``descriptor`` and stream ``vmdk`` into it; return its MoRef."""
content = self._si.RetrieveContent()
pool = self._mo(vim.ResourcePool, place.resource_pool)
params = vim.OvfManager.CreateImportSpecParams(
entityName=name,
diskProvisioning="thin",
networkMapping=[
vim.OvfManager.NetworkMapping(name=f"net{i}", network=self._mo(vim.Network, ref))
for i, ref in enumerate(place.networks)
],
)
spec = invoke(
content.ovfManager.CreateImportSpec,
descriptor,
pool,
self._mo(vim.Datastore, place.datastore),
params,
)
if spec.error:
raise RuntimeError("; ".join(e.msg or type(e).__name__ for e in spec.error))
lease = invoke(
pool.ImportVApp,
spec.importSpec,
self._mo(vim.Folder, place.folder),
self._mo(vim.HostSystem, place.host),
)
try:
wait_until(
lambda: (
self._inventory.properties(lease, ["state"]).get("state") in ("ready", "error")
),
timeout,
"the import lease",
)
info = self._inventory.properties(lease, ["state", "error", "info"])
if info.get("state") == "error":
raise RuntimeError(fault_message(info.get("error") or {}))
url = transfer.lease_url(info["info"]["deviceUrl"][0]["url"], self.hostname, self._port)
transfer.upload_disk(
url,
vmdk,
transfer.session_cookie(self._si),
self._verify_ssl,
report=lambda pct: lease.HttpNfcLeaseProgress(pct),
)
lease.HttpNfcLeaseComplete()
except Exception:
try:
lease.HttpNfcLeaseAbort()
except Exception: # noqa: BLE001 - the original error is the one to report
pass
raise
return info["info"]["entity"]
def _finish_vm(
self,
vm_ref: str,
name: str,
place: Placement,
nics: list[NICConfigDict],
user: str,
*,
disk_gb: int | None,
image_size: int,
) -> None:
"""MACs, disk size and the cloud-init seed: everything the OVF could not say."""
vm_mo = self._mo(vim.VirtualMachine, vm_ref)
props = self._inventory.properties(
vm_mo, ["config.hardware.device", "config.files.vmPathName"], raw=True
)
# Live device objects: an "edit" replaces the device with what is sent,
# so it has to be the whole device, backing included.
devices = list(props.get("config.hardware.device") or [])
adapters = [d for d in devices if dev.is_nic(to_plain(d))]
changes = _mac_changes(adapters, nics) + _disk_growth(devices, disk_gb, image_size)
macs = [
(nic.get("mac") or adapter.macAddress or "").lower()
for adapter, nic in zip(adapters, nics)
]
seed = nocloud_iso(
user,
meta_data(name, f"iid-{uuid.uuid4()}"),
network_config(
[(mac, nic.get("dhcp", i == 0)) for i, (mac, nic) in enumerate(zip(macs, nics))]
),
)
vm_dir = posixpath.dirname(str(props["config.files.vmPathName"]).split("] ", 1)[1])
self._upload_seed(place, f"{vm_dir}/{_SEED_ISO}", seed)
iso_path = f"[{place.datastore_name}] {vm_dir}/{_SEED_ISO}"
changes += _seed_cdrom(iso_path)
self._run_task(vm_mo.ReconfigVM_Task, spec=vim.vm.ConfigSpec(deviceChange=changes))
def _upload_seed(self, place: Placement, path: str, data: bytes) -> None: # pragma: no cover
transfer.upload_file(
f"https://{self.hostname}:{self._port}",
place.datacenter_name,
place.datastore_name,
path,
data,
transfer.session_cookie(self._si),
self._verify_ssl,
)
def _discard(self, vm_ref: str) -> None:
vm_mo = self._mo(vim.VirtualMachine, vm_ref)
try:
state = self._inventory.properties(vm_mo, ["runtime.powerState"]).get(
"runtime.powerState"
)
if state == "poweredOn":
self._run_task(vm_mo.PowerOffVM_Task)
self._run_task(vm_mo.Destroy_Task)
except Exception as exc: # noqa: BLE001 - report the provisioning error, not this one
logger.error("Could not remove half-provisioned VM %s: %s", vm_ref, exc)
# -- after creation ------------------------------------------------------------
def destroy_vm(self, vmid: str, *, remove_disk: bool = True, timeout: int = 60) -> None:
try:
vm = self._find_vm(vmid)
except ValueError as exc:
raise RuntimeError(str(exc)) from exc
vm_mo = self._mo(vim.VirtualMachine, vm["_moref"])
if vm.get("runtime.powerState") == "poweredOn":
self._run_task(vm_mo.PowerOffVM_Task)
if remove_disk:
self._run_task(vm_mo.Destroy_Task)
else:
invoke(vm_mo.UnregisterVM)
def get_vm_status(
self, vmid: str, *, wait_for_ip: bool = False, timeout: int = 300, poll_interval: int = 5
) -> VMStatusDict:
deadline = time.monotonic() + timeout
while True:
status = self._vm_status(vmid)
if not wait_for_ip or status.get("ip_address"):
return status
if time.monotonic() >= deadline:
raise RuntimeError(f"VM {vmid} reported no IP address within {timeout}s")
time.sleep(poll_interval)
def _vm_status(self, vmid: str) -> VMStatusDict:
try:
vm = self._find_vm(vmid)
except ValueError as exc:
raise RuntimeError(str(exc)) from exc
hosts = self._hosts()
(entry,) = vm_list([vm], hosts, self._index(hosts))
status: VMStatusDict = {"status": _STATUS.get(vm.get("runtime.powerState", ""), "unknown")}
if entry["ipv4"]:
status["ip_address"] = entry["ipv4"]
if vm.get("guest.hostName"):
status["hostname"] = vm["guest.hostName"]
first_nic = next(iter(entry["interfaces"].values()), None)
if first_nic and first_nic["mac_address"]:
status["mac_address"] = first_nic["mac_address"]
return status
def _mac_changes(adapters: list[Any], nics: list[NICConfigDict]) -> list[Any]:
"""Device edits pinning the MACs the caller asked for."""
changes = []
for adapter, nic in zip(adapters, nics):
if not nic.get("mac"):
continue
adapter.addressType = "manual"
adapter.macAddress = nic["mac"].lower()
changes.append(vim.vm.device.VirtualDeviceSpec(operation="edit", device=adapter))
return changes
def _disk_growth(devices: list[Any], disk_gb: int | None, image_size: int) -> list[Any]:
if not disk_gb or disk_gb * _GB <= image_size:
return []
disk = next(d for d in devices if isinstance(d, vim.vm.device.VirtualDisk))
disk.capacityInKB = disk_gb * 1024 * 1024
# Both fields, or the live object carries the old size in the other one;
# pyVmomi's stubs type capacityInBytes as None, hence setattr.
setattr(disk, "capacityInBytes", disk_gb * _GB) # noqa: B010
return [vim.vm.device.VirtualDeviceSpec(operation="edit", device=disk)]
def _seed_cdrom(iso_path: str) -> list[Any]:
"""A SATA controller and a CD-ROM with the seed ISO, both new.
The OVF describes only a SCSI controller; a CD-ROM needs IDE or SATA.
Negative keys are placeholders vSphere replaces, and let the CD-ROM name
its not-yet-existing controller.
"""
controller = vim.vm.device.VirtualAHCIController(key=-101, busNumber=0)
cdrom = vim.vm.device.VirtualCdrom(
key=-102,
controllerKey=-101,
unitNumber=0,
backing=vim.vm.device.VirtualCdrom.IsoBackingInfo(fileName=iso_path),
connectable=vim.vm.device.VirtualDevice.ConnectInfo(
startConnected=True, connected=True, allowGuestControl=True
),
)
return [
vim.vm.device.VirtualDeviceSpec(operation="add", device=controller),
vim.vm.device.VirtualDeviceSpec(operation="add", device=cdrom),
]
+2 -1
View File
@@ -14,9 +14,10 @@ from napalm_device_types import FingerprintRule
from napalm_vmware.actions import VmwareActionsMixin
from napalm_vmware.base import VmwareBaseDriver
from napalm_vmware.parse.facts import vcenter_facts
from napalm_vmware.provisioning import VmwareProvisioningMixin
class VmwareVcenterDriver(VmwareActionsMixin, VmwareBaseDriver):
class VmwareVcenterDriver(VmwareProvisioningMixin, VmwareActionsMixin, VmwareBaseDriver):
"""A vCenter and every VM, datastore and port group it manages."""
DRIVER_NAME = "vmware_vcenter"