"""PropertyCollector helpers: filter specs, updates, folder traversal props.""" from __future__ import annotations import json import re import secrets from typing import Any from xml.sax.saxutils import escape from app.db.pool import Database from app.vsphere import inventory from app.vsphere.inventory import ManagedObject # Types that inherit ManagedEntity (govmomi Ancestors / Finder propSet Type=ManagedEntity). _MANAGED_ENTITY_TYPES = frozenset( { "Folder", "Datacenter", "VirtualMachine", "HostSystem", "ClusterComputeResource", "ComputeResource", "ResourcePool", "VirtualApp", "StoragePod", "Datastore", "Network", "DistributedVirtualPortgroup", "VmwareDistributedVirtualSwitch", "DistributedVirtualSwitch", "OpaqueNetwork", } ) def _decode(value: Any) -> Any: current = value while isinstance(current, str): try: current = json.loads(current) except json.JSONDecodeError: break return current async def clear_pc_state(database: Database | None = None) -> None: """Drop durable PC views/tokens/versions (and orphan ContainerView inventory rows).""" if database is None: return pool = database.pool # type: ignore[attr-defined] async with pool.acquire() as conn: await conn.execute( """ DO $$ BEGIN IF to_regclass('public.vsphere_pc_state') IS NOT NULL THEN DELETE FROM vsphere_pc_state; END IF; END $$; """ ) await conn.execute("DELETE FROM vsphere_objects WHERE type = 'ContainerView'") async def _pc_put(database: Database, kind: str, key: str, payload: Any) -> None: pool = database.pool # type: ignore[attr-defined] async with pool.acquire() as conn: await conn.execute( """ INSERT INTO vsphere_pc_state (kind, key, payload, updated_at) VALUES ($1, $2, $3::jsonb, now()) ON CONFLICT (kind, key) DO UPDATE SET payload = EXCLUDED.payload, updated_at = now() """, kind, key, json.dumps(payload), ) async def _pc_get(database: Database, kind: str, key: str) -> Any | None: pool = database.pool # type: ignore[attr-defined] async with pool.acquire() as conn: row = await conn.fetchrow( "SELECT payload FROM vsphere_pc_state WHERE kind = $1 AND key = $2", kind, key, ) if row is None: return None return _decode(row["payload"]) async def _pc_pop(database: Database, kind: str, key: str) -> Any | None: pool = database.pool # type: ignore[attr-defined] async with pool.acquire() as conn: row = await conn.fetchrow( """ DELETE FROM vsphere_pc_state WHERE kind = $1 AND key = $2 RETURNING payload """, kind, key, ) if row is None: return None return _decode(row["payload"]) def parse_path_sets(body: str) -> list[str]: return re.findall(r"<(?:\w+:)?pathSet[^>]*>([^<]+)", body) def parse_obj_refs(body: str) -> list[tuple[str, str]]: return [ (m.group(1), m.group(2)) for m in re.finditer( r'<(?:\w+:)?obj[^>]*type="([^"]+)"[^>]*>([^<]+)', body, ) ] def parse_view_types(body: str) -> list[str]: return re.findall(r"<(?:\w+:)?type[^>]*>([^<]+)", body) def parse_prop_types(body: str) -> list[str]: """Types listed under propSet (what the client wants returned).""" return re.findall( r"<(?:\w+:)?propSet\b[^>]*>\s*<(?:\w+:)?type[^>]*>([^<]+)", body, flags=re.DOTALL, ) def parse_continue_token(body: str) -> str | None: match = re.search(r"<(?:\w+:)?token[^>]*>([^<]+)", body) return match.group(1).strip() if match else None async def next_view_id(database: Database) -> str: pool = database.pool # type: ignore[attr-defined] async with pool.acquire() as conn: row = await conn.fetchrow( """ INSERT INTO vsphere_pc_state (kind, key, payload, updated_at) VALUES ('meta', 'view_seq', '{"n": 1}'::jsonb, now()) ON CONFLICT (kind, key) DO UPDATE SET payload = jsonb_build_object( 'n', COALESCE((vsphere_pc_state.payload->>'n')::int, 0) + 1 ), updated_at = now() RETURNING payload """ ) n = int((_decode(row["payload"]) or {}).get("n") or 1) return f"view-{n}" async def register_container_view(database: Database, view_id: str, moids: list[str]) -> None: payload = {"moids": list(moids)} await _pc_put(database, "view", view_id, payload) await inventory.upsert_object( database, moid=view_id, type_name="ContainerView", name=view_id, parent_moid=None, props={"view_moids": list(moids)}, ) def view_moids_from_object(obj: ManagedObject) -> list[str]: props = obj.props or {} return [str(m) for m in (props.get("view_moids") or [])] async def store_page_token( database: Database, remaining_moids: list[str], path_sets: list[str] | None = None, ) -> str: token = f"token-{secrets.token_hex(8)}" await _pc_put( database, "token", token, {"moids": list(remaining_moids), "path_sets": list(path_sets or [])}, ) return token async def take_page_token(database: Database, token: str) -> list[str] | None: payload = await _pc_pop(database, "token", token) if payload is None: return None if isinstance(payload, list): return payload return list(payload.get("moids") or []) async def take_page_token_full( database: Database, token: str ) -> tuple[list[str], list[str]] | None: payload = await _pc_pop(database, "token", token) if payload is None: return None if isinstance(payload, list): return payload, [] return list(payload.get("moids") or []), list(payload.get("path_sets") or []) def _xsi(type_name: str) -> str: return f'xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:type="{escape(type_name)}"' def _prop(name: str, value: str, *, xsi_type: str | None = None) -> str: # govmomi SOAP AnyType decode requires xsi:type; bare leaves Name empty → Ancestors panic. type_name = xsi_type or "xsd:string" return ( f"{escape(name)}" f"{escape(value)}" ) def _prop_bool(name: str, value: bool) -> str: return _prop(name, "true" if value else "false", xsi_type="xsd:boolean") def _prop_int(name: str, value: int) -> str: return _prop(name, str(value), xsi_type="xsd:int") def _prop_long(name: str, value: int) -> str: return _prop(name, str(value), xsi_type="xsd:long") def _prop_raw(name: str, inner_xml: str, *, xsi_type: str) -> str: return f"{escape(name)}{inner_xml}" def _prop_mor(name: str, type_name: str, moid: str) -> str: return ( f"{escape(name)}" f'' f"{escape(moid)}" ) def _prop_array_mor(name: str, items: list[tuple[str, str]]) -> str: if not items: return ( f"{escape(name)}" f"" ) inner = "".join( f'{escape(m)}' for t, m in items ) return ( f"{escape(name)}" f"{inner}" ) def _vim_power(state: str) -> str: return { "POWERED_ON": "poweredOn", "POWERED_OFF": "poweredOff", "SUSPENDED": "suspended", }.get(state, "poweredOff") def _parent_type(moid: str) -> str: if moid.startswith("group-"): return "Folder" if moid.startswith("datacenter-"): return "Datacenter" if moid.startswith("domain-"): return "ClusterComputeResource" if moid.startswith("host-"): return "HostSystem" if moid.startswith("resgroup-"): return "ResourcePool" if moid.startswith("network-"): return "Network" if moid.startswith("dvportgroup-"): return "DistributedVirtualPortgroup" if moid.startswith("dvs-"): return "VmwareDistributedVirtualSwitch" if moid.startswith("datastore-"): return "Datastore" return "ManagedEntity" def _tools_status(props: dict[str, Any]) -> str: raw = str(props.get("tools_status") or "") mapping = { "GUEST_TOOLS_RUNNING": "toolsOk", "GUEST_TOOLS_NOT_RUNNING": "toolsNotRunning", "GUEST_TOOLS_OLD": "toolsOld", } return mapping.get( raw, "toolsOk" if props.get("power_state") == "POWERED_ON" else "toolsNotRunning" ) def _guest_id(props: dict[str, Any]) -> str: guest = str(props.get("guest_OS") or "OTHER_GUEST_64") # REST enums often use UBUNTU_64_GUEST; SOAP guestId uses ubuntu64Guest. mapping = { "UBUNTU_64_GUEST": "ubuntu64Guest", "CENTOS_64_GUEST": "centos64Guest", "RHEL_8_64_GUEST": "rhel8_64Guest", "WINDOWS_2019_64_GUEST": "windows2019srv_64Guest", "OTHER_GUEST_64": "otherGuest64", "OTHER_GUEST": "otherGuest", } if guest in mapping: return mapping[guest] if guest.endswith("_GUEST"): return guest.lower().replace("_guest", "Guest").replace("_64", "64") return guest def _instance_uuid(obj: ManagedObject) -> str: identity = obj.props.get("identity") or {} if identity.get("instance_uuid"): return str(identity["instance_uuid"]) # Stable synthetic UUID from moid digits. digits = "".join(ch for ch in obj.moid if ch.isdigit()) or "0" n = int(digits) % 10_000_000_000_000 return f"5029aaaa-bbbb-cccc-dddd-{n:012d}" def _bios_uuid(obj: ManagedObject) -> str: identity = obj.props.get("identity") or {} if identity.get("bios_uuid"): return str(identity["bios_uuid"]) digits = "".join(ch for ch in obj.moid if ch.isdigit()) or "0" n = int(digits) % 10_000_000_000_000 return f"4200aaaa-bbbb-cccc-dddd-{n:012d}" def _folder_child_types(obj: ManagedObject, children: list[ManagedObject]) -> list[str]: """Declare Folder.childType the way vCenter does (not ManagedEntity).""" from_children = sorted({c.type for c in children}) name = (obj.name or "").lower() moid = obj.moid or "" if name in {"datacenters", "datacenter"} or moid.startswith("group-d"): return ["Folder", "Datacenter"] if name == "vm" or moid.startswith("group-v"): return ["Folder", "VirtualMachine", "VirtualApp"] if name == "host" or moid.startswith("group-h"): return ["Folder", "ComputeResource", "ClusterComputeResource", "HostSystem"] if name in {"datastore", "datastores"} or moid.startswith("group-s"): return ["Folder", "Datastore", "StoragePod"] if name == "network" or moid.startswith("group-n"): return [ "Folder", "Network", "DistributedVirtualPortgroup", "VmwareDistributedVirtualSwitch", "DistributedVirtualSwitch", ] if from_children: return ["Folder", *from_children] return ["Folder", "VirtualMachine"] def _vm_devices_xml(props: dict[str, Any]) -> str: """Minimal ArrayOfVirtualDevice suitable for govmomi / Terraform device reads.""" xsi = 'xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"' datastore = str(props.get("datastore") or "datastore-31") devices: list[str] = [] # PCI root + SCSI — Terraform device sorter requires controller assignment. devices.append( f'' "100" "PCI controller 0" "0" ) devices.append( f'' "1000" "LSI Logic" "1003" "0noSharing" "7" "" ) devices.append( f'' "200IDE 0" "0" ) # Terraform/pulumi-vsphere getVirtualMachine requires a video card device. devices.append( f'' "500" "Video card" "1000" "4096" "1false" "false" "" ) disks = list(props.get("disks") or []) if not disks: disks = [{"key": "2000", "value": {"label": "Hard disk 1", "capacity": 42949672960}}] for disk in disks: key = int(str(disk.get("key") or 2000)) value = disk.get("value") or disk capacity = int(value.get("capacity") or 42949672960) capacity_kb = max(capacity // 1024, 1) label = escape(str(value.get("label") or f"Hard disk {key}")) uuid = escape(str(value.get("uuid") or f"6000C29{key:08x}-0000-0000-0000-000000000001")) file_name = escape(str(value.get("file_name") or f"[datastore1] disk-{key}.vmdk")) devices.append( f'' f"{key}" f"{capacity_kb} KB" f"1000{max(key - 2000, 0)}" f"{capacity_kb}" f"{capacity}" f'' f"{file_name}" f'{escape(datastore)}' f"persistent" f"true" f"false" f"false" f"{uuid}" f"" ) nics = list(props.get("nics") or []) if not nics: nics = [ { "key": "4000", "value": { "label": "Network adapter 1", "mac": "00:50:56:01:00:01", "state": "CONNECTED", "backing": {"network": "network-41"}, }, } ] for nic in nics: key = int(str(nic.get("key") or 4000)) value = nic.get("value") or nic backing = value.get("backing") or {} network = str(backing.get("network") or "network-41") mac = escape(str(value.get("mac") or "00:50:56:00:00:01")) label = escape(str(value.get("label") or f"Network adapter {key}")) connected = "true" if str(value.get("state") or "").upper() == "CONNECTED" else "false" devices.append( f'' f"{key}" f"VMXNET3" f"1007" f"{mac}" f"assigned" f"{connected}" f"truetrue" f"" "0-1" "50normal" "" f'' f"VM Network" f'{escape(network)}' f"" ) for cdrom in props.get("cdroms") or []: key = int(str(cdrom.get("cdrom") or cdrom.get("key") or 3000)) label = escape(str(cdrom.get("label") or "CD/DVD drive 1")) devices.append( f'' f"{key}" f"CD/DVD" f"2000" f'' f"[datastore1] ISO/ubuntu.iso" f"" ) return "".join(devices) def build_prop_map( obj: ManagedObject, *, children: list[ManagedObject], all_by_moid: dict[str, ManagedObject], ) -> dict[str, str]: """Map property path -> already-escaped XML fragment for ….""" props = obj.props out: dict[str, str] = { "name": _prop("name", obj.name), "overallStatus": _prop("overallStatus", "green"), } if obj.parent_moid: out["parent"] = _prop_mor("parent", _parent_type(obj.parent_moid), obj.parent_moid) child_items = [(c.type, c.moid) for c in children] if obj.type == "Folder": out["childEntity"] = _prop_array_mor("childEntity", child_items) # Folder.childType is ArrayOfString (govmomi folder helpers reject scalar "ManagedEntity"). child_types = _folder_child_types(obj, children) out["childType"] = _prop_raw( "childType", "".join(f"{escape(t)}" for t in child_types), xsi_type="ArrayOfString", ) if obj.type == "Datacenter": out["hostFolder"] = _prop_mor( "hostFolder", "Folder", str(props.get("host_folder") or "group-h23") ) out["vmFolder"] = _prop_mor( "vmFolder", "Folder", str(props.get("vm_folder") or "group-v23") ) out["datastoreFolder"] = _prop_mor( "datastoreFolder", "Folder", str(props.get("datastore_folder") or "group-s23") ) out["networkFolder"] = _prop_mor( "networkFolder", "Folder", str(props.get("network_folder") or "group-n23") ) datastores = [("Datastore", o.moid) for o in all_by_moid.values() if o.type == "Datastore"] networks = [ (o.type, o.moid) for o in all_by_moid.values() if o.type in {"Network", "DistributedVirtualPortgroup"} ] out["datastore"] = _prop_array_mor("datastore", datastores) out["network"] = _prop_array_mor("network", networks) if obj.type == "ClusterComputeResource": hosts = [ ("HostSystem", c.moid) for c in all_by_moid.values() if c.type == "HostSystem" and c.parent_moid == obj.moid ] out["host"] = _prop_array_mor("host", hosts) out["resourcePool"] = _prop_mor( "resourcePool", "ResourcePool", str(props.get("resource_pool") or "resgroup-22") ) # Terraform CreateVM loads default devices via EnvironmentBrowser. env_browser = str(props.get("environment_browser") or f"envbrowser-{obj.moid}") out["environmentBrowser"] = _prop_mor( "environmentBrowser", "EnvironmentBrowser", env_browser ) out["datastore"] = _prop_array_mor( "datastore", [("Datastore", o.moid) for o in all_by_moid.values() if o.type == "Datastore"][:8], ) out["network"] = _prop_array_mor( "network", [ (o.type, o.moid) for o in all_by_moid.values() if o.type in {"Network", "DistributedVirtualPortgroup"} ][:20], ) out["name"] = _prop("name", obj.name) out["summary.effectiveCpu"] = _prop_int("summary.effectiveCpu", 48000) out["summary.effectiveMemory"] = _prop_long("summary.effectiveMemory", 256 * 1024) out["summary.numHosts"] = _prop_int("summary.numHosts", len(hosts)) out["configuration.drsConfig.enabled"] = _prop_bool("configuration.drsConfig.enabled", True) out["configuration.dasConfig.enabled"] = _prop_bool( "configuration.dasConfig.enabled", False ) if obj.type == "ResourcePool": out["owner"] = _prop_mor("owner", "ClusterComputeResource", obj.parent_moid or "domain-c21") vms = [ ("VirtualMachine", o.moid) for o in all_by_moid.values() if o.type == "VirtualMachine" and str(o.props.get("resource_pool") or "") == obj.moid ] if not vms: vms = [ ("VirtualMachine", o.moid) for o in all_by_moid.values() if o.type == "VirtualMachine" and o.parent_moid and "group-v" in (o.parent_moid or "") ][:50] out["vm"] = _prop_array_mor("vm", vms) out["resourcePool"] = _prop_array_mor("resourcePool", []) out["config.cpuAllocation.limit"] = _prop_long("config.cpuAllocation.limit", -1) out["config.memoryAllocation.limit"] = _prop_long("config.memoryAllocation.limit", -1) if obj.type == "VirtualMachine": power = _vim_power(str(props.get("power_state", "POWERED_OFF"))) cpu = int(props.get("cpu_count") or 1) mem = int(props.get("memory_size_mib") or 1024) guest_id = _guest_id(props) inst_uuid = _instance_uuid(obj) bios_uuid = _bios_uuid(obj) host = str(props.get("host") or "") datastore = str(props.get("datastore") or "datastore-31") ds_obj = all_by_moid.get(datastore) datastore_name = ds_obj.name if ds_obj is not None else "datastore1" pool = str(props.get("resource_pool") or "resgroup-22") networks = [str(n) for n in (props.get("networks") or ["network-41"])] guest_ip = str(props.get("guest_ip") or props.get("ip_address") or "") hostname = str((props.get("identity") or {}).get("name") or obj.name) template = bool(props.get("template")) tools = _tools_status(props) out["runtime.powerState"] = _prop("runtime.powerState", power) out["summary.runtime.powerState"] = _prop("summary.runtime.powerState", power) out["config.hardware.numCPU"] = _prop_int("config.hardware.numCPU", cpu) out["config.hardware.memoryMB"] = _prop_int("config.hardware.memoryMB", mem) out["config.hardware.numCoresPerSocket"] = _prop_int("config.hardware.numCoresPerSocket", 1) out["config.uuid"] = _prop("config.uuid", bios_uuid) out["config.instanceUuid"] = _prop("config.instanceUuid", inst_uuid) out["config.guestId"] = _prop("config.guestId", guest_id) out["config.guestFullName"] = _prop( "config.guestFullName", str(props.get("guest_OS") or guest_id) ) out["config.name"] = _prop("config.name", obj.name) out["config.version"] = _prop( "config.version", str(props.get("hardware_version") or "vmx-19").lower() ) out["config.template"] = _prop_bool("config.template", template) out["config.files.vmPathName"] = _prop( "config.files.vmPathName", f"[{datastore_name}] {obj.name}/{obj.name}.vmx" ) out["config.hardware.device"] = _prop_raw( "config.hardware.device", _vm_devices_xml(props), xsi_type="ArrayOfVirtualDevice" ) # Terraform flattenVirtualMachineConfigInfo requires non-nil Tools / allocations / bootOptions. out["config.tools"] = _prop_raw( "config.tools", "manual" "true" "true" "true" "true" "true" "false" "true", xsi_type="ToolsConfigInfo", ) out["config.firmware"] = _prop("config.firmware", str(props.get("firmware") or "bios")) out["config.changeVersion"] = _prop("config.changeVersion", "1") out["config.annotation"] = _prop("config.annotation", str(props.get("annotation") or "")) out["config.alternateGuestName"] = _prop("config.alternateGuestName", "") out["config.memoryHotAddEnabled"] = _prop_bool("config.memoryHotAddEnabled", False) out["config.cpuHotAddEnabled"] = _prop_bool("config.cpuHotAddEnabled", False) out["config.cpuHotRemoveEnabled"] = _prop_bool("config.cpuHotRemoveEnabled", False) out["config.memoryReservationLockedToMax"] = _prop_bool( "config.memoryReservationLockedToMax", False ) out["config.nestedHVEnabled"] = _prop_bool("config.nestedHVEnabled", False) out["config.vPMCEnabled"] = _prop_bool("config.vPMCEnabled", False) out["config.swapPlacement"] = _prop("config.swapPlacement", "inherit") out["config.flags"] = _prop_raw( "config.flags", "true" "false" "automatic" "hvAuto" "false" "release" "false" "false", xsi_type="VirtualMachineFlagInfo", ) alloc = ( "0-1" "1000normal" ) out["config.cpuAllocation"] = _prop_raw( "config.cpuAllocation", alloc, xsi_type="ResourceAllocationInfo" ) out["config.memoryAllocation"] = _prop_raw( "config.memoryAllocation", "0-1" "20480normal", xsi_type="ResourceAllocationInfo", ) out["config.bootOptions"] = _prop_raw( "config.bootOptions", "0" "false" "false" "10000" "false", xsi_type="VirtualMachineBootOptions", ) out["summary.config.numCpu"] = _prop_int("summary.config.numCpu", cpu) out["summary.config.memorySizeMB"] = _prop_int("summary.config.memorySizeMB", mem) out["summary.config.name"] = _prop("summary.config.name", obj.name) out["summary.config.uuid"] = _prop("summary.config.uuid", bios_uuid) out["summary.config.instanceUuid"] = _prop("summary.config.instanceUuid", inst_uuid) out["summary.config.guestId"] = _prop("summary.config.guestId", guest_id) out["summary.guest.ipAddress"] = _prop("summary.guest.ipAddress", guest_ip) out["summary.guest.hostName"] = _prop("summary.guest.hostName", hostname) out["summary.guest.toolsStatus"] = _prop("summary.guest.toolsStatus", tools) out["guest.ipAddress"] = _prop("guest.ipAddress", guest_ip) out["guest.hostName"] = _prop("guest.hostName", hostname) out["guest.toolsStatus"] = _prop("guest.toolsStatus", tools) out["guest.guestId"] = _prop("guest.guestId", guest_id) out["guest.guestState"] = _prop( "guest.guestState", "running" if power == "poweredOn" else "notRunning" ) if host: out["runtime.host"] = _prop_mor("runtime.host", "HostSystem", host) out["summary.runtime.host"] = _prop_mor("summary.runtime.host", "HostSystem", host) out["resourcePool"] = _prop_mor("resourcePool", "ResourcePool", pool) out["datastore"] = _prop_array_mor("datastore", [("Datastore", datastore)]) net_items: list[tuple[str, str]] = [] for net_id in networks: net_obj = all_by_moid.get(net_id) net_items.append((net_obj.type if net_obj else "Network", net_id)) out["network"] = _prop_array_mor("network", net_items) out["layoutEx.file"] = _prop_raw( "layoutEx.file", f"0" f"[{escape(datastore_name)}] {escape(obj.name)}/{escape(obj.name)}.vmx" f"config4096", xsi_type="ArrayOfVirtualMachineFileLayoutExFileInfo", ) if obj.type == "HostSystem": cpu_mhz = int(props.get("cpu_mhz") or 2400) cpu_cores = int(props.get("cpu_cores") or props.get("num_cpu_cores") or 16) mem_bytes = int(props.get("memory_bytes") or props.get("memory_size") or 137438953472) out["runtime.connectionState"] = _prop("runtime.connectionState", "connected") out["runtime.powerState"] = _prop("runtime.powerState", "poweredOn") out["runtime.inMaintenanceMode"] = _prop_bool( "runtime.inMaintenanceMode", bool(props.get("maintenance_mode")) ) out["summary.config.name"] = _prop("summary.config.name", obj.name) out["summary.hardware.uuid"] = _prop( "summary.hardware.uuid", str(props.get("uuid") or f"host-uuid-{obj.moid}") ) out["summary.hardware.memorySize"] = _prop_long("summary.hardware.memorySize", mem_bytes) out["summary.hardware.numCpuCores"] = _prop_int("summary.hardware.numCpuCores", cpu_cores) out["summary.hardware.cpuMhz"] = _prop_int("summary.hardware.cpuMhz", cpu_mhz) out["hardware.memorySize"] = _prop_long("hardware.memorySize", mem_bytes) out["hardware.cpuInfo.numCpuCores"] = _prop_int("hardware.cpuInfo.numCpuCores", cpu_cores) out["hardware.cpuInfo.hz"] = _prop_long("hardware.cpuInfo.hz", cpu_mhz * 1_000_000) out["config.product.version"] = _prop( "config.product.version", str(props.get("version") or "8.0.0") ) out["config.product.fullName"] = _prop( "config.product.fullName", f"VMware ESXi {props.get('version') or '8.0.0'}" ) ds_ids = ( props.get("datastores") or [o.moid for o in all_by_moid.values() if o.type == "Datastore"][:4] ) out["datastore"] = _prop_array_mor("datastore", [("Datastore", str(d)) for d in ds_ids]) net_ids = ( props.get("networks") or [o.moid for o in all_by_moid.values() if o.type == "Network"][:4] ) out["network"] = _prop_array_mor("network", [("Network", str(n)) for n in net_ids]) vms = [ ("VirtualMachine", o.moid) for o in all_by_moid.values() if o.type == "VirtualMachine" and str(o.props.get("host") or "") == obj.moid ] out["vm"] = _prop_array_mor("vm", vms) out["environmentBrowser"] = _prop_mor( "environmentBrowser", "EnvironmentBrowser", str(props.get("environment_browser") or f"envbrowser-{obj.moid}"), ) if obj.type == "EnvironmentBrowser" or obj.moid.startswith("envbrowser-"): out["name"] = _prop("name", obj.name) if obj.type == "Datastore": capacity = int(props.get("capacity") or 0) free = int(props.get("free_space") or 0) out["summary.capacity"] = _prop_long("summary.capacity", capacity) out["summary.freeSpace"] = _prop_long("summary.freeSpace", free) out["summary.type"] = _prop("summary.type", str(props.get("type") or "VMFS")) out["summary.name"] = _prop("summary.name", obj.name) out["summary.url"] = _prop("summary.url", f"ds:///vmfs/volumes/{obj.moid}/") out["summary.accessible"] = _prop_bool("summary.accessible", True) out["summary.multipleHostAccess"] = _prop_bool("summary.multipleHostAccess", True) out["info.name"] = _prop("info.name", obj.name) out["info.url"] = _prop("info.url", f"ds:///vmfs/volumes/{obj.moid}/") out["info.freeSpace"] = _prop_long("info.freeSpace", free) out["info.maxFileSize"] = _prop_long("info.maxFileSize", 62 * 1024**4) hosts = [o.moid for o in all_by_moid.values() if o.type == "HostSystem"][:20] # Datastore.host is ArrayOfDatastoreHostMount — not ArrayOfManagedObjectReference. host_mounts = "".join( "" f'{escape(hid)}' "" f"/vmfs/volumes/{escape(obj.moid)}" "readWrite" "true" "true" "" "" for hid in hosts ) out["host"] = _prop_raw("host", host_mounts, xsi_type="ArrayOfDatastoreHostMount") vms = [ ("VirtualMachine", o.moid) for o in all_by_moid.values() if o.type == "VirtualMachine" and str(o.props.get("datastore") or "") == obj.moid ][:100] out["vm"] = _prop_array_mor("vm", vms) if obj.type in {"Network", "DistributedVirtualPortgroup", "VmwareDistributedVirtualSwitch"}: out["summary.name"] = _prop("summary.name", obj.name) out["summary.accessible"] = _prop_bool("summary.accessible", True) hosts = [("HostSystem", o.moid) for o in all_by_moid.values() if o.type == "HostSystem"][ :20 ] out["host"] = _prop_array_mor("host", hosts) vms = [ ("VirtualMachine", o.moid) for o in all_by_moid.values() if o.type == "VirtualMachine" and ( obj.moid in (o.props.get("networks") or []) or obj.name in (o.props.get("networks") or []) ) ][:100] out["vm"] = _prop_array_mor("vm", vms) if obj.type == "DistributedVirtualPortgroup": key = str(props.get("key") or obj.moid) out["key"] = _prop("key", key) out["config.name"] = _prop("config.name", obj.name) out["config.key"] = _prop("config.key", key) dvs = str(props.get("dvs") or props.get("distributed_switch") or "dvs-41") out["config.distributedVirtualSwitch"] = _prop_mor( "config.distributedVirtualSwitch", "VmwareDistributedVirtualSwitch", dvs ) if obj.type == "VmwareDistributedVirtualSwitch": out["uuid"] = _prop("uuid", str(props.get("uuid") or f"dvs-uuid-{obj.moid}")) out["summary.name"] = _prop("summary.name", obj.name) out["config.uuid"] = _prop( "config.uuid", str(props.get("uuid") or f"dvs-uuid-{obj.moid}") ) pgs = [ ("DistributedVirtualPortgroup", o.moid) for o in all_by_moid.values() if o.type == "DistributedVirtualPortgroup" ] out["portgroup"] = _prop_array_mor("portgroup", pgs) if obj.type == "ContainerView" or obj.moid.startswith("view-"): moids = view_moids_from_object(obj) items = [] for moid in moids: child = all_by_moid.get(moid) if child: items.append((child.type, child.moid)) out["view"] = _prop_array_mor("view", items) if obj.type == "Task": state = str(props.get("state") or "success") progress = 100 if state == "success" else 50 xsi = 'xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"' result_xml = "" entity_xml = "" if props.get("result_moid"): rtype = str(props.get("result_type") or "VirtualMachine") rmoid = escape(str(props["result_moid"])) result_xml = ( f'' f"{rmoid}" ) entity_xml = ( f'' f"{rmoid}{escape(rtype)}" ) info_inner = ( f"{escape(obj.moid)}" f'{escape(obj.moid)}' f"{escape(state)}" "falsefalse" f"{escape(str(props.get('operation') or 'task'))}" f"{progress}" f"{entity_xml}{result_xml}" ) out["info"] = _prop_raw("info", info_inner, xsi_type="TaskInfo") out["info.state"] = _prop("info.state", state) out["info.descriptionId"] = _prop( "info.descriptionId", str(props.get("operation") or "task") ) if props.get("entity"): out["info.entity"] = _prop_mor("info.entity", "VirtualMachine", str(props["entity"])) out["info.progress"] = _prop_int("info.progress", progress) if props.get("result_moid"): out["info.result"] = _prop_mor( "info.result", str(props.get("result_type") or "VirtualMachine"), str(props["result_moid"]), ) if obj.type == "HttpNfcLease" or obj.moid.startswith("lease-"): # Lease document is durable in vsphere_objects.props (+ vsphere_nfc_leases). lease = props if isinstance(props, dict) else {} state = str(lease.get("state") or "ready") out["state"] = _prop("state", state) out["info.state"] = _prop("info.state", state) out["info.entity"] = _prop_mor( "info.entity", "VirtualMachine", str(lease.get("entity") or "vm-101") ) out["info.initializeProgress"] = _prop_int( "info.initializeProgress", int(lease.get("initializeProgress") or 100) ) out["info.transferProgress"] = _prop_int( "info.transferProgress", int(lease.get("transferProgress") or 0) ) urls = (lease.get("info") or {}).get("deviceUrl") or [] url_xml = "".join( f"{escape(str(u.get('key')))}" f"{escape(str(u.get('importKey')))}" f"{escape(str(u.get('url')))}" f"{escape(str(u.get('sslThumbprint') or ''))}" f"" for u in urls ) out["info.deviceUrl"] = _prop_raw( "info.deviceUrl", url_xml, xsi_type="ArrayOfHttpNfcLeaseDeviceUrl" ) return out def object_content_xml( obj: ManagedObject, *, children: list[ManagedObject], all_by_moid: dict[str, ManagedObject], path_sets: list[str] | None = None, ) -> str: if obj.props.get("_not_found"): xsi = 'xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"' # Shape must decode to types.ManagedObjectNotFound so terraform # viapi.IsManagedObjectNotFoundError is true (CreateVM vs CreateChildVM). # MissingProperty.fault is the MethodFault itself (no nested ). path = (path_sets[0] if path_sets else None) or "name" return ( f'{escape(obj.moid)}' f"{escape(path)}" f'' f'{escape(obj.moid)}' f"" f"" ) prop_map = build_prop_map(obj, children=children, all_by_moid=all_by_moid) if path_sets: selected = [] for path in path_sets: if path in prop_map: selected.append(prop_map[path]) else: # prefix match e.g. runtime / config for key, frag in prop_map.items(): if key == path or key.startswith(path + "."): selected.append(frag) prop_xml = selected or list(prop_map.values()) else: prop_xml = list(prop_map.values()) return ( f'{escape(obj.moid)}' + "".join(prop_xml) + "" ) def resolve_inventory_path(path: str, objects: list[ManagedObject]) -> ManagedObject | None: """Resolve inventory paths (govmomi omits root ``Datacenters`` folder).""" parts = [p for p in path.strip("/").split("/") if p] if not parts: return None by_moid = {o.moid: o for o in objects} def _walk(start_parent: str | None) -> ManagedObject | None: current_parent = start_parent match: ManagedObject | None = None for part in parts: found = None if current_parent: found = find_child(current_parent, part, objects) if found is None: found = next( (o for o in objects if o.parent_moid == current_parent and o.name == part), None, ) if found is None and current_parent is None: hits = [o for o in objects if o.name == part] if len(hits) == 1: found = hits[0] # Standard DC folder names even if the Folder.name was renamed in lab. if found is None and current_parent and current_parent in by_moid: parent = by_moid[current_parent] if parent.type == "Datacenter": key = { "vm": "vm_folder", "host": "host_folder", "datastore": "datastore_folder", "network": "network_folder", }.get(part.lower()) if key and parent.props.get(key): found = by_moid.get(str(parent.props[key])) if found is None: return None match = found current_parent = found.moid return match hit = _walk(None) if hit is not None: return hit for root in (o for o in objects if o.parent_moid is None): hit = _walk(root.moid) if hit is not None: return hit return None def find_child(parent_moid: str, name: str, objects: list[ManagedObject]) -> ManagedObject | None: """SearchIndex.FindChild — direct child by name.""" hits = [o for o in objects if o.parent_moid == parent_moid and o.name == name] if len(hits) == 1: return hits[0] # Datacenter folder props (vmFolder etc.) are not parent_moid edges for DC itself. parent = next((o for o in objects if o.moid == parent_moid), None) if parent and parent.type == "Datacenter": for key in ("host_folder", "vm_folder", "datastore_folder", "network_folder"): folder_moid = parent.props.get(key) if not folder_moid: continue folder = next((o for o in objects if o.moid == folder_moid and o.name == name), None) if folder: return folder if hits: return hits[0] return None def _descendants( roots: list[ManagedObject], all_objects: list[ManagedObject], *, max_depth: int = 32, ) -> list[ManagedObject]: """BFS descendants by parent_moid (folder/datacenter/cluster inventory walk).""" by_parent: dict[str | None, list[ManagedObject]] = {} for obj in all_objects: by_parent.setdefault(obj.parent_moid, []).append(obj) seen: dict[str, ManagedObject] = {o.moid: o for o in roots} frontier = list(roots) depth = 0 while frontier and depth < max_depth: nxt: list[ManagedObject] = [] for obj in frontier: for child in by_parent.get(obj.moid, []): if child.moid not in seen: seen[child.moid] = child nxt.append(child) # Datacenter folders live as props, not parent_moid edges. if obj.type == "Datacenter": for key in ("host_folder", "vm_folder", "datastore_folder", "network_folder"): folder_moid = obj.props.get(key) if not folder_moid: continue folder = next((o for o in all_objects if o.moid == folder_moid), None) if folder and folder.moid not in seen: seen[folder.moid] = folder nxt.append(folder) if obj.type == "ClusterComputeResource": rp = obj.props.get("resource_pool") if rp: pool = next((o for o in all_objects if o.moid == rp), None) if pool and pool.moid not in seen: seen[pool.moid] = pool nxt.append(pool) frontier = nxt depth += 1 return list(seen.values()) async def select_objects_for_retrieve( database: Any, body: str, ) -> tuple[list[ManagedObject], list[str]]: path_sets = parse_path_sets(body) prop_types = parse_prop_types(body) continue_token = parse_continue_token(body) all_objects = await inventory.list_objects(database) by_moid = {o.moid: o for o in all_objects} if continue_token: remaining = await take_page_token(database, continue_token) if remaining is None: return [], path_sets selected = [by_moid[m] for m in remaining if m in by_moid] return selected, path_sets refs = parse_obj_refs(body) selected: list[ManagedObject] = [] if refs: for type_name, moid in refs: if type_name == "ContainerView" or moid.startswith("view-"): view_obj = by_moid.get(moid) if view_obj is not None: selected.append(view_obj) else: stored = await _pc_get(database, "view", moid) moids = ( list((stored or {}).get("moids") or []) if isinstance(stored, dict) else [] ) selected.append( ManagedObject( moid=moid, type="ContainerView", name=moid, parent_moid=None, props={"view_moids": moids}, ) ) elif type_name == "Task" or moid.startswith("task-"): selected.append( ManagedObject( moid=moid, type="Task", name=moid, parent_moid=None, props={"task": True}, ) ) elif type_name == "HttpNfcLease" or moid.startswith("lease-"): selected.append( ManagedObject( moid=moid, type="HttpNfcLease", name=moid, parent_moid=None, props={"lease": True}, ) ) elif type_name == "EnvironmentBrowser" or moid.startswith("envbrowser-"): selected.append( ManagedObject( moid=moid, type="EnvironmentBrowser", name=moid, parent_moid=None, props={"environment_browser": True}, ) ) elif moid in by_moid: obj = by_moid[moid] # Honor requested MOR type — VirtualApp:resgroup-X must NOT resolve a ResourcePool. if type_name and not _type_compatible(type_name, obj.type): # Return the underlying ResourcePool for VirtualApp probes so the # provider proceeds to CreateChildVM_Task (which we implement) instead # of failing/hanging on ManagedObjectNotFound fault decoding. if type_name == "VirtualApp" and obj.type == "ResourcePool": selected.append(obj) continue selected.append( ManagedObject( moid=moid, type=type_name, name=moid, parent_moid=None, props={"_not_found": True}, ) ) continue selected.append(obj) elif moid in {"propertyCollector", "TaskManager", "SearchIndex", "ViewManager"}: continue else: pass wants_traversal = "selectSet" in body or "TraversalSpec" in body if wants_traversal and selected: # Include ContainerView — govmomi view.Retrieve traverses Path=view from the view MOR. inventory_roots = [o for o in selected if o.type != "Task"] if inventory_roots: if _wants_parent_traversal(body): expanded = _ancestors(inventory_roots, all_objects) elif _is_recursive_traversal(body): expanded = _descendants(inventory_roots, all_objects) else: # govmomi list.Lister TraversalSpec is one hop (e.g. Folder.childEntity). expanded = _one_level_traverse(inventory_roots, all_objects, body) skip_root = _object_set_skip(body) if prop_types: type_set = set(prop_types) filtered = [o for o in expanded if _type_matches(o.type, type_set)] if not skip_root: for root in inventory_roots: if root.type == "ContainerView": continue if root.moid not in {o.moid for o in filtered}: if _type_matches(root.type, type_set) or not type_set: filtered.append(root) selected = ( filtered if filtered else [o for o in expanded if o.type != "ContainerView"] ) else: selected = [o for o in expanded if o.type != "ContainerView" or not skip_root] if not skip_root: for root in inventory_roots: if root.type == "ContainerView": continue if root.moid not in {o.moid for o in selected}: selected.append(root) else: selected = all_objects if prop_types: type_set = set(prop_types) selected = [o for o in selected if _type_matches(o.type, type_set)] return selected, path_sets def children_of(moid: str, all_objects: list[ManagedObject]) -> list[ManagedObject]: return [o for o in all_objects if o.parent_moid == moid] def _type_compatible(requested: str, actual: str) -> bool: """Whether an inventory object may satisfy a client MOR of ``requested`` type.""" if not requested or requested == actual: return True if requested == "ManagedEntity": return actual in _MANAGED_ENTITY_TYPES if requested == "ComputeResource" and actual in {"ComputeResource", "ClusterComputeResource"}: return True if requested == "ResourcePool" and actual in {"ResourcePool", "VirtualApp"}: return True # Plain ResourcePool must NOT satisfy VirtualApp (CreateChildVM vs CreateVM). # Exception: none — keep strict. CreateChildVM path is handled separately. if requested == "VirtualApp": return actual == "VirtualApp" if requested == "DistributedVirtualSwitch" and actual in { "VmwareDistributedVirtualSwitch", "DistributedVirtualSwitch", }: return True if requested == "Network" and actual in { "Network", "DistributedVirtualPortgroup", "OpaqueNetwork", }: return True return False def _type_matches(obj_type: str, type_set: set[str]) -> bool: if obj_type in type_set: return True return any(_type_compatible(t, obj_type) for t in type_set) def _traversal_paths(body: str) -> list[str]: """Paths from TraversalSpec (), not PropertySpec ().""" return re.findall(r"<(?:\w+:)?path(?:\s[^>]*)?>([^<]+)", body) def _wants_parent_traversal(body: str) -> bool: """govmomi mo.Ancestors walks parent/parentVApp — not inventory descendants.""" paths = {p.strip() for p in _traversal_paths(body)} if not paths: return "traverseParent" in body parentish = {"parent", "parentVApp"} childish = { "childEntity", "hostFolder", "vmFolder", "datastoreFolder", "networkFolder", "host", "resourcePool", "vm", "datastore", "network", "view", "portgroup", } return bool(paths & parentish) and not bool(paths & childish) def _is_recursive_traversal(body: str) -> bool: """True when a TraversalSpec nests a SelectionSpec (name-only selectSet) to recurse.""" if _wants_parent_traversal(body): return True # Nested SelectionSpec: inside TraversalSpec. # Do not match the TraversalSpec's own sibling of . return bool( re.search( r'xsi:type="TraversalSpec"[^>]*>[\s\S]*?' r"<(?:\w+:)?selectSet(?:\s[^>]*)?>\s*" r"<(?:\w+:)?name(?:\s[^>]*)?>[^<]+\s*" r"", body, flags=re.IGNORECASE, ) ) def _object_set_skip(body: str) -> bool: """objectSet/@skip — ListFolder uses skip=true so the folder itself is omitted.""" match = re.search( r"<(?:\w+:)?objectSet\b[^>]*>.*?<(?:\w+:)?skip[^>]*>([^<]+)", body, flags=re.DOTALL | re.IGNORECASE, ) if not match: return False return match.group(1).strip().lower() == "true" def _ancestors(roots: list[ManagedObject], all_objects: list[ManagedObject]) -> list[ManagedObject]: by_moid = {o.moid: o for o in all_objects} seen: dict[str, ManagedObject] = {} for root in roots: cur: ManagedObject | None = root while cur is not None: if cur.moid in seen: break seen[cur.moid] = cur parent_moid = cur.parent_moid if not parent_moid: break cur = by_moid.get(parent_moid) return list(seen.values()) def _follow_path( root: ManagedObject, path: str, all_objects: list[ManagedObject], by_moid: dict[str, ManagedObject], ) -> list[ManagedObject]: """One PropertyCollector TraversalSpec hop (govmomi list.Lister).""" if path == "childEntity": return children_of(root.moid, all_objects) if path in {"parent", "parentVApp"}: if not root.parent_moid: return [] parent = by_moid.get(root.parent_moid) return [parent] if parent else [] folder_keys = { "vmFolder": "vm_folder", "hostFolder": "host_folder", "datastoreFolder": "datastore_folder", "networkFolder": "network_folder", } if path in folder_keys: folder_moid = root.props.get(folder_keys[path]) if not folder_moid: return [] folder = by_moid.get(str(folder_moid)) return [folder] if folder else [] if path == "resourcePool": rp = root.props.get("resource_pool") or root.props.get("resourcePool") if rp and str(rp) in by_moid: return [by_moid[str(rp)]] return [o for o in all_objects if o.type == "ResourcePool" and o.parent_moid == root.moid] if path == "host": return [o for o in all_objects if o.type == "HostSystem" and o.parent_moid == root.moid] if path == "vm": if root.type == "ResourcePool": return [ o for o in all_objects if o.type == "VirtualMachine" and str(o.props.get("resource_pool") or "") == root.moid ] if root.type == "HostSystem": return [ o for o in all_objects if o.type == "VirtualMachine" and str(o.props.get("host") or "") == root.moid ] return children_of(root.moid, all_objects) if path == "datastore": if root.type == "Datacenter": return [o for o in all_objects if o.type == "Datastore"] ds_ids = root.props.get("datastores") or root.props.get("datastore") or [] if isinstance(ds_ids, str): ds_ids = [ds_ids] out = [by_moid[str(d)] for d in ds_ids if str(d) in by_moid] if out: return out return [o for o in all_objects if o.type == "Datastore"][:20] if path == "network": return [ o for o in all_objects if o.type in {"Network", "DistributedVirtualPortgroup", "OpaqueNetwork"} ][:50] if path == "portgroup": return [o for o in all_objects if o.type == "DistributedVirtualPortgroup"] if path == "view": moids = view_moids_from_object(root) return [by_moid[m] for m in moids if m in by_moid] return [] def _one_level_traverse( roots: list[ManagedObject], all_objects: list[ManagedObject], body: str, ) -> list[ManagedObject]: paths = [p.strip() for p in _traversal_paths(body)] or ["childEntity"] by_moid = {o.moid: o for o in all_objects} seen: dict[str, ManagedObject] = {} for root in roots: for path in paths: for hit in _follow_path(root, path, all_objects, by_moid): seen[hit.moid] = hit return list(seen.values()) def _vim_task_state(status: str) -> str: return { "SUCCEEDED": "success", "FAILED": "error", "RUNNING": "running", "PENDING": "queued", }.get(status, "error") def _parse_max_wait_seconds(body: str) -> float: match = re.search( r"<(?:\w+:)?maxWaitSeconds[^>]*>([^<]*)", body, ) if not match: return 0.0 try: return max(0.0, float(match.group(1).strip() or "0")) except ValueError: return 0.0 def _task_object_set_xml(task: dict[str, Any]) -> str: """PropertyCollector update so govmomi ``task.Wait`` observes completion. govmomi waits on the whole ``info`` property (TaskInfo), not ``info.state``. """ task_id = str(task.get("task") or "") state = _vim_task_state(str(task.get("status") or "FAILED")) progress = 100 if state in {"success", "error"} else int(task.get("progress") or 50) result = task.get("result") if isinstance(task.get("result"), dict) else {} vm_moid = str((result or {}).get("vm") or "") xsi = 'xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"' result_xml = "" entity_xml = "" if vm_moid: result_xml = ( f'' f"{escape(vm_moid)}" ) entity_xml = ( f'' f"{escape(vm_moid)}" "VirtualMachine" ) info_xml = ( f'' f"{escape(task_id)}" f'{escape(task_id)}' f"{escape(state)}" f"false" f"false" f"{escape(str(task.get('operation') or 'task'))}" f"{progress}" f"{entity_xml}{result_xml}" f"" ) return ( "" f'{escape(task_id)}' "modify" f"infoassign{info_xml}" "" ) def _vim_power_state(raw: str) -> str: return { "POWERED_ON": "poweredOn", "POWERED_OFF": "poweredOff", "SUSPENDED": "suspended", }.get(raw, "poweredOff") def _vm_runtime_object_set_xml(obj: ManagedObject) -> str: """Push VM runtime/guest props so post-CreateVM waiters can finish.""" props = obj.props if isinstance(obj.props, dict) else {} power = _vim_power_state(str(props.get("power_state") or "POWERED_OFF")) identity = props.get("identity") if isinstance(props.get("identity"), dict) else {} uuid = str(identity.get("instance_uuid") or props.get("uuid") or f"uuid-{obj.moid}") bios = str(identity.get("bios_uuid") or f"bios-{obj.moid}") xsi = 'xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"' changes = [ f"nameassign{escape(obj.name)}", f"runtime.powerStateassign" f'{escape(power)}', f"summary.runtime.powerStateassign" f'{escape(power)}', f"config.uuidassign{escape(uuid)}", f"config.instanceUuidassign" f"{escape(uuid)}", f"summary.config.uuidassign" f"{escape(uuid)}", f"summary.config.instanceUuidassign" f"{escape(uuid)}", f"summary.config.guestIdassign" f"otherGuest64", f"guest.toolsStatusassign" f'toolsOk', f"guest.toolsRunningStatusassign" f"guestToolsRunning", f"guest.ipAddressassign" f"192.168.1.50", f"summary.guest.ipAddressassign" f"192.168.1.50", f"config.hardware.deviceChangeassign" f'false', ] del bios # reserved for future BIOS UUID prop parity return ( "" f'{escape(obj.moid)}' "modify" + "".join(changes) + "" ) async def wait_updates_xml( database: Database, *, session_key: str, body: str, objects: list[ManagedObject], ) -> str: from app.vsphere.domain import tasks as task_store version_match = re.search( r"<(?:\w+:)?version[^>]*>([^<]*)", body, ) client_version = (version_match.group(1).strip() if version_match else "") or "" stored = await _pc_get(database, "version", session_key) last = int((stored or {}).get("n") or 0) if isinstance(stored, dict) else int(stored or 0) notified_raw = await _pc_get(database, "meta", f"tasks_notified:{session_key}") notified: set[str] = set() if isinstance(notified_raw, dict): notified = {str(x) for x in (notified_raw.get("ids") or [])} elif isinstance(notified_raw, list): notified = {str(x) for x in notified_raw} all_tasks = await task_store.list_tasks(database) new_tasks = [t for t in all_tasks if str(t.get("task") or "") not in notified] synced = bool(client_version) and int(client_version or "0") >= last and last > 0 max_wait = _parse_max_wait_seconds(body) # Late task.Wait callers often subscribe after CreateVM already SUCCEEDED and after # an earlier WaitForUpdates marked that task "notified". Re-assert the newest task # once while the client is blocking (maxWaitSeconds > 0). replay_tasks: list[dict[str, Any]] = [] if synced and not new_tasks and max_wait > 0: newest = next( (t for t in all_tasks if str(t.get("status") or "") in {"SUCCEEDED", "FAILED"}), None, ) if newest is None: return f"{last}" # Cap replay storms: only bump when client still has maxWait (blocking wait). replay_tasks = [newest] elif synced and not new_tasks: return f"{last}" next_ver = max(last, 0) + 1 await _pc_put(database, "version", session_key, {"n": next_ver}) objects_xml: list[str] = [] by_moid = {o.moid: o for o in objects} # First update for a session: seed inventory name props (Finder / ContainerView). if last == 0: limited: list[ManagedObject] = [] vm_count = 0 for obj in objects: if obj.type == "VirtualMachine": if vm_count >= 200: continue vm_count += 1 limited.append(obj) for obj in limited: objects_xml.append( "" f'{escape(obj.moid)}' "enter" f"nameassign" f"{escape(obj.name)}" "" ) emit_tasks = new_tasks or replay_tasks vm_ids: set[str] = set() for task in emit_tasks: objects_xml.append(_task_object_set_xml(task)) result = task.get("result") if isinstance(task.get("result"), dict) else {} vm_moid = str((result or {}).get("vm") or "") if vm_moid: vm_ids.add(vm_moid) for vm_moid in sorted(vm_ids): obj = by_moid.get(vm_moid) if obj is not None: objects_xml.append(_vm_runtime_object_set_xml(obj)) if new_tasks: notified.update(str(t.get("task") or "") for t in new_tasks) await _pc_put(database, "meta", f"tasks_notified:{session_key}", {"ids": sorted(notified)}) return ( "" f"{next_ver}" 'filter-1' + "".join(objects_xml) + "" )