diff --git a/atlas/api/models.py b/atlas/api/models.py index 943dd81df..f74ab98ec 100644 --- a/atlas/api/models.py +++ b/atlas/api/models.py @@ -6,7 +6,16 @@ import frappe from frappe.utils import get_datetime, get_system_timezone -from pydantic import AnyHttpUrl, BaseModel, ConfigDict, Field, StringConstraints, model_validator +from pydantic import ( + AnyHttpUrl, + BaseModel, + ConfigDict, + Discriminator, + Field, + StringConstraints, + Tag, + model_validator, +) from atlas.api.core.base import ListQuery, PatchPayload, StrictModel from atlas.api.core.errors import ApiErrorField @@ -27,6 +36,7 @@ FirewallRule, VirtualMachineCreateRequest, ) +from atlas.vm.core.placement.affinity import AffinityOperator, AffinityResource, AffinityRules if TYPE_CHECKING: from atlas.metal_server.doctype.public_ip_allocation.public_ip_allocation import ( @@ -80,8 +90,16 @@ class PlacementBusyError(BaseModel): fields: list[ApiErrorField] = Field(description="Invalid request fields, or an empty list.") +class AffinityUnsatisfiedError(BaseModel): + """No host with room meets the affinity rules of the VM.""" + + code: Literal["affinity_unsatisfied"] = Field(description="Stable machine-readable error code.") + message: str = Field(description="Safe description of the failure.") + fields: list[ApiErrorField] = Field(description="Invalid request fields, or an empty list.") + + CapacityError = Annotated[ - OutOfCapacityError | PlacementBusyError, + OutOfCapacityError | PlacementBusyError | AffinityUnsatisfiedError, Field(discriminator="code"), ] @@ -410,6 +428,58 @@ def validate_effective_rule_count(self) -> FirewallUpdatePayload: return self +class AffinityRulePayload(StrictModel): + """One rule on the tags of the candidate Metal Server, or of a VM that runs on it.""" + + resource: AffinityResource = Field( + description="`metal_server` checks the candidate host. `virtual_machine` checks the VMs on the candidate host." + ) + operator: AffinityOperator = Field( + description="`has` needs a resource with every tag pair. `has_not` rejects such a resource." + ) + tags: TagMap = Field(min_length=1, description="Tag pairs that must all be on one resource.") + within: Annotated[str, StringConstraints(min_length=1, max_length=MAXIMUM_TAG_KEY_LENGTH)] | None = Field( + default=None, + description="A host tag key, such as `rack`. A `virtual_machine` rule then reads the VMs on every host that has the same value for this key as the candidate host. A host without the key fails the rule.", + ) + + +class AffinityAnyOfPayload(StrictModel): + """A group that holds when at least one of its rules or groups holds.""" + + any_of: list[AffinityNodePayload] = Field(min_length=1, description="Rules or groups. One must hold.") + + +class AffinityAllOfPayload(StrictModel): + """A group that holds when every one of its rules or groups holds.""" + + all_of: list[AffinityNodePayload] = Field( + min_length=1, description="Rules or groups. Every one must hold." + ) + + +def _affinity_node_kind(value: object) -> str | None: + """Pick the node model by its keys, so an invalid node reports the errors of one model.""" + if isinstance(value, BaseModel): + keys = type(value).model_fields + elif isinstance(value, dict): + keys = value + else: + return None + + return next((kind for kind in ("any_of", "all_of") if kind in keys), "rule") + + +AffinityNodePayload = Annotated[ + Annotated[AffinityRulePayload, Tag("rule")] + | Annotated[AffinityAnyOfPayload, Tag("any_of")] + | Annotated[AffinityAllOfPayload, Tag("all_of")], + Discriminator(_affinity_node_kind), +] +AffinityAnyOfPayload.model_rebuild() +AffinityAllOfPayload.model_rebuild() + + class CreateVirtualMachinePayload(StrictModel): """Values that create one virtual machine.""" @@ -451,6 +521,11 @@ class CreateVirtualMachinePayload(StrictModel): firewall: FirewallPayload = Field( default_factory=FirewallPayload, description="Desired firewall configuration." ) + tags: TagMap = Field(default_factory=dict) + placement_rules: list[AffinityNodePayload] = Field( + default_factory=list, + description="Rules that limit the Metal Servers for the virtual machine. Every listed rule or group must hold. Placement uses only the hosts that meet them.", + ) @model_validator(mode="after") def validate_ipv4_internet_access(self) -> CreateVirtualMachinePayload: @@ -458,6 +533,12 @@ def validate_ipv4_internet_access(self) -> CreateVirtualMachinePayload: raise ValueError("A public IPv4 address needs ipv4_internet_access.") return self + @model_validator(mode="after") + def validate_placement_rules(self) -> CreateVirtualMachinePayload: + """Apply the shared affinity rule limits.""" + AffinityRules.from_value([node.model_dump() for node in self.placement_rules]) + return self + def to_domain_request(self, tenant_id: int, image_name: str) -> VirtualMachineCreateRequest: """Build the domain request for this API payload.""" return VirtualMachineCreateRequest( @@ -481,6 +562,8 @@ def to_domain_request(self, tenant_id: int, image_name: str) -> VirtualMachineCr firewall=FirewallConfiguration.from_value(self.firewall.model_dump()), public_ipv4=self.public_ipv4, public_ipv6=self.public_ipv6, + tags=self.tags, + placement_rules=AffinityRules.from_value([node.model_dump() for node in self.placement_rules]), ) diff --git a/atlas/api/routes/virtual_machines.py b/atlas/api/routes/virtual_machines.py index e4657d7bc..2823239f9 100644 --- a/atlas/api/routes/virtual_machines.py +++ b/atlas/api/routes/virtual_machines.py @@ -88,7 +88,8 @@ def request_virtual_machine_power_state( "The virtual machine was not placed. `error.code` is `out_of_capacity` when no host " "can hold it, which needs more capacity in the region, or `placement_busy` when " "every candidate host was held by another placement, which only needs a retry. A " - "busy response carries `Retry-After` in seconds." + "busy response carries `Retry-After` in seconds. `affinity_unsatisfied` means that no " + "host with room meets the affinity rules of the VM." ), "model": CapacityUnavailableResponse, "headers": { @@ -108,6 +109,8 @@ def create_virtual_machine( Creates a tenant VM from an image and requests the specified compute, disk, network, and guest configuration. Only tenant 0 can set `is_privileged`, which lets the VM reach every tenant through the mesh. Set `is_termination_protected` to refuse deletion of the new VM. The termination protection route changes it later. + + Use `tags` to label the VM, for example, `{"role": "cargo-server"}`. `placement_rules` limits the Metal Servers for the VM by host tags and by the tags of other VMs on the host. Atlas Settings selects whether a VM that no host with room can satisfy fails with `affinity_unsatisfied` or goes to any host. """ image = get_owned_image(payload.image_id) request = payload.to_domain_request(get_current_tenant_id(), image.name) diff --git a/atlas/api/tests/test_virtual_machines.py b/atlas/api/tests/test_virtual_machines.py index 35527c631..c6de70cce 100644 --- a/atlas/api/tests/test_virtual_machines.py +++ b/atlas/api/tests/test_virtual_machines.py @@ -36,6 +36,7 @@ from atlas.vm.core.metal_models import MetalFirewall, MetalFirewallRule from atlas.vm.core.models import DEFAULT_ROUTES, Route from atlas.vm.core.placement import OutOfCapacity +from atlas.vm.core.placement.affinity import MAXIMUM_AFFINITY_RULES CREATE_BODY = { "image_id": "system-image", @@ -267,6 +268,41 @@ def test_create_passes_the_firewall(self) -> None: self.assertTrue(firewall.enabled) self.assertEqual(firewall.inbound[0].cidrs, ("203.0.113.0/24",)) + def test_create_passes_tags_and_placement_rules(self) -> None: + rules = [ + { + "any_of": [ + {"resource": "metal_server", "operator": "has", "tags": {"rack": "a"}}, + {"resource": "metal_server", "operator": "has", "tags": {"rack": "b"}}, + ] + }, + {"resource": "virtual_machine", "operator": "has_not", "tags": {"role": "cargo-server"}}, + ] + + status, _, create = self.create( + {**CREATE_BODY, "tags": {"role": "cargo-server"}, "placement_rules": rules} + ) + + self.assertEqual(status, 202) + request = create.call_args.args[0] + self.assertEqual(request.tags, {"role": "cargo-server"}) + self.assertEqual(request.placement_rules.as_list(), rules) + + def test_create_rejects_invalid_placement_rules(self) -> None: + rule = {"resource": "virtual_machine", "operator": "has_not", "tags": {"role": "cargo-server"}} + for placement_rules in ( + [{**rule, "weight": 1}], + [{"resource": "metal_server", "operator": "has", "tags": {"rack": "a"}, "within": "rack"}], + [{"any_of": []}], + [rule] * (MAXIMUM_AFFINITY_RULES + 1), + ): + with self.subTest(placement_rules=placement_rules): + status, body, create = self.create({**CREATE_BODY, "placement_rules": placement_rules}) + + self.assertEqual(status, 400) + self.assertEqual(body["error"]["code"], "invalid_request") + create.assert_not_called() + def test_create_passes_both_public_ip_selectors(self) -> None: status, _, create = self.create({**CREATE_BODY, "public_ipv4": "auto", "public_ipv6": "auto"}) self.assertEqual(status, 202) diff --git a/atlas/atlas/doctype/atlas_settings/atlas_settings.json b/atlas/atlas/doctype/atlas_settings/atlas_settings.json index 6826c6ccd..1269e6004 100644 --- a/atlas/atlas/doctype/atlas_settings/atlas_settings.json +++ b/atlas/atlas/doctype/atlas_settings/atlas_settings.json @@ -116,6 +116,7 @@ "vm_scheduler_tab", "placement_strategy_section", "placement_strategy", + "affinity_matching", "column_break_dxsw", "use_dedicated_sleepy_vm_hosts", "sleepy_vm_overcommit_factor", @@ -738,6 +739,15 @@ "reqd": 1, "show_description_on_click": 1 }, + { + "default": "Enforced", + "description": "Enforced fails a VM whose affinity rules no host with capacity meets. Preferred places it on any host instead.", + "fieldname": "affinity_matching", + "fieldtype": "Select", + "label": "Affinity Matching", + "options": "Enforced\nPreferred", + "show_description_on_click": 1 + }, { "default": "1.0", "description": "Memory weight for sleepy VM placement. Higher values allow denser placement; 1.0 means no overcommit.", @@ -819,7 +829,7 @@ "index_web_pages_for_search": 1, "issingle": 1, "links": [], - "modified": "2026-09-24 04:43:01.505736", + "modified": "2026-10-01 00:00:00.000000", "modified_by": "Administrator", "module": "Atlas", "name": "Atlas Settings", diff --git a/atlas/atlas/doctype/atlas_settings/atlas_settings.py b/atlas/atlas/doctype/atlas_settings/atlas_settings.py index 13ff6c130..7c0f95179 100644 --- a/atlas/atlas/doctype/atlas_settings/atlas_settings.py +++ b/atlas/atlas/doctype/atlas_settings/atlas_settings.py @@ -50,6 +50,7 @@ class AtlasSettings(Document): if TYPE_CHECKING: from frappe.types import DF + affinity_matching: DF.Literal["Enforced", "Preferred"] atlas_tls_certificate: DF.Password | None atlas_tls_private_key: DF.Password | None auto_spawn_metal_server: DF.Check diff --git a/atlas/metal_server/doctype/metal_server/metal_server.json b/atlas/metal_server/doctype/metal_server/metal_server.json index 1d4595b66..674e51afa 100644 --- a/atlas/metal_server/doctype/metal_server/metal_server.json +++ b/atlas/metal_server/doctype/metal_server/metal_server.json @@ -37,6 +37,8 @@ "wireguard_public_key", "section_break_metadata", "provider_metadata", + "tags_section", + "tags", "hidden_data_section", "column_break_xhax", "metald_tls_certificate", @@ -252,6 +254,19 @@ { "fieldname": "column_break_lnik", "fieldtype": "Column Break" + }, + { + "fieldname": "tags_section", + "fieldtype": "Section Break", + "label": "Tags" + }, + { + "description": "Your own key and value labels for this Metal Server, such as env and local. Use them to find it again.", + "fieldname": "tags", + "fieldtype": "Table", + "label": "Tags", + "options": "Atlas Tag", + "show_description_on_click": 1 } ], "grid_page_length": 50, @@ -273,7 +288,7 @@ "link_fieldname": "server" } ], - "modified": "2026-09-24 01:22:16.439715", + "modified": "2026-10-01 00:00:00.000000", "modified_by": "Administrator", "module": "Metal Server", "name": "Metal Server", diff --git a/atlas/metal_server/doctype/metal_server/metal_server.py b/atlas/metal_server/doctype/metal_server/metal_server.py index fa2ce8782..591724047 100644 --- a/atlas/metal_server/doctype/metal_server/metal_server.py +++ b/atlas/metal_server/doctype/metal_server/metal_server.py @@ -13,6 +13,7 @@ from atlas.atlas.core.background_jobs import run_as_admin from atlas.atlas.core.server_providers.base import ServerCreateRequest, ServerPowerAction +from atlas.atlas.core.tags import validate_tags from atlas.atlas.core.tls.metal import CERTIFICATE_RENEWAL_WINDOW_DAYS, is_certificate_authority_expiring from atlas.atlas.doctype.ssh_task.ssh_task import SSHTask from atlas.metal_server.core.disk_inventory import DiskInventory @@ -42,6 +43,7 @@ class MetalServer(Document): if TYPE_CHECKING: from frappe.types import DF + from atlas.atlas.doctype.atlas_tag.atlas_tag import AtlasTag from atlas.metal_server.doctype.metal_server_disk.metal_server_disk import MetalServerDisk architecture: DF.Literal["amd64", "arm64"] @@ -62,6 +64,7 @@ class MetalServer(Document): server_image: DF.Link server_size: DF.Link status: DF.Literal["Pending", "Installing", "Running", "Stopped", "Failed", "Deleted"] + tags: DF.Table[AtlasTag] title: DF.Data wireguard_ip_address: DF.Data | None wireguard_public_key: DF.Data | None @@ -116,7 +119,8 @@ def ensure_provider_server(self) -> None: self.provider_metadata = frappe.as_json(provider_server.provider_metadata) def validate(self) -> None: - """Fill the mesh address.""" + """Fill the mesh address and check the tags.""" + validate_tags(self) self._set_wireguard_ip_address_if_not_set() def after_insert(self) -> None: diff --git a/atlas/metal_server/doctype/metal_server/test_metal_server.py b/atlas/metal_server/doctype/metal_server/test_metal_server.py index 63ba90e3d..6f87ac3df 100644 --- a/atlas/metal_server/doctype/metal_server/test_metal_server.py +++ b/atlas/metal_server/doctype/metal_server/test_metal_server.py @@ -116,6 +116,24 @@ def test_before_validate_takes_the_architecture_without_provider_creation(self) self.assertEqual(server.architecture, "amd64") provider.ensure_server.assert_not_called() + def test_validate_trims_the_tags_and_fills_the_mesh_address(self) -> None: + tags = [SimpleNamespace(key=" env ", value=" local ")] + server = SimpleNamespace(get=lambda fieldname: tags, _set_wireguard_ip_address_if_not_set=Mock()) + + MetalServer.validate(server) + + self.assertEqual((tags[0].key, tags[0].value), ("env", "local")) + server._set_wireguard_ip_address_if_not_set.assert_called_once_with() + + def test_validate_rejects_a_repeated_tag_key(self) -> None: + tags = [SimpleNamespace(key="env", value="local"), SimpleNamespace(key="env", value="dev")] + server = SimpleNamespace(get=lambda fieldname: tags, _set_wireguard_ip_address_if_not_set=Mock()) + + with self.assertRaises(frappe.ValidationError): + MetalServer.validate(server) + + server._set_wireguard_ip_address_if_not_set.assert_not_called() + def test_ensure_provider_server_identifies_the_host_by_name(self) -> None: provider = SimpleNamespace( ensure_server=Mock( diff --git a/atlas/vm/SPEC.md b/atlas/vm/SPEC.md index 847a8d140..4b30157e0 100644 --- a/atlas/vm/SPEC.md +++ b/atlas/vm/SPEC.md @@ -12,6 +12,7 @@ This file identifies code owners and rules that a code change must preserve. The | `core/vm_service.py` | Changes that span an Atlas record and a Metal host. | | `core/reconciliation.py` | Draft and termination checks after an uncertain Metal result. | | `core/placement/` | Host selection, capacity reservations, and placement strategies. | +| `core/placement/affinity.py` | Affinity rule types, their validation, and their stored JSON form. | | `core/vm_migration.py` and `core/vm_resize.py` | Host moves and shape changes. | | `core/metal_client.py` and `core/metal_models.py` | Metal transport, errors, and typed responses. | | `core/vm_image_transfer.py`, `core/vm_image_deletion.py`, and `core/multipart_upload.py` | Image publication and removal. | @@ -25,6 +26,10 @@ This file identifies code owners and rules that a code change must preserve. The - Metal is the authority for current host state. `Virtual Machine State` is a cache for lists and image checks. A failed Metal read must not appear as a guest state. - Placement checks capacity under a MariaDB named lock with READ COMMITTED isolation. Keep the lock until the draft or migration reservation commits. CPU can be oversubscribed. Memory and disk cannot. - A VM copies its image name and architecture at creation. Image references are immutable. A later VM action does not need the image record. +- The draft commits its tags and affinity rules in the same transaction as its host reservation. +- `PlacementContext` owns the affinity rules, so the strategies do not know about them. It removes the hosts that fail the rules from the snapshot, and checks the rules again under the host lock. Keep both checks: the snapshot can be 1 second old. A failed check under the lock counts as contention, so the next attempt reads a fresh snapshot. +- A migration to a destination that an operator names does not apply the affinity rules. +- A `within` rule reads the VMs of every host in the group of the candidate host. Placement locks each of those hosts without waiting, and keeps the locks until the draft commits. - Atlas changes one network value, then sends the complete network object to Metal. Public address requests own their corresponding default routes. - Guest-specific keys, metadata, and mesh addresses go through Metal and guest metadata. Do not bake them into a shared image. - A protected VM cannot be terminated. An unprotected Atlas record is deleted only after Metal confirms that the VM is absent. @@ -33,6 +38,7 @@ This file identifies code owners and rules that a code change must preserve. The ## Add or change behavior - For a new host ranking rule, add a `PlacementStrategy` subclass under `core/placement/strategies/` and register it. Keep capacity checks in the shared placement path. Read [host selection](../../docs/compute/placement.md). +- For an affinity rule type, change `core/placement/affinity.py` and the matching payload models in `api/models.py` together, then regenerate the API client. Read [affinity rules](../../docs/compute/placement.md#affinity-rules). - For create, resize, move, or delete rules, start in the owning service above. Read [VM lifecycle](../../docs/compute/index.md) and [migration](../../docs/compute/migration.md). - For image transfer or deletion, start in the image owner above. Read [image records](../../docs/storage/image-records.md). - For guest network values and public addresses, read [VM records](../../docs/compute/vm-records.md), [host networking](../../docs/networking/host-networking.md), and [public IPs](../../docs/networking/public-ips.md). diff --git a/atlas/vm/core/models.py b/atlas/vm/core/models.py index a43b5f2e8..b7336f27c 100644 --- a/atlas/vm/core/models.py +++ b/atlas/vm/core/models.py @@ -10,6 +10,7 @@ from atlas.atlas.core.mesh_address import MESH_NETWORK from atlas.atlas.core.parsing import strict_bool +from atlas.vm.core.placement.affinity import AffinityRules ROUTE_VIA_HOST = "host" IPV4_INTERNET_DESTINATION = "0.0.0.0/0" @@ -258,6 +259,8 @@ class VirtualMachineCreateRequest: public_ipv6: str | None = None metadata: dict[str, str] = field(default_factory=dict) firewall: FirewallConfiguration = field(default_factory=FirewallConfiguration) + tags: dict[str, str] = field(default_factory=dict) + placement_rules: AffinityRules = field(default_factory=AffinityRules) @classmethod def from_value(cls, value: str | dict[str, Any]) -> VirtualMachineCreateRequest: @@ -315,6 +318,8 @@ def from_value(cls, value: str | dict[str, Any]) -> VirtualMachineCreateRequest: public_ipv6=public_ipv6, metadata=cls.metadata_map(payload), firewall=FirewallConfiguration.from_value(payload.get("firewall")), + tags=cls.tag_map(payload), + placement_rules=AffinityRules.from_value(payload.get("placement_rules")), ) @staticmethod @@ -336,6 +341,17 @@ def non_negative_integer(payload: dict[str, Any], field_name: str) -> int: raise ValueError(f"{field_name} must be a non-negative integer.") return value + @staticmethod + def tag_map(payload: dict[str, Any]) -> dict[str, str]: + """Return the VM tags. The Virtual Machine record trims them and applies the tag limits.""" + value = payload.get("tags") or {} + if not isinstance(value, dict) or any( + not isinstance(key, str) or not isinstance(item, str) for key, item in value.items() + ): + raise ValueError("Tags must be a string-to-string map.") + + return dict(value) + @staticmethod def metadata_map(payload: dict[str, Any]) -> dict[str, str]: """Return validated metadata with clean keys.""" diff --git a/atlas/vm/core/placement/__init__.py b/atlas/vm/core/placement/__init__.py index 44b27979a..d9fafcdef 100644 --- a/atlas/vm/core/placement/__init__.py +++ b/atlas/vm/core/placement/__init__.py @@ -1,5 +1,11 @@ """Public host placement API.""" +from atlas.vm.core.placement.affinity import ( + AffinityGroup, + AffinityRule, + AffinityRules, + AffinityUnsatisfied, +) from atlas.vm.core.placement.models import CurrentPlacement, PlacementRequirements from atlas.vm.core.placement.strategies.base import ( OutOfCapacity, @@ -9,6 +15,10 @@ ) __all__ = [ + "AffinityGroup", + "AffinityRule", + "AffinityRules", + "AffinityUnsatisfied", "CurrentPlacement", "OutOfCapacity", "PlacementBusy", diff --git a/atlas/vm/core/placement/affinity.py b/atlas/vm/core/placement/affinity.py new file mode 100644 index 000000000..1224f418d --- /dev/null +++ b/atlas/vm/core/placement/affinity.py @@ -0,0 +1,453 @@ +"""Affinity rules that limit the Metal Servers that can hold a virtual machine.""" + +from __future__ import annotations + +import json +from collections.abc import Iterator, Sequence +from dataclasses import dataclass, field +from typing import Literal, assert_never, cast, get_args + +import frappe + +from atlas.atlas.core.exceptions import AtlasUserError +from atlas.atlas.core.tags import ( + MAXIMUM_TAG_KEY_LENGTH, + MAXIMUM_TAG_VALUE_LENGTH, + MAXIMUM_TAGS, + read_tags_for, +) + +AffinityResource = Literal["metal_server", "virtual_machine"] +AffinityOperator = Literal["has", "has_not"] +AffinityGroupKind = Literal["any_of", "all_of"] +AffinityMatching = Literal["Enforced", "Preferred"] + +AFFINITY_RULE_FIELDS = ("resource", "operator", "tags") +OPTIONAL_AFFINITY_RULE_FIELDS = ("within",) +MAXIMUM_AFFINITY_RULES = 16 +MAXIMUM_AFFINITY_GROUP_DEPTH = 5 +# A migration in these states reserves its destination host, as in the capacity query. +RESERVED_MIGRATION_STATUSES = ("preparing", "copying", "cutting_over", "starting", "finalizing", "canceling") + + +class AffinityUnsatisfied(AtlasUserError): + """No Metal Server that meets the affinity rules can hold the requested VM.""" + + code = "affinity_unsatisfied" + http_status_code = 503 + + +def load_affinity_matching() -> AffinityMatching: + """Read how strictly placement applies affinity rules. An unsaved setting uses its default.""" + value = frappe.get_cached_value("Atlas Settings", "Atlas Settings", "affinity_matching") + if not value: + value = frappe.get_meta("Atlas Settings").get_field("affinity_matching").default + if value not in get_args(AffinityMatching): + raise ValueError(f"Unknown affinity matching setting: {value}.") + + return cast(AffinityMatching, value) + + +@dataclass(frozen=True, slots=True) +class AffinityHost: + """The tags that affinity rules read on one candidate Metal Server.""" + + name: str + tags: dict[str, str] + virtual_machine_tags: tuple[dict[str, str], ...] + # The VM tags on every host that shares this host's value, keyed by `within` tag key. + # A key that this host does not have is missing. + virtual_machine_tags_within: dict[str, tuple[dict[str, str], ...]] = field(default_factory=dict) + + @classmethod + def load_many( + cls, + host_names: Sequence[str], + tenant_id: int, + excluded_virtual_machine: str | None = None, + within_keys: Sequence[str] = (), + ) -> list[AffinityHost]: + """Load each host's tags and the tags of the tenant VMs that it, and its groups, hold. + + A VM counts on its assigned host, also as a draft, and on the destination of its active + migration. A VM that is being terminated and `excluded_virtual_machine` never count. + """ + host_tags = read_tags_for("Metal Server", list(host_names)) + groups = {key: find_hosts_by_tag_value(key) for key in within_keys} + hosts_to_load = set(host_names) + for hosts_by_value in groups.values(): + for members in hosts_by_value.values(): + hosts_to_load.update(members) + tags_by_host = cls._load_virtual_machine_tags( + sorted(hosts_to_load), tenant_id, excluded_virtual_machine + ) + + hosts = [] + for name in host_names: + virtual_machine_tags_within = {} + for key, hosts_by_value in groups.items(): + if key not in host_tags[name]: + continue + group_tags = [] + # A tag that changed between the two reads leaves the host in a group of its own. + for member in hosts_by_value.get(host_tags[name][key], [name]): + group_tags.extend(tags_by_host[member]) + virtual_machine_tags_within[key] = tuple(group_tags) + hosts.append(cls(name, host_tags[name], tuple(tags_by_host[name]), virtual_machine_tags_within)) + + return hosts + + @classmethod + def _load_virtual_machine_tags( + cls, host_names: Sequence[str], tenant_id: int, excluded_virtual_machine: str | None + ) -> dict[str, list[dict[str, str]]]: + """Return the tags of each tenant VM, grouped by the host that counts it.""" + placements = cls._find_virtual_machine_hosts(host_names, tenant_id, excluded_virtual_machine) + virtual_machine_tags = read_tags_for("Virtual Machine", sorted({name for name, _ in placements})) + + tags_by_host: dict[str, list[dict[str, str]]] = {name: [] for name in host_names} + for virtual_machine, host_name in placements: + tags_by_host[host_name].append(virtual_machine_tags[virtual_machine]) + + return tags_by_host + + @staticmethod + def _find_virtual_machine_hosts( + host_names: Sequence[str], tenant_id: int, excluded_virtual_machine: str | None + ) -> list[tuple[str, str]]: + """Return a (VM, host) pair for each assigned host and each active migration destination.""" + if not host_names: + return [] + + virtual_machine = frappe.qb.DocType("Virtual Machine") + migration = frappe.qb.DocType("Virtual Machine Migration") + assigned = ( + frappe.qb.from_(virtual_machine) + .select(virtual_machine.name, virtual_machine.server) + .where( + virtual_machine.server.isin(list(host_names)) + & (virtual_machine.tenant_id == tenant_id) + & (virtual_machine.is_terminating == 0) + ) + ) + incoming = ( + frappe.qb.from_(migration) + .join(virtual_machine) + .on(virtual_machine.name == migration.virtual_machine) + .select(virtual_machine.name, migration.destination_metal_server) + .where( + migration.destination_metal_server.isin(list(host_names)) + & migration.status.isin(RESERVED_MIGRATION_STATUSES) + & (virtual_machine.tenant_id == tenant_id) + & (virtual_machine.is_terminating == 0) + ) + ) + if excluded_virtual_machine: + assigned = assigned.where(virtual_machine.name != excluded_virtual_machine) + incoming = incoming.where(virtual_machine.name != excluded_virtual_machine) + + return [*assigned.run(), *incoming.run()] + + +@dataclass(frozen=True, slots=True) +class AffinityRule: + """Require or reject tags on the candidate Metal Server, or on one VM that runs on it.""" + + resource: AffinityResource + operator: AffinityOperator + tags: dict[str, str] + # A host tag key. The rule then reads the VMs on every host with the candidate's value. + within: str | None = None + + def __post_init__(self) -> None: + """Reject a rule whose fields do not match their types.""" + if self.resource not in get_args(AffinityResource): + raise ValueError( + f"Affinity rule resource must be one of: {', '.join(get_args(AffinityResource))}." + ) + if self.operator not in get_args(AffinityOperator): + raise ValueError( + f"Affinity rule operator must be one of: {', '.join(get_args(AffinityOperator))}." + ) + if not isinstance(self.tags, dict) or any( + not isinstance(key, str) or not isinstance(value, str) for key, value in self.tags.items() + ): + raise TypeError("Affinity rule tags must be a dict of strings to strings.") + if not self.tags: + raise ValueError("Affinity rule tags must be an object with at least one key.") + if self.within is not None: + self._check_within() + + def _check_within(self) -> None: + if self.resource != "virtual_machine": + raise ValueError("Only a virtual_machine affinity rule takes within.") + if not isinstance(self.within, str): + raise TypeError("Affinity rule within must be a host tag key.") + if not self.within or len(self.within) > MAXIMUM_TAG_KEY_LENGTH: + raise ValueError(f"Affinity rule within takes 1 to {MAXIMUM_TAG_KEY_LENGTH} characters.") + + @classmethod + def from_value(cls, value: dict[str, object]) -> AffinityRule: + """Parse one rule. Every tag pair must be on the same resource.""" + unknown = set(value) - set(AFFINITY_RULE_FIELDS) - set(OPTIONAL_AFFINITY_RULE_FIELDS) + if unknown: + raise ValueError(f"Unknown affinity rule field: {sorted(unknown)[0]}.") + missing = [name for name in AFFINITY_RULE_FIELDS if name not in value] + if missing: + raise ValueError(f"An affinity rule needs {missing[0]}.") + within = value.get("within") + if within is not None and not isinstance(within, str): + raise ValueError("Affinity rule within must be a host tag key.") + + # The constructor checks the resource, the operator, and within. + return cls( + resource=cast(AffinityResource, value["resource"]), + operator=cast(AffinityOperator, value["operator"]), + tags=cls._parse_tags(value["tags"]), + within=within.strip() if within is not None else None, + ) + + @classmethod + def _parse_tags(cls, value: object) -> dict[str, str]: + """Trim the tag pairs and apply the Atlas Tag limits.""" + if not isinstance(value, dict) or not value: + raise ValueError("Affinity rule tags must be an object with at least one key.") + if len(value) > MAXIMUM_TAGS: + raise ValueError(f"An affinity rule takes at most {MAXIMUM_TAGS} tags.") + + tags: dict[str, str] = {} + for raw_key, raw_value in value.items(): + key, tag_value = cls._parse_tag(raw_key, raw_value) + if key in tags: + raise ValueError(f"Affinity rule tag key {key} is repeated.") + tags[key] = tag_value + + return tags + + @staticmethod + def _parse_tag(key: object, value: object) -> tuple[str, str]: + if not isinstance(key, str) or not isinstance(value, str): + raise ValueError("Affinity rule tag keys and values must be strings.") + + key, value = key.strip(), value.strip() + if not key: + raise ValueError("An affinity rule tag needs a key.") + if len(key) > MAXIMUM_TAG_KEY_LENGTH: + raise ValueError(f"An affinity rule tag key takes at most {MAXIMUM_TAG_KEY_LENGTH} characters.") + if len(value) > MAXIMUM_TAG_VALUE_LENGTH: + raise ValueError( + f"The value of affinity rule tag key {key} takes at most {MAXIMUM_TAG_VALUE_LENGTH} characters." + ) + + return key, value + + def is_satisfied_by(self, host: AffinityHost) -> bool: + """Return whether the host meets this rule. Every pair must be on the same resource.""" + if self.resource == "metal_server": + has_tags = self._is_found_in(host.tags) + elif self.resource == "virtual_machine": + if self.within is not None and self.within not in host.virtual_machine_tags_within: + # A host without the `within` key cannot show which group it is in. + return False + has_tags = self._is_found_on_a_virtual_machine(host) + else: + assert_never(self.resource) + + if self.operator == "has": + return has_tags + elif self.operator == "has_not": + return not has_tags + else: + assert_never(self.operator) + + def _is_found_on_a_virtual_machine(self, host: AffinityHost) -> bool: + """Return whether one VM on the host, or in its `within` group, has every pair.""" + if self.within is None: + virtual_machine_tags = host.virtual_machine_tags + else: + virtual_machine_tags = host.virtual_machine_tags_within[self.within] + + for tags in virtual_machine_tags: + if self._is_found_in(tags): + return True + + return False + + def _is_found_in(self, tags: dict[str, str]) -> bool: + """Return whether `tags` holds every key of this rule, with the same value.""" + for key in self.tags: + if not (key in tags and self.tags[key] == tags[key]): + return False + + return True + + def as_dict(self) -> dict[str, object]: + """Return the JSON-compatible rule.""" + value: dict[str, object] = { + "resource": self.resource, + "operator": self.operator, + "tags": dict(self.tags), + } + if self.within is not None: + value["within"] = self.within + return value + + +@dataclass(frozen=True, slots=True) +class AffinityGroup: + """Combine nodes. `any_of` holds when one node holds, and `all_of` when every node holds.""" + + kind: AffinityGroupKind + nodes: tuple[AffinityNode, ...] + + def __post_init__(self) -> None: + """Reject a group whose kind or nodes do not match their types.""" + if self.kind not in get_args(AffinityGroupKind): + raise ValueError(f"Affinity group kind must be one of: {', '.join(get_args(AffinityGroupKind))}.") + if not isinstance(self.nodes, tuple) or any( + not isinstance(node, AffinityNode) for node in self.nodes + ): + raise TypeError("An affinity group holds a tuple of rules and groups.") + if not self.nodes: + raise ValueError(f"Affinity group {self.kind} must be a list with at least one rule or group.") + + @classmethod + def from_value(cls, kind: AffinityGroupKind, value: object, depth: int) -> AffinityGroup: + """Parse one group. `depth` is 1 for a top-level group.""" + if depth > MAXIMUM_AFFINITY_GROUP_DEPTH: + raise ValueError(f"Affinity groups nest at most {MAXIMUM_AFFINITY_GROUP_DEPTH} deep.") + if not isinstance(value, list): + raise ValueError(f"Affinity group {kind} must be a list with at least one rule or group.") + + # The constructor checks the kind and rejects an empty group. + return cls(kind=kind, nodes=tuple(_parse_node(node, depth) for node in value)) + + def is_satisfied_by(self, host: AffinityHost) -> bool: + """Return whether the host meets one node (`any_of`) or every node (`all_of`).""" + results = (node.is_satisfied_by(host) for node in self.nodes) + if self.kind == "any_of": + return any(results) + elif self.kind == "all_of": + return all(results) + else: + assert_never(self.kind) + + def as_dict(self) -> dict[str, object]: + """Return the JSON-compatible group.""" + return {self.kind: [node.as_dict() for node in self.nodes]} + + +AffinityNode = AffinityRule | AffinityGroup + + +@dataclass(frozen=True, slots=True) +class AffinityRules: + """The affinity rules of one virtual machine. Every top-level node must hold.""" + + nodes: tuple[AffinityNode, ...] = () + + def __post_init__(self) -> None: + """Reject top-level nodes that are not rules or groups.""" + if not isinstance(self.nodes, tuple) or any( + not isinstance(node, AffinityNode) for node in self.nodes + ): + raise TypeError("Affinity rules hold a tuple of rules and groups.") + + @classmethod + def from_value(cls, value: object | None) -> AffinityRules: + """Parse and validate a JSON rule list. None means no rules.""" + if value is None: + return cls() + if not isinstance(value, list): + raise ValueError("Affinity rules must be a list.") + + nodes = tuple(_parse_node(node, depth=0) for node in value) + if len(list(_iter_rules(nodes))) > MAXIMUM_AFFINITY_RULES: + raise ValueError(f"A virtual machine takes at most {MAXIMUM_AFFINITY_RULES} affinity rules.") + + return cls(nodes=nodes) + + @classmethod + def from_json(cls, value: str | None) -> AffinityRules: + """Parse the rules that a Virtual Machine stores. An empty value means no rules.""" + return cls.from_value(json.loads(value) if value else None) + + def filter_hosts( + self, host_names: Sequence[str], tenant_id: int, excluded_virtual_machine: str | None = None + ) -> list[str]: + """Return the hosts that meet every rule, in the given order. + + VM rules read only the VMs of `tenant_id`. `excluded_virtual_machine` is the VM that + placement moves, so it never matches its own rules. + """ + if not self.nodes: + return list(host_names) + + hosts = AffinityHost.load_many(host_names, tenant_id, excluded_virtual_machine, self.within_keys) + return [host.name for host in hosts if self.is_satisfied_by(host)] + + @property + def within_keys(self) -> tuple[str, ...]: + """Return the host tag keys that the `within` rules read, sorted.""" + keys = {rule.within for rule in _iter_rules(self.nodes) if rule.within is not None} + return tuple(sorted(keys)) + + def find_related_hosts(self, host_name: str) -> list[str]: + """Return the other hosts whose VMs the `within` rules read for this host, sorted.""" + within_keys = self.within_keys + if not within_keys: + return [] + + host_tags = read_tags_for("Metal Server", [host_name])[host_name] + related: set[str] = set() + for key in within_keys: + if key in host_tags: + related.update(find_hosts_by_tag_value(key).get(host_tags[key], [])) + + related.discard(host_name) + return sorted(related) + + def is_satisfied_by(self, host: AffinityHost) -> bool: + """Return whether the host meets every top-level node.""" + return all(node.is_satisfied_by(host) for node in self.nodes) + + def as_list(self) -> list[dict[str, object]]: + """Return the JSON-compatible rule list.""" + return [node.as_dict() for node in self.nodes] + + +def _parse_node(value: object, depth: int) -> AffinityNode: + """Parse one rule or group. `depth` counts the groups that enclose the node.""" + if not isinstance(value, dict) or any(not isinstance(key, str) for key in value): + raise ValueError("An affinity rule or group must be an object.") + + kinds = [kind for kind in get_args(AffinityGroupKind) if kind in value] + if not kinds: + return AffinityRule.from_value(value) + if len(value) != 1: + raise ValueError("An affinity group takes exactly one key: any_of or all_of.") + + return AffinityGroup.from_value(kinds[0], value[kinds[0]], depth + 1) + + +def _iter_rules(nodes: tuple[AffinityNode, ...]) -> Iterator[AffinityRule]: + for node in nodes: + if isinstance(node, AffinityRule): + yield node + else: + yield from _iter_rules(node.nodes) + + +def find_hosts_by_tag_value(key: str) -> dict[str, list[str]]: + """Return the Metal Servers for each value of one host tag key, whatever their status.""" + rows = frappe.get_all( + "Atlas Tag", + filters={"parenttype": "Metal Server", "parentfield": "tags", "key": key}, + fields=["parent", "value"], + order_by="parent", + ) + hosts_by_value: dict[str, list[str]] = {} + for row in rows: + hosts_by_value.setdefault(row.value or "", []).append(row.parent) + + return hosts_by_value diff --git a/atlas/vm/core/placement/context.py b/atlas/vm/core/placement/context.py index d45a636da..ee34185a9 100644 --- a/atlas/vm/core/placement/context.py +++ b/atlas/vm/core/placement/context.py @@ -4,10 +4,13 @@ from collections.abc import Iterable from datetime import timedelta from hashlib import sha256 +from typing import assert_never import frappe +from frappe import _ from frappe.utils import now_datetime +from atlas.vm.core.placement.affinity import AffinityUnsatisfied, load_affinity_matching from atlas.vm.core.placement.models import ( CurrentPlacement, FleetUsage, @@ -81,6 +84,8 @@ def __init__( self.current_host_name = current_placement.host_name if current_placement else None self.sleepy_vm_overcommit_factor = sleepy_vm_overcommit_factor self.use_dedicated_sleepy_vm_hosts = use_dedicated_sleepy_vm_hosts + # The affinity filter sets this when it removes hosts from the snapshot. + self.apply_affinity = False self.has_contended_hosts = False self.probe_count = 0 self.last_probe_was_contended = False @@ -127,20 +132,59 @@ def try_select(self, host_name: str, *, wait: bool = False) -> bool: self.last_probe_was_contended = True return False + held_locks = [lock_name] try: - has_capacity = self._host_has_capacity(host_name) + is_selectable = self._check_locked_host(host_name, held_locks) except Exception: - self._release_host_lock(lock_name) + self._release_host_locks(held_locks) raise - if not has_capacity: - self._release_host_lock(lock_name) + if not is_selectable: + self._release_host_locks(held_locks) return False - self._keep_host_lock_for_transaction(lock_name) + for held_lock in held_locks: + self._keep_host_lock_for_transaction(held_lock) self._selected_host = host_name return True + def _check_locked_host(self, host_name: str, held_locks: list[str]) -> bool: + """Check capacity and affinity while holding the host lock. + + A `within` rule reads the VMs of every host in the group of this host. Each of those + hosts is locked too, so no VM lands there before this placement commits. Each extra + lock is added to `held_locks`. + """ + if not self._host_has_capacity(host_name): + return False + + for member in self._find_affinity_group(host_name): + member_lock = self._host_lock_name(member) + # Do not wait, so two placements that lock the same group cannot deadlock. + if not self._acquire_host_lock(member_lock, wait=False): + self.has_contended_hosts = True + return False + held_locks.append(member_lock) + + if not self._host_meets_affinity(host_name): + # The snapshot missed a VM that a concurrent placement committed. Treat the + # host as contended, so the next attempt reads a fresh snapshot. + self.has_contended_hosts = True + return False + + return True + + def _find_affinity_group(self, host_name: str) -> list[str]: + """Return the other hosts whose VMs the `within` rules read for this host.""" + if not self.apply_affinity or host_name == self.current_host_name: + return [] + + return self.requirements.placement_rules.find_related_hosts(host_name) + + def _release_host_locks(self, lock_names: list[str]) -> None: + for lock_name in lock_names: + self._release_host_lock(lock_name) + @staticmethod def _host_lock_name(host_name: str) -> str: """Return a server-global lock name scoped to this site database.""" @@ -282,10 +326,11 @@ def _find_snapshot_host(self, host_name: str) -> HostUsage | None: return next((host for host in self.usage.hosts if host.name == host_name), None) def _load_fleet_usage(self, cache_snapshot: bool) -> FleetUsage: - """Build a fleet view from cached or current rows.""" - if not cache_snapshot: - return self._build_fleet_usage(self._load_snapshot_rows()) + """Build a fleet view from cached or current rows, without hosts that fail the affinity rules.""" + rows = self._load_cached_snapshot_rows() if cache_snapshot else self._load_snapshot_rows() + return self._build_fleet_usage(self._filter_by_affinity(rows)) + def _load_cached_snapshot_rows(self) -> list[frappe._dict]: cache = frappe.cache() pool = int(self.requirements.is_sleepy) if self.use_dedicated_sleepy_vm_hosts else "any" key = "atlas:placement-snapshot:{0}:{1}:{2}".format( @@ -296,7 +341,57 @@ def _load_fleet_usage(self, cache_snapshot: bool) -> FleetUsage: rows = self._load_snapshot_rows() cache.set_value(key, rows, expires_in_sec=SNAPSHOT_CACHE_SECONDS) - return self._build_fleet_usage(rows) + return rows + + def _filter_by_affinity(self, rows: list[frappe._dict]) -> list[frappe._dict]: + """Keep only the hosts that meet the affinity rules, when one of them has room. + + Without such a host, `Enforced` matching fails and `Preferred` matching keeps every + host. The current host of a resize always stays, so an in-place resize ignores the rules. + """ + if not self.requirements.placement_rules.nodes: + return rows + + allowed = set(self._find_affinity_hosts([row.name for row in rows])) + matching = [row for row in rows if row.name in allowed or row.name == self.current_host_name] + if self.current_host_name or any(self._has_room(row) for row in matching): + self.apply_affinity = True + return matching + + affinity_matching = load_affinity_matching() + if affinity_matching == "Enforced": + frappe.throw( + _("No Metal Server that meets the affinity rules has capacity."), exc=AffinityUnsatisfied + ) + elif affinity_matching == "Preferred": + return rows + else: + assert_never(affinity_matching) + + def _has_room(self, row: frappe._dict) -> bool: + """Return whether a snapshot row reports room for the request, as a strategy checks it.""" + return ( + row.name not in self._excluded_servers + and row.free_memory_mib >= self.requirements.memory_mib + and row.free_storage_mib >= self.requirements.disk_mib + ) + + def _host_meets_affinity(self, host_name: str) -> bool: + """Check the affinity rules again under the host lock, against committed VMs. + + Concurrent placements can share one snapshot. The host lock orders them, so the + second placement reads the committed draft of the first one here. + """ + if not self.apply_affinity or host_name == self.current_host_name: + return True + + return bool(self._find_affinity_hosts([host_name])) + + def _find_affinity_hosts(self, host_names: list[str]) -> list[str]: + requirements = self.requirements + return requirements.placement_rules.filter_hosts( + host_names, requirements.tenant_id, requirements.virtual_machine + ) def _load_snapshot_rows(self) -> list[frappe._dict]: """Load current capacity and rank data for the required host pool.""" diff --git a/atlas/vm/core/placement/models.py b/atlas/vm/core/placement/models.py index 4ef387378..b7c1aa8e5 100644 --- a/atlas/vm/core/placement/models.py +++ b/atlas/vm/core/placement/models.py @@ -1,8 +1,10 @@ from __future__ import annotations -from dataclasses import dataclass +from dataclasses import dataclass, field from datetime import datetime +from atlas.vm.core.placement.affinity import AffinityRules + @dataclass(frozen=True, slots=True) class Resources: @@ -30,6 +32,9 @@ class PlacementRequirements: architecture: str tenant_id: int is_sleepy: bool + placement_rules: AffinityRules = field(default_factory=AffinityRules) + # The VM that placement moves. It never matches its own affinity rules. + virtual_machine: str | None = None @dataclass(frozen=True, slots=True) diff --git a/atlas/vm/core/placement/test_affinity.py b/atlas/vm/core/placement/test_affinity.py new file mode 100644 index 000000000..578461dec --- /dev/null +++ b/atlas/vm/core/placement/test_affinity.py @@ -0,0 +1,376 @@ +import re +from unittest.mock import patch +from uuid import uuid7 + +import frappe +from frappe.tests import IntegrationTestCase, UnitTestCase + +from atlas.atlas.core.tags import MAXIMUM_TAG_KEY_LENGTH, MAXIMUM_TAG_VALUE_LENGTH, MAXIMUM_TAGS +from atlas.vm.core.placement.affinity import ( + MAXIMUM_AFFINITY_GROUP_DEPTH, + MAXIMUM_AFFINITY_RULES, + AffinityGroup, + AffinityHost, + AffinityRule, + AffinityRules, + load_affinity_matching, +) + +STORAGE_HOST = {"resource": "metal_server", "operator": "has", "tags": {"type": "storage-optimised"}} +RACK_A_HOST = {"resource": "metal_server", "operator": "has", "tags": {"rack": "a"}} +NO_CARGO_SERVER = {"resource": "virtual_machine", "operator": "has_not", "tags": {"role": "cargo-server"}} + + +def nest(node: dict, depth: int) -> dict: + """Wrap one node in `depth` groups that alternate between any_of and all_of.""" + for level in range(depth): + node = {("any_of", "all_of")[level % 2]: [node]} + return node + + +def rule(resource: str, operator: str, **tags: str) -> dict: + return {"resource": resource, "operator": operator, "tags": tags} + + +def insert(doctype: str, **values: object) -> str: + """Insert one row and its child rows without document hooks, and return its name.""" + document = frappe.get_doc({"doctype": doctype, **values}) + document.db_insert() + for row in document.get_all_children(): + row.db_insert() + return document.name + + +class TestAffinityRules(UnitTestCase): + def test_nested_rules_keep_their_shape(self) -> None: + value = [ + {"any_of": [{"all_of": [STORAGE_HOST, RACK_A_HOST]}, STORAGE_HOST]}, + NO_CARGO_SERVER, + ] + + rules = AffinityRules.from_value(value) + + storage = AffinityRule("metal_server", "has", {"type": "storage-optimised"}) + rack_a = AffinityRule("metal_server", "has", {"rack": "a"}) + self.assertEqual( + rules, + AffinityRules( + ( + AffinityGroup("any_of", (AffinityGroup("all_of", (storage, rack_a)), storage)), + AffinityRule("virtual_machine", "has_not", {"role": "cargo-server"}), + ) + ), + ) + self.assertEqual(rules.as_list(), value) + + def test_absent_and_empty_rules_mean_no_rules(self) -> None: + self.assertEqual(AffinityRules.from_value(None), AffinityRules()) + self.assertEqual(AffinityRules.from_value([]), AffinityRules()) + + def test_tag_pairs_are_trimmed(self) -> None: + rules = AffinityRules.from_value([{**RACK_A_HOST, "tags": {" rack ": " a "}}]) + + self.assertEqual(rules.as_list(), [RACK_A_HOST]) + + def test_within_is_trimmed_kept_and_listed(self) -> None: + value = [{"any_of": [{**NO_CARGO_SERVER, "within": " rack "}, {**NO_CARGO_SERVER, "within": "zone"}]}] + + rules = AffinityRules.from_value(value) + + self.assertEqual(rules.as_list()[0]["any_of"][0], {**NO_CARGO_SERVER, "within": "rack"}) + self.assertNotIn("within", AffinityRules.from_value([NO_CARGO_SERVER]).as_list()[0]) + self.assertEqual(rules.within_keys, ("rack", "zone")) + + def test_rules_at_the_limits_are_accepted(self) -> None: + AffinityRules.from_value([NO_CARGO_SERVER] * MAXIMUM_AFFINITY_RULES) + AffinityRules.from_value([nest(NO_CARGO_SERVER, MAXIMUM_AFFINITY_GROUP_DEPTH)]) + + def test_construction_checks_every_field_type(self) -> None: + rule = AffinityRule("metal_server", "has", {"rack": "a"}) + cases = [ + ("unknown resource", ValueError, lambda: AffinityRule("rack", "has", {"rack": "a"})), + ( + "within on a host rule", + ValueError, + lambda: AffinityRule("metal_server", "has", {"rack": "a"}, "rack"), + ), + ( + "within that is not a string", + TypeError, + lambda: AffinityRule("virtual_machine", "has", {"a": "b"}, 5), + ), + ("unknown operator", ValueError, lambda: AffinityRule("metal_server", "in", {"rack": "a"})), + ("empty tags", ValueError, lambda: AffinityRule("metal_server", "has", {})), + ( + "tag value that is not a string", + TypeError, + lambda: AffinityRule("metal_server", "has", {"rack": 1}), + ), + ("unknown group kind", ValueError, lambda: AffinityGroup("none_of", (rule,))), + ("empty group", ValueError, lambda: AffinityGroup("any_of", ())), + ("group nodes in a list", TypeError, lambda: AffinityGroup("any_of", [rule])), + ("group node that is not a rule", TypeError, lambda: AffinityGroup("all_of", ("rule",))), + ("top-level node that is a dict", TypeError, lambda: AffinityRules((STORAGE_HOST,))), + ] + for label, error, build in cases: + with self.subTest(label), self.assertRaises(error): + build() + + def test_invalid_rules_are_rejected(self) -> None: + cases = [ + ("Affinity rules must be a list.", NO_CARGO_SERVER), + ("An affinity rule or group must be an object.", ["has_not"]), + ("exactly one key", [{"any_of": [STORAGE_HOST], "all_of": [STORAGE_HOST]}]), + ("exactly one key", [{"any_of": [STORAGE_HOST], "resource": "metal_server"}]), + # An unknown group key is not a group, so the rule parser rejects it. + ("Unknown affinity rule field: none_of.", [{"none_of": [STORAGE_HOST]}]), + ("at least one rule or group", [{"all_of": []}]), + ("Unknown affinity rule field: weight.", [{**NO_CARGO_SERVER, "weight": 1}]), + ("Only a virtual_machine affinity rule takes within.", [{**RACK_A_HOST, "within": "rack"}]), + ("within takes 1 to", [{**NO_CARGO_SERVER, "within": " "}]), + ("within takes 1 to", [{**NO_CARGO_SERVER, "within": "k" * (MAXIMUM_TAG_KEY_LENGTH + 1)}]), + ("within must be a host tag key", [{**NO_CARGO_SERVER, "within": 5}]), + ("An affinity rule needs operator.", [{"resource": "metal_server", "tags": {"rack": "a"}}]), + ("resource must be one of", [{**STORAGE_HOST, "resource": "metal-server"}]), + ("operator must be one of", [{**STORAGE_HOST, "operator": "in"}]), + ("at least one key", [{**STORAGE_HOST, "tags": {}}]), + ("must be strings", [{**STORAGE_HOST, "tags": {"rack": 1}}]), + ("needs a key", [{**STORAGE_HOST, "tags": {" ": "a"}}]), + ( + "Affinity rule tag key rack is repeated.", + [{**STORAGE_HOST, "tags": {"rack": "a", " rack": "b"}}], + ), + ("key takes at most", [{**STORAGE_HOST, "tags": {"k" * (MAXIMUM_TAG_KEY_LENGTH + 1): "a"}}]), + ( + "value of affinity rule tag key rack", + [{**STORAGE_HOST, "tags": {"rack": "v" * (MAXIMUM_TAG_VALUE_LENGTH + 1)}}], + ), + ( + f"at most {MAXIMUM_TAGS} tags", + [{**STORAGE_HOST, "tags": {f"key-{index}": "a" for index in range(MAXIMUM_TAGS + 1)}}], + ), + ("nest at most", [nest(NO_CARGO_SERVER, MAXIMUM_AFFINITY_GROUP_DEPTH + 1)]), + ("affinity rules", [NO_CARGO_SERVER] * (MAXIMUM_AFFINITY_RULES + 1)), + # Rules inside groups count toward the same limit. + ("affinity rules", [{"any_of": [NO_CARGO_SERVER] * (MAXIMUM_AFFINITY_RULES + 1)}]), + ] + for message, value in cases: + with self.subTest(message=message), self.assertRaisesRegex(ValueError, re.escape(message)): + AffinityRules.from_value(value) + + +class TestAffinityMatching(UnitTestCase): + def test_an_unsaved_setting_uses_its_default(self) -> None: + with patch("atlas.vm.core.placement.affinity.frappe.get_cached_value", return_value=None): + self.assertEqual(load_affinity_matching(), "Enforced") + + def test_an_unknown_setting_is_rejected(self) -> None: + with ( + patch("atlas.vm.core.placement.affinity.frappe.get_cached_value", return_value="Sometimes"), + self.assertRaises(ValueError), + ): + load_affinity_matching() + + +class TestAffinityEvaluation(UnitTestCase): + def test_rules_read_the_host_and_its_virtual_machines(self) -> None: + storage_in_rack_a = AffinityHost("a", {"type": "storage-optimised", "rack": "a"}, ()) + memory_host = AffinityHost("b", {"type": "memory-optimised"}, ()) + split_pairs = AffinityHost("c", {}, ({"role": "db"}, {"env": "prod"})) + one_vm_with_both = AffinityHost("d", {}, ({"role": "db", "env": "prod"},)) + cases = [ + ("a host rule reads the host tags", [STORAGE_HOST], storage_in_rack_a, True), + ("a host rule needs the same value", [STORAGE_HOST], memory_host, False), + ( + "has_not rejects a host with every pair", + [rule("metal_server", "has_not", type="storage-optimised")], + storage_in_rack_a, + False, + ), + ( + "has_not accepts a host without a pair", + [rule("metal_server", "has_not", type="storage-optimised", rack="b")], + storage_in_rack_a, + True, + ), + ( + "a VM rule needs every pair on one VM", + [rule("virtual_machine", "has", role="db", env="prod")], + split_pairs, + False, + ), + ( + "a VM rule accepts one VM with every pair", + [rule("virtual_machine", "has", role="db", env="prod")], + one_vm_with_both, + True, + ), + ( + "all_of accepts pairs on different VMs", + [ + { + "all_of": [ + rule("virtual_machine", "has", role="db"), + rule("virtual_machine", "has", env="prod"), + ] + } + ], + split_pairs, + True, + ), + ("has_not accepts a host without VMs", [NO_CARGO_SERVER], memory_host, True), + ("any_of needs one node", [{"any_of": [STORAGE_HOST, RACK_A_HOST]}], memory_host, False), + ( + "any_of accepts one matching node", + [{"any_of": [rule("metal_server", "has", type="memory-optimised"), RACK_A_HOST]}], + memory_host, + True, + ), + ( + "every top-level node must hold", + [STORAGE_HOST, rule("metal_server", "has", rack="b")], + storage_in_rack_a, + False, + ), + ("no rules accept every host", [], memory_host, True), + ] + for label, value, host, expected in cases: + with self.subTest(label): + self.assertEqual(AffinityRules.from_value(value).is_satisfied_by(host), expected) + + def test_within_reads_every_virtual_machine_in_the_group(self) -> None: + # A replica runs on another host of rack a. This host runs none itself. + rack_a = AffinityHost("a", {"rack": "a"}, (), {"rack": ({"role": "replica"},)}) + rack_b = AffinityHost("b", {"rack": "b"}, (), {"rack": ()}) + no_rack = AffinityHost("c", {}, (), {}) + avoid = { + "resource": "virtual_machine", + "operator": "has_not", + "tags": {"role": "replica"}, + "within": "rack", + } + join = {**avoid, "operator": "has"} + cases = [ + ("has_not rejects a rack with a replica", [avoid], rack_a, False), + ("has_not accepts a rack without one", [avoid], rack_b, True), + ("has accepts a rack with a replica", [join], rack_a, True), + ("has rejects a rack without one", [join], rack_b, False), + ("has_not fails on a host without the key", [avoid], no_rack, False), + ("has fails on a host without the key", [join], no_rack, False), + ("without within, only the host counts", [{**avoid, "within": None}], rack_a, True), + ] + for label, value, host, expected in cases: + with self.subTest(label): + self.assertEqual(AffinityRules.from_value(value).is_satisfied_by(host), expected) + + +class TestAffinityHostFilter(IntegrationTestCase): + def setUp(self) -> None: + super().setUp() + # Metal Server and Virtual Machine Migration names are UUID columns, so they reject other values. + self.storage_host = insert( + "Metal Server", + name=str(uuid7()), + tags=[{"key": "type", "value": "storage-optimised"}, {"key": "rack", "value": self.rack("a")}], + ) + self.memory_host = insert( + "Metal Server", + name=str(uuid7()), + tags=[{"key": "type", "value": "memory-optimised"}, {"key": "rack", "value": self.rack("b")}], + ) + self.plain_host = insert( + "Metal Server", name=str(uuid7()), tags=[{"key": "rack", "value": self.rack("a")}] + ) + self.hosts = [self.storage_host, self.memory_host, self.plain_host] + + self.cargo_server = self.virtual_machine(0, self.storage_host, role="cargo-server") + self.virtual_machine(7, self.memory_host, role="cargo-server") + moving_database = self.virtual_machine(0, self.plain_host, role="db") + self.migration(moving_database, self.memory_host, "copying") + moved_archive = self.virtual_machine(0, self.plain_host, role="archive") + self.migration(moved_archive, self.storage_host, "completed") + + def rack(self, name: str) -> str: + """Return a rack value unique to this test, so other stored hosts never join its racks.""" + if not hasattr(self, "_rack_suffix"): + self._rack_suffix = frappe.generate_hash(length=8) + return f"{name}-{self._rack_suffix}" + + def virtual_machine(self, tenant_id: int, host: str, *, is_terminating: int = 0, **tags: str) -> str: + return insert( + "Virtual Machine", + name=f"test-vm-{frappe.generate_hash(length=8)}", + tenant_id=tenant_id, + server=host, + is_terminating=is_terminating, + tags=[{"key": key, "value": value} for key, value in tags.items()], + ) + + def migration(self, virtual_machine: str, destination: str, status: str) -> None: + insert( + "Virtual Machine Migration", + name=str(uuid7()), + virtual_machine=virtual_machine, + destination_metal_server=destination, + status=status, + ) + + def filter_hosts(self, value: list, excluded_virtual_machine: str | None = None) -> list[str]: + return AffinityRules.from_value(value).filter_hosts(self.hosts, 0, excluded_virtual_machine) + + def test_host_rules_keep_the_given_order(self) -> None: + rules = AffinityRules.from_value( + [{"any_of": [STORAGE_HOST, rule("metal_server", "has", type="memory-optimised")]}] + ) + + self.assertEqual( + rules.filter_hosts([self.memory_host, self.plain_host, self.storage_host], 0), + [self.memory_host, self.storage_host], + ) + + def test_virtual_machine_rules_read_only_the_tenant_virtual_machines(self) -> None: + self.assertEqual(self.filter_hosts([NO_CARGO_SERVER]), [self.memory_host, self.plain_host]) + + def test_the_excluded_virtual_machine_never_matches(self) -> None: + self.assertEqual(self.filter_hosts([NO_CARGO_SERVER], self.cargo_server), self.hosts) + + def test_an_active_migration_counts_on_its_destination(self) -> None: + self.assertEqual( + self.filter_hosts([rule("virtual_machine", "has", role="db")]), + [self.memory_host, self.plain_host], + ) + self.assertEqual( + self.filter_hosts([rule("virtual_machine", "has", role="archive")]), [self.plain_host] + ) + + def test_within_reads_the_virtual_machines_of_the_whole_rack(self) -> None: + # The tenant-0 Cargo Server runs on the storage host in rack a, next to the plain host. + rule = {**NO_CARGO_SERVER, "within": "rack"} + + self.assertEqual(self.filter_hosts([rule]), [self.memory_host]) + + def test_within_on_a_key_that_no_host_has_rejects_every_host(self) -> None: + self.assertEqual(self.filter_hosts([{**NO_CARGO_SERVER, "within": "zone"}]), []) + + def test_the_related_hosts_share_the_rack(self) -> None: + rules = AffinityRules.from_value([{**NO_CARGO_SERVER, "within": "rack"}]) + + self.assertEqual(rules.find_related_hosts(self.storage_host), [self.plain_host]) + self.assertEqual(rules.find_related_hosts(self.memory_host), []) + self.assertEqual( + AffinityRules.from_value([NO_CARGO_SERVER]).find_related_hosts(self.storage_host), [] + ) + + def test_a_terminating_virtual_machine_does_not_count(self) -> None: + self.virtual_machine(0, self.plain_host, is_terminating=1, role="cargo-server") + leaving = self.virtual_machine(0, self.storage_host, is_terminating=1, role="archive") + self.migration(leaving, self.memory_host, "copying") + + self.assertEqual(self.filter_hosts([NO_CARGO_SERVER]), [self.memory_host, self.plain_host]) + self.assertEqual( + self.filter_hosts([rule("virtual_machine", "has", role="archive")]), [self.plain_host] + ) + + def test_no_rules_keep_every_host_without_a_query(self) -> None: + with self.assertQueryCount(0): + self.assertEqual(AffinityRules().filter_hosts(self.hosts, 0), self.hosts) diff --git a/atlas/vm/core/placement/test_context.py b/atlas/vm/core/placement/test_context.py index 05d6dc467..586679deb 100644 --- a/atlas/vm/core/placement/test_context.py +++ b/atlas/vm/core/placement/test_context.py @@ -1,4 +1,5 @@ import time +from contextlib import contextmanager from datetime import datetime, timedelta from types import SimpleNamespace from unittest.mock import patch @@ -7,6 +8,7 @@ import frappe from frappe.tests import IntegrationTestCase, UnitTestCase +from atlas.vm.core.placement.affinity import AffinityRules, AffinityUnsatisfied from atlas.vm.core.placement.context import ( CAPACITY_MAXIMUM_AGE, LOCK_WAIT_MAXIMUM_SECONDS, @@ -71,6 +73,132 @@ def context(self, rows: object, **kwargs: object) -> PlacementContext: with patch("atlas.vm.core.placement.context.frappe.db.sql", side_effect=rows): return PlacementContext(self.requirements(), 1.0, **kwargs) # type: ignore[arg-type] + def affinity_placement(self, current_host_name: str | None = None) -> PlacementContext: + """Return a placement whose rules only host "a" meets.""" + placement = object.__new__(PlacementContext) + rules = AffinityRules.from_value( + [{"resource": "metal_server", "operator": "has", "tags": {"rack": "a"}}] + ) + placement.requirements = self.requirements(placement_rules=rules) + placement.current_host_name = current_host_name + placement.apply_affinity = False + placement._excluded_servers = frozenset() + return placement + + @staticmethod + def affinity_rows(free_memory_of_a: int = 4096) -> list[frappe._dict]: + return [ + frappe._dict(name="a", free_memory_mib=free_memory_of_a, free_storage_mib=20480), + frappe._dict(name="b", free_memory_mib=4096, free_storage_mib=20480), + frappe._dict(name="c", free_memory_mib=0, free_storage_mib=0), + ] + + def filter_by_affinity( + self, placement: PlacementContext, rows: list, affinity_matching: str + ) -> list[str]: + with ( + patch.object(PlacementContext, "_find_affinity_hosts", return_value=["a"]), + patch("atlas.vm.core.placement.context.load_affinity_matching", return_value=affinity_matching), + ): + return [row.name for row in placement._filter_by_affinity(rows)] + + def test_the_affinity_filter_keeps_only_matching_hosts_when_one_has_room(self) -> None: + placement = self.affinity_placement() + + self.assertEqual(self.filter_by_affinity(placement, self.affinity_rows(), "Enforced"), ["a"]) + self.assertTrue(placement.apply_affinity) + + def test_the_affinity_filter_keeps_the_current_host_of_a_resize(self) -> None: + placement = self.affinity_placement(current_host_name="c") + + self.assertEqual(self.filter_by_affinity(placement, self.affinity_rows(0), "Enforced"), ["a", "c"]) + + def test_enforced_matching_fails_when_no_matching_host_has_room(self) -> None: + with self.assertRaises(AffinityUnsatisfied): + self.filter_by_affinity(self.affinity_placement(), self.affinity_rows(0), "Enforced") + + def test_preferred_matching_keeps_every_host_when_no_matching_host_has_room(self) -> None: + placement = self.affinity_placement() + + self.assertEqual( + self.filter_by_affinity(placement, self.affinity_rows(0), "Preferred"), ["a", "b", "c"] + ) + self.assertFalse(placement.apply_affinity) + + def test_a_rule_that_fails_under_the_lock_counts_as_contention(self) -> None: + for has_capacity, is_contended in ((True, True), (False, False)): + placement = object.__new__(PlacementContext) + placement._selected_host = None + placement._excluded_servers = frozenset() + placement.probe_count = 0 + placement.has_contended_hosts = False + placement.last_probe_was_contended = False + with ( + self.subTest(has_capacity=has_capacity), + patch.object(PlacementContext, "_find_snapshot_host", return_value=True), + patch.object(PlacementContext, "_host_lock_name", return_value="lock"), + patch.object(PlacementContext, "_acquire_host_lock", return_value=True), + patch.object(PlacementContext, "_release_host_lock"), + patch.object(PlacementContext, "_host_has_capacity", return_value=has_capacity), + patch.object(PlacementContext, "_find_affinity_group", return_value=[]), + patch.object(PlacementContext, "_host_meets_affinity", return_value=False), + ): + self.assertFalse(placement.try_select("a")) + + self.assertEqual(placement.has_contended_hosts, is_contended) + self.assertFalse(placement.last_probe_was_contended) + + @contextmanager + def locked_group(self, acquired: list[bool], meets_affinity: bool = True): + """Patch the locks of host "a" and of its within group, hosts "b" and "c".""" + with ( + patch.object(PlacementContext, "_find_snapshot_host", return_value=True), + patch.object(PlacementContext, "_host_lock_name", side_effect=lambda name: f"lock-{name}"), + patch.object(PlacementContext, "_acquire_host_lock", side_effect=acquired) as acquire, + patch.object(PlacementContext, "_release_host_lock") as release, + patch.object(PlacementContext, "_keep_host_lock_for_transaction") as keep, + patch.object(PlacementContext, "_host_has_capacity", return_value=True), + patch.object(PlacementContext, "_find_affinity_group", return_value=["b", "c"]), + patch.object(PlacementContext, "_host_meets_affinity", return_value=meets_affinity), + ): + yield acquire, release, keep + + @staticmethod + def unselected_placement() -> PlacementContext: + placement = object.__new__(PlacementContext) + placement._selected_host = None + placement._excluded_servers = frozenset() + placement.probe_count = 0 + placement.has_contended_hosts = False + placement.last_probe_was_contended = False + return placement + + def test_the_within_group_locks_are_kept_until_the_transaction_ends(self) -> None: + placement = self.unselected_placement() + with self.locked_group([True, True, True]) as (acquire, release, keep): + self.assertTrue(placement.try_select("a")) + + self.assertEqual([c.kwargs["wait"] for c in acquire.call_args_list], [False, False, False]) + self.assertEqual([c.args[0] for c in keep.call_args_list], ["lock-a", "lock-b", "lock-c"]) + release.assert_not_called() + + def test_a_busy_host_in_the_within_group_counts_as_contention(self) -> None: + placement = self.unselected_placement() + with self.locked_group([True, True, False]) as (_acquire, release, keep): + self.assertFalse(placement.try_select("a")) + + self.assertTrue(placement.has_contended_hosts) + self.assertEqual([c.args[0] for c in release.call_args_list], ["lock-a", "lock-b"]) + keep.assert_not_called() + + def test_a_rule_that_fails_under_the_group_locks_releases_them_all(self) -> None: + placement = self.unselected_placement() + with self.locked_group([True, True, True], meets_affinity=False) as (_acquire, release, _keep): + self.assertFalse(placement.try_select("a")) + + self.assertTrue(placement.has_contended_hosts) + self.assertEqual([c.args[0] for c in release.call_args_list], ["lock-a", "lock-b", "lock-c"]) + def test_a_cached_snapshot_is_shared_between_placements(self) -> None: rows = [host_row("a")] store: dict[str, object] = {} @@ -475,9 +603,11 @@ def placement( placement = object.__new__(PlacementContext) placement.requirements = TestPlacementContext.requirements(**overrides) placement.current_placement = current_placement + placement.current_host_name = current_placement.host_name if current_placement else None placement._created_at = NOW placement._deadline = None placement._excluded_servers = frozenset() + placement.apply_affinity = False placement._selected_host = None placement.has_contended_hosts = False placement.last_probe_was_contended = False @@ -490,6 +620,31 @@ def test_capacity_query_accepts_a_fitting_host(self) -> None: with self.primary_connection(): self.assertTrue(self.placement()._host_has_capacity(self.host_name)) + def test_the_locked_affinity_check_reads_a_vm_that_breaks_a_rule(self) -> None: + rules = AffinityRules.from_value( + [{"resource": "virtual_machine", "operator": "has_not", "tags": {"role": "cargo-server"}}] + ) + with self.primary_connection(): + placement = self.placement(placement_rules=rules) + placement.apply_affinity = True + self.assertTrue(placement._host_meets_affinity(self.host_name)) + + # Another placement committed this VM after the snapshot was taken. + virtual_machine = frappe.get_doc( + { + "doctype": "Virtual Machine", + "name": f"test-vm-{frappe.generate_hash(length=8)}", + "tenant_id": 7, + "server": self.host_name, + "tags": [{"key": "role", "value": "cargo-server"}], + } + ) + virtual_machine.db_insert() + for row in virtual_machine.get_all_children(): + row.db_insert() + + self.assertFalse(placement._host_meets_affinity(self.host_name)) + def test_capacity_query_rejects_a_host_without_capacity(self) -> None: with self.primary_connection(): self.assertFalse(self.placement(memory_mib=8192)._host_has_capacity(self.host_name)) diff --git a/atlas/vm/core/test_vm_migration.py b/atlas/vm/core/test_vm_migration.py index f86d44154..dd3f9c29a 100644 --- a/atlas/vm/core/test_vm_migration.py +++ b/atlas/vm/core/test_vm_migration.py @@ -45,6 +45,7 @@ def source_vm(**overrides: object) -> SimpleNamespace: "memory_mib": 512, "disk_mib": 1024, "sleep_after_idle_seconds": 0, + "placement_rules": None, } values.update(overrides) return SimpleNamespace(**values) @@ -92,7 +93,48 @@ def test_aborted_destination_becomes_failed_after_a_metal_error(self) -> None: service.settle.assert_called_once_with("failed") +AFFINITY_RULES = '[{"resource": "virtual_machine", "operator": "has_not", "tags": {"role": "db"}}]' + + class TestDestinationMetalServerSelection(UnitTestCase): + def test_a_chosen_destination_follows_the_placement_rules(self) -> None: + service = MigrationService(migration_doc(status="scheduled", destination_metal_server=None)) + service.update = Mock() + + with ( + patch( + "atlas.vm.core.vm_migration.frappe.get_doc", + return_value=source_vm(placement_rules=AFFINITY_RULES), + ), + patch("atlas.vm.core.vm_migration.now_datetime", return_value="2026-09-21 12:00:00"), + patch( + "atlas.vm.core.vm_migration.PlacementStrategy.find_server", return_value="metal-2" + ) as find_server, + ): + self.assertTrue(service.select_destination_metal_server()) + + requirements = find_server.call_args.args[0] + self.assertEqual(requirements.placement_rules.as_list()[0]["tags"], {"role": "db"}) + self.assertEqual(requirements.virtual_machine, "vm-00001") + + def test_a_named_destination_ignores_the_placement_rules(self) -> None: + service = MigrationService(migration_doc(status="scheduled", destination_metal_server="metal-3")) + service.update = Mock() + + with ( + patch( + "atlas.vm.core.vm_migration.frappe.get_doc", + return_value=source_vm(placement_rules=AFFINITY_RULES), + ), + patch("atlas.vm.core.vm_migration.now_datetime", return_value="2026-09-21 12:00:00"), + patch( + "atlas.vm.core.vm_migration.PlacementStrategy.reserve_server", return_value="metal-3" + ) as reserve_server, + ): + self.assertTrue(service.select_destination_metal_server()) + + self.assertEqual(reserve_server.call_args.args[0].placement_rules.nodes, ()) + def test_automatic_selection_retries_after_a_capacity_error(self) -> None: service = MigrationService(migration_doc(destination_metal_server=None)) service.record_destination_metal_server_selection_failure = Mock() diff --git a/atlas/vm/core/test_vm_resize.py b/atlas/vm/core/test_vm_resize.py index cf575f781..88f04169e 100644 --- a/atlas/vm/core/test_vm_resize.py +++ b/atlas/vm/core/test_vm_resize.py @@ -33,6 +33,7 @@ def setUp(self) -> None: memory_mib=2048, disk_mib=20480, sleep_after_idle_seconds=0, + placement_rules=None, db_set=Mock(), ) self.resize = VirtualMachineResize(self.virtual_machine) diff --git a/atlas/vm/core/test_vm_service.py b/atlas/vm/core/test_vm_service.py index 363eacdb4..5dccf9df2 100644 --- a/atlas/vm/core/test_vm_service.py +++ b/atlas/vm/core/test_vm_service.py @@ -1,3 +1,4 @@ +import json from types import SimpleNamespace from unittest.mock import Mock, patch @@ -6,7 +7,7 @@ from atlas.atlas.core.exceptions import AtlasConflictError from atlas.vm.core.metal_client import MetalClientError -from atlas.vm.core.models import Route +from atlas.vm.core.models import Route, VirtualMachineCreateRequest from atlas.vm.core.placement import OutOfCapacity, PlacementStrategy from atlas.vm.core.vm_service import ( InsufficientHostCapacity, @@ -15,6 +16,8 @@ ) from atlas.vm.doctype.virtual_machine_image.virtual_machine_image import VirtualMachineImage +AFFINITY_RULES = [{"resource": "virtual_machine", "operator": "has_not", "tags": {"role": "cargo-server"}}] + def build_image(tenant_id: int, image_type: str = "machine") -> VirtualMachineImage: """Return one image document that answers the tenant visibility rule.""" @@ -40,6 +43,39 @@ def test_out_of_capacity_does_not_create_a_vm_draft(self) -> None: insert_draft.assert_not_called() + def test_any_tenant_can_send_placement_rules(self) -> None: + image = SimpleNamespace(architecture="amd64", validate_compatibility=Mock()) + # Tenant 7, on a VM that is not privileged. + request = {**self.request(), "placement_rules": AFFINITY_RULES} + with ( + patch.object(VirtualMachineService, "get_image", return_value=image), + patch.object( + PlacementStrategy, "find_server", side_effect=OutOfCapacity("retry later") + ) as find_server, + self.assertRaises(OutOfCapacity), + ): + VirtualMachineService.create(request) + + self.assertEqual(find_server.call_args.args[0].placement_rules.as_list(), AFFINITY_RULES) + + def test_the_draft_stores_the_tags_and_placement_rules(self) -> None: + image = SimpleNamespace(name="image-1", architecture="amd64") + for placement_rules in (AFFINITY_RULES, []): + request = VirtualMachineCreateRequest.from_value( + {**self.request(), "tags": {"role": "cargo-server"}, "placement_rules": placement_rules} + ) + with ( + self.subTest(placement_rules=placement_rules), + patch("atlas.vm.core.vm_service.frappe.get_doc") as get_doc, + ): + VirtualMachineService.insert_draft(request, image, "metal-1") + + values = get_doc.call_args.args[0] + self.assertEqual(values["tags"], [{"key": "role", "value": "cargo-server"}]) + # A VM without rules stores nothing. + stored_rules = values["placement_rules"] + self.assertEqual(json.loads(stored_rules) if stored_rules else [], placement_rules) + def test_creation_commits_the_draft_before_the_metal_request(self) -> None: operations: list[str] = [] image = SimpleNamespace( diff --git a/atlas/vm/core/vm_migration.py b/atlas/vm/core/vm_migration.py index f382ea5c3..9177754f5 100644 --- a/atlas/vm/core/vm_migration.py +++ b/atlas/vm/core/vm_migration.py @@ -13,7 +13,7 @@ from atlas.vm.core.metal_client import MetalClient, MetalClientError from atlas.vm.core.metal_models import timestamp_field from atlas.vm.core.models import VirtualMachineShape -from atlas.vm.core.placement import PlacementRequirements, PlacementStrategy +from atlas.vm.core.placement import AffinityRules, PlacementRequirements, PlacementStrategy from atlas.vm.core.placement.transaction import use_read_committed from atlas.vm.core.vm_state import LIVE_STATES @@ -165,6 +165,13 @@ def select_destination_metal_server(self) -> bool: "VirtualMachine", frappe.get_doc("Virtual Machine", self.migration.virtual_machine) ) shape = self.destination_shape(virtual_machine) + requested_destination = self.migration.destination_metal_server + # A destination that the operator names is not limited by the affinity rules. + placement_rules = ( + AffinityRules() + if requested_destination + else AffinityRules.from_json(virtual_machine.placement_rules) + ) requirements = PlacementRequirements( shape.cpu_millicores, shape.memory_mib, @@ -172,8 +179,9 @@ def select_destination_metal_server(self) -> bool: cast(str, virtual_machine.architecture), virtual_machine.tenant_id, shape.sleep_after_idle_seconds > 0, + placement_rules=placement_rules, + virtual_machine=virtual_machine.name, ) - requested_destination = self.migration.destination_metal_server exclude_servers = {self.migration.source_metal_server} try: diff --git a/atlas/vm/core/vm_resize.py b/atlas/vm/core/vm_resize.py index d6fd14bc8..761325e8e 100644 --- a/atlas/vm/core/vm_resize.py +++ b/atlas/vm/core/vm_resize.py @@ -13,7 +13,7 @@ MINIMUM_CPU_MILLICORES, VirtualMachineShape, ) -from atlas.vm.core.placement import CurrentPlacement, PlacementRequirements, PlacementStrategy +from atlas.vm.core.placement import AffinityRules, CurrentPlacement, PlacementRequirements, PlacementStrategy from atlas.vm.core.placement.transaction import use_read_committed from atlas.vm.core.vm_migration import MigrationService from atlas.vm.core.vm_service import VirtualMachineService @@ -132,6 +132,8 @@ def _get_requirements(self, target: VirtualMachineShape) -> PlacementRequirement cast(str, self.virtual_machine.architecture), self.virtual_machine.tenant_id, target.sleep_after_idle_seconds > 0, + placement_rules=AffinityRules.from_json(self.virtual_machine.placement_rules), + virtual_machine=self.virtual_machine.name, ) def _set_idle_shutdown(self, current: VirtualMachineShape, sleep_after_idle_seconds: int) -> None: diff --git a/atlas/vm/core/vm_service.py b/atlas/vm/core/vm_service.py index eb30263c3..b04d52386 100644 --- a/atlas/vm/core/vm_service.py +++ b/atlas/vm/core/vm_service.py @@ -85,6 +85,7 @@ def create(cls, value: str | dict[str, Any] | VirtualMachineCreateRequest) -> Vi image.architecture, request.tenant_id, request.sleep_after_idle_seconds > 0, + placement_rules=request.placement_rules, ) server_name = PlacementStrategy.find_server(requirements) virtual_machine = cls.insert_draft(request, image, server_name) @@ -144,6 +145,10 @@ def insert_draft( "is_privileged": request.is_privileged, "is_termination_protected": request.is_termination_protected, "sleep_after_idle_seconds": request.sleep_after_idle_seconds, + "tags": [{"key": key, "value": value} for key, value in request.tags.items()], + "placement_rules": frappe.as_json(request.placement_rules.as_list()) + if request.placement_rules.nodes + else None, } ) virtual_machine.flags.created_by_virtual_machine_api = True diff --git a/atlas/vm/doctype/virtual_machine/test_virtual_machine.py b/atlas/vm/doctype/virtual_machine/test_virtual_machine.py index 409b16cd4..f1083a60f 100644 --- a/atlas/vm/doctype/virtual_machine/test_virtual_machine.py +++ b/atlas/vm/doctype/virtual_machine/test_virtual_machine.py @@ -110,6 +110,36 @@ def test_request_parses_a_firewall(self) -> None: (FirewallRule(protocol="tcp", ports="22", cidrs=("203.0.113.0/24",)),), ) + def test_request_parses_tags_and_placement_rules(self) -> None: + rules = [{"any_of": [{"resource": "metal_server", "operator": "has", "tags": {"rack": "a"}}]}] + request = VirtualMachineCreateRequest.from_value( + { + "virtual_machine_image": "Ubuntu 24.04", + "cpu_millicores": 1500, + "memory_mib": 2048, + "disk_mib": 10240, + "tenant_id": 0, + "tags": {"role": "cargo-server"}, + "placement_rules": rules, + } + ) + + self.assertEqual(request.tags, {"role": "cargo-server"}) + self.assertEqual(request.placement_rules.as_list(), rules) + + def test_request_rejects_tags_that_are_not_strings(self) -> None: + with self.assertRaisesRegex(ValueError, "string-to-string"): + VirtualMachineCreateRequest.from_value( + { + "virtual_machine_image": "Ubuntu 24.04", + "cpu_millicores": 1500, + "memory_mib": 2048, + "disk_mib": 10240, + "tenant_id": 0, + "tags": {"role": 1}, + } + ) + def test_request_rejects_a_noncanonical_firewall_cidr(self) -> None: with self.assertRaisesRegex(ValueError, "canonical"): VirtualMachineCreateRequest.from_value( diff --git a/atlas/vm/doctype/virtual_machine/virtual_machine.json b/atlas/vm/doctype/virtual_machine/virtual_machine.json index 5604de288..d5f13bfbc 100644 --- a/atlas/vm/doctype/virtual_machine/virtual_machine.json +++ b/atlas/vm/doctype/virtual_machine/virtual_machine.json @@ -41,6 +41,7 @@ "ssh_keys", "firewall_summary", "metadata", + "placement_rules", "tags" ], "fields": [ @@ -304,6 +305,15 @@ "options": "JSON", "read_only": 1 }, + { + "description": "Rules that limit the Metal Servers for this virtual machine. Placement applies them at creation, at a migration to a host that placement chooses, and at a resize that moves the virtual machine.", + "fieldname": "placement_rules", + "fieldtype": "Code", + "label": "Placement Rules", + "options": "JSON", + "read_only": 1, + "show_description_on_click": 1 + }, { "fieldname": "column_break_status", "fieldtype": "Column Break" @@ -355,7 +365,7 @@ "link_fieldname": "virtual_machine" } ], - "modified": "2026-09-26 03:05:51.552687", + "modified": "2026-10-02 00:00:00.000000", "modified_by": "Administrator", "module": "VM", "name": "Virtual Machine", diff --git a/atlas/vm/doctype/virtual_machine/virtual_machine.py b/atlas/vm/doctype/virtual_machine/virtual_machine.py index 68e094560..a3033dccc 100644 --- a/atlas/vm/doctype/virtual_machine/virtual_machine.py +++ b/atlas/vm/doctype/virtual_machine/virtual_machine.py @@ -55,6 +55,7 @@ class VirtualMachine(Document): is_termination_protected: DF.Check memory_mib: DF.Int metadata: DF.Code | None + placement_rules: DF.Code | None routes: DF.Code | None server: DF.Link sleep_after_idle_seconds: DF.Int diff --git a/clients/atlas-client/atlas_client/api/vm_lifecycle/create_virtual_machine.py b/clients/atlas-client/atlas_client/api/vm_lifecycle/create_virtual_machine.py index 32644cac6..2dd8d8d70 100644 --- a/clients/atlas-client/atlas_client/api/vm_lifecycle/create_virtual_machine.py +++ b/clients/atlas-client/atlas_client/api/vm_lifecycle/create_virtual_machine.py @@ -120,6 +120,11 @@ def sync_detailed( Set `is_termination_protected` to refuse deletion of the new VM. The termination protection route changes it later. + Use `tags` to label the VM, for example, `{"role": "cargo-server"}`. `placement_rules` limits the + Metal Servers for the VM by host tags and by the tags of other VMs on the host. Atlas Settings + selects whether a VM that no host with room can satisfy fails with `affinity_unsatisfied` or goes to + any host. + Args: x_tenant_id (int): body (CreateVirtualMachinePayload): Values that create one virtual machine. @@ -161,6 +166,11 @@ def sync( Set `is_termination_protected` to refuse deletion of the new VM. The termination protection route changes it later. + Use `tags` to label the VM, for example, `{"role": "cargo-server"}`. `placement_rules` limits the + Metal Servers for the VM by host tags and by the tags of other VMs on the host. Atlas Settings + selects whether a VM that no host with room can satisfy fails with `affinity_unsatisfied` or goes to + any host. + Args: x_tenant_id (int): body (CreateVirtualMachinePayload): Values that create one virtual machine. @@ -197,6 +207,11 @@ async def asyncio_detailed( Set `is_termination_protected` to refuse deletion of the new VM. The termination protection route changes it later. + Use `tags` to label the VM, for example, `{"role": "cargo-server"}`. `placement_rules` limits the + Metal Servers for the VM by host tags and by the tags of other VMs on the host. Atlas Settings + selects whether a VM that no host with room can satisfy fails with `affinity_unsatisfied` or goes to + any host. + Args: x_tenant_id (int): body (CreateVirtualMachinePayload): Values that create one virtual machine. @@ -238,6 +253,11 @@ async def asyncio( Set `is_termination_protected` to refuse deletion of the new VM. The termination protection route changes it later. + Use `tags` to label the VM, for example, `{"role": "cargo-server"}`. `placement_rules` limits the + Metal Servers for the VM by host tags and by the tags of other VMs on the host. Atlas Settings + selects whether a VM that no host with room can satisfy fails with `affinity_unsatisfied` or goes to + any host. + Args: x_tenant_id (int): body (CreateVirtualMachinePayload): Values that create one virtual machine. diff --git a/clients/atlas-client/atlas_client/models/__init__.py b/clients/atlas-client/atlas_client/models/__init__.py index 45169180b..325cd660c 100644 --- a/clients/atlas-client/atlas_client/models/__init__.py +++ b/clients/atlas-client/atlas_client/models/__init__.py @@ -1,5 +1,12 @@ """ Contains all the data models used in inputs/outputs """ +from .affinity_all_of_payload import AffinityAllOfPayload +from .affinity_any_of_payload import AffinityAnyOfPayload +from .affinity_rule_payload import AffinityRulePayload +from .affinity_rule_payload_operator import AffinityRulePayloadOperator +from .affinity_rule_payload_resource import AffinityRulePayloadResource +from .affinity_rule_payload_tags import AffinityRulePayloadTags +from .affinity_unsatisfied_error import AffinityUnsatisfiedError from .api_error_detail import ApiErrorDetail from .api_error_field import ApiErrorField from .api_error_response import ApiErrorResponse @@ -11,6 +18,7 @@ from .console_token_response_mode import ConsoleTokenResponseMode from .create_virtual_machine_payload import CreateVirtualMachinePayload from .create_virtual_machine_payload_metadata import CreateVirtualMachinePayloadMetadata +from .create_virtual_machine_payload_tags import CreateVirtualMachinePayloadTags from .disk_update_payload import DiskUpdatePayload from .download_image_artifact import DownloadImageArtifact from .firewall_payload import FirewallPayload @@ -71,6 +79,13 @@ from .webhook_configuration_response import WebhookConfigurationResponse __all__ = ( + "AffinityAllOfPayload", + "AffinityAnyOfPayload", + "AffinityRulePayload", + "AffinityRulePayloadOperator", + "AffinityRulePayloadResource", + "AffinityRulePayloadTags", + "AffinityUnsatisfiedError", "ApiErrorDetail", "ApiErrorField", "ApiErrorResponse", @@ -82,6 +97,7 @@ "ConsoleTokenResponseMode", "CreateVirtualMachinePayload", "CreateVirtualMachinePayloadMetadata", + "CreateVirtualMachinePayloadTags", "DiskUpdatePayload", "DownloadImageArtifact", "FirewallPayload", diff --git a/clients/atlas-client/atlas_client/models/affinity_all_of_payload.py b/clients/atlas-client/atlas_client/models/affinity_all_of_payload.py new file mode 100644 index 000000000..7fbe6dd42 --- /dev/null +++ b/clients/atlas-client/atlas_client/models/affinity_all_of_payload.py @@ -0,0 +1,115 @@ +from __future__ import annotations + +from collections.abc import Mapping +from typing import Any, TypeVar, BinaryIO, TextIO, TYPE_CHECKING, Generator + +from attrs import define as _attrs_define +from attrs import field as _attrs_field + +from ..types import UNSET, Unset + +from typing import cast + +if TYPE_CHECKING: + from ..models.affinity_any_of_payload import AffinityAnyOfPayload + from ..models.affinity_rule_payload import AffinityRulePayload + + + + + +T = TypeVar("T", bound="AffinityAllOfPayload") + + + +@_attrs_define +class AffinityAllOfPayload: + """ A group that holds when every one of its rules or groups holds. + + Attributes: + all_of (list[AffinityAllOfPayload | AffinityAnyOfPayload | AffinityRulePayload]): Rules or groups. Every one + must hold. + """ + + all_of: list[AffinityAllOfPayload | AffinityAnyOfPayload | AffinityRulePayload] + + + + + + def to_dict(self) -> dict[str, Any]: + from ..models.affinity_any_of_payload import AffinityAnyOfPayload # noqa: PLC0415 + from ..models.affinity_rule_payload import AffinityRulePayload # noqa: PLC0415 + all_of = [] + for all_of_item_data in self.all_of: + all_of_item: dict[str, Any] + if isinstance(all_of_item_data, AffinityRulePayload): + all_of_item = all_of_item_data.to_dict() + elif isinstance(all_of_item_data, AffinityAnyOfPayload): + all_of_item = all_of_item_data.to_dict() + else: + all_of_item = all_of_item_data.to_dict() + + all_of.append(all_of_item) + + + + + field_dict: dict[str, Any] = {} + + field_dict.update({ + "all_of": all_of, + }) + + return field_dict + + + + @classmethod + def from_dict(cls: type[T], src_dict: Mapping[str, Any]) -> T: + from ..models.affinity_any_of_payload import AffinityAnyOfPayload # noqa: PLC0415 + from ..models.affinity_rule_payload import AffinityRulePayload # noqa: PLC0415 + d = dict(src_dict) + all_of = [] + _all_of = d.pop("all_of") + for all_of_item_data in (_all_of): + def _parse_all_of_item(data: object) -> AffinityAllOfPayload | AffinityAnyOfPayload | AffinityRulePayload: + try: + if not isinstance(data, dict): + raise TypeError() + all_of_item_type_0 = AffinityRulePayload.from_dict(data) + + + + return all_of_item_type_0 + except (TypeError, ValueError, AttributeError, KeyError): + pass + try: + if not isinstance(data, dict): + raise TypeError() + all_of_item_type_1 = AffinityAnyOfPayload.from_dict(data) + + + + return all_of_item_type_1 + except (TypeError, ValueError, AttributeError, KeyError): + pass + if not isinstance(data, dict): + raise TypeError() + all_of_item_type_2 = AffinityAllOfPayload.from_dict(data) + + + + return all_of_item_type_2 + + all_of_item = _parse_all_of_item(all_of_item_data) + + all_of.append(all_of_item) + + + affinity_all_of_payload = cls( + all_of=all_of, + ) + + return affinity_all_of_payload + diff --git a/clients/atlas-client/atlas_client/models/affinity_any_of_payload.py b/clients/atlas-client/atlas_client/models/affinity_any_of_payload.py new file mode 100644 index 000000000..cb2e85993 --- /dev/null +++ b/clients/atlas-client/atlas_client/models/affinity_any_of_payload.py @@ -0,0 +1,115 @@ +from __future__ import annotations + +from collections.abc import Mapping +from typing import Any, TypeVar, BinaryIO, TextIO, TYPE_CHECKING, Generator + +from attrs import define as _attrs_define +from attrs import field as _attrs_field + +from ..types import UNSET, Unset + +from typing import cast + +if TYPE_CHECKING: + from ..models.affinity_all_of_payload import AffinityAllOfPayload + from ..models.affinity_rule_payload import AffinityRulePayload + + + + + +T = TypeVar("T", bound="AffinityAnyOfPayload") + + + +@_attrs_define +class AffinityAnyOfPayload: + """ A group that holds when at least one of its rules or groups holds. + + Attributes: + any_of (list[AffinityAllOfPayload | AffinityAnyOfPayload | AffinityRulePayload]): Rules or groups. One must + hold. + """ + + any_of: list[AffinityAllOfPayload | AffinityAnyOfPayload | AffinityRulePayload] + + + + + + def to_dict(self) -> dict[str, Any]: + from ..models.affinity_all_of_payload import AffinityAllOfPayload # noqa: PLC0415 + from ..models.affinity_rule_payload import AffinityRulePayload # noqa: PLC0415 + any_of = [] + for any_of_item_data in self.any_of: + any_of_item: dict[str, Any] + if isinstance(any_of_item_data, AffinityRulePayload): + any_of_item = any_of_item_data.to_dict() + elif isinstance(any_of_item_data, AffinityAnyOfPayload): + any_of_item = any_of_item_data.to_dict() + else: + any_of_item = any_of_item_data.to_dict() + + any_of.append(any_of_item) + + + + + field_dict: dict[str, Any] = {} + + field_dict.update({ + "any_of": any_of, + }) + + return field_dict + + + + @classmethod + def from_dict(cls: type[T], src_dict: Mapping[str, Any]) -> T: + from ..models.affinity_all_of_payload import AffinityAllOfPayload # noqa: PLC0415 + from ..models.affinity_rule_payload import AffinityRulePayload # noqa: PLC0415 + d = dict(src_dict) + any_of = [] + _any_of = d.pop("any_of") + for any_of_item_data in (_any_of): + def _parse_any_of_item(data: object) -> AffinityAllOfPayload | AffinityAnyOfPayload | AffinityRulePayload: + try: + if not isinstance(data, dict): + raise TypeError() + any_of_item_type_0 = AffinityRulePayload.from_dict(data) + + + + return any_of_item_type_0 + except (TypeError, ValueError, AttributeError, KeyError): + pass + try: + if not isinstance(data, dict): + raise TypeError() + any_of_item_type_1 = AffinityAnyOfPayload.from_dict(data) + + + + return any_of_item_type_1 + except (TypeError, ValueError, AttributeError, KeyError): + pass + if not isinstance(data, dict): + raise TypeError() + any_of_item_type_2 = AffinityAllOfPayload.from_dict(data) + + + + return any_of_item_type_2 + + any_of_item = _parse_any_of_item(any_of_item_data) + + any_of.append(any_of_item) + + + affinity_any_of_payload = cls( + any_of=any_of, + ) + + return affinity_any_of_payload + diff --git a/clients/atlas-client/atlas_client/models/affinity_rule_payload.py b/clients/atlas-client/atlas_client/models/affinity_rule_payload.py new file mode 100644 index 000000000..5bac7a152 --- /dev/null +++ b/clients/atlas-client/atlas_client/models/affinity_rule_payload.py @@ -0,0 +1,116 @@ +from __future__ import annotations + +from collections.abc import Mapping +from typing import Any, TypeVar, BinaryIO, TextIO, TYPE_CHECKING, Generator + +from attrs import define as _attrs_define +from attrs import field as _attrs_field + +from ..types import UNSET, Unset + +from ..models.affinity_rule_payload_operator import AffinityRulePayloadOperator +from ..models.affinity_rule_payload_resource import AffinityRulePayloadResource +from ..types import UNSET, Unset +from typing import cast + +if TYPE_CHECKING: + from ..models.affinity_rule_payload_tags import AffinityRulePayloadTags + + + + + +T = TypeVar("T", bound="AffinityRulePayload") + + + +@_attrs_define +class AffinityRulePayload: + """ One rule on the tags of the candidate Metal Server, or of a VM that runs on it. + + Attributes: + operator (AffinityRulePayloadOperator): `has` needs a resource with every tag pair. `has_not` rejects such a + resource. + resource (AffinityRulePayloadResource): `metal_server` checks the candidate host. `virtual_machine` checks the + VMs on the candidate host. + tags (AffinityRulePayloadTags): Tag pairs that must all be on one resource. + within (None | str | Unset): A host tag key, such as `rack`. A `virtual_machine` rule then reads the VMs on + every host that has the same value for this key as the candidate host. A host without the key fails the rule. + """ + + operator: AffinityRulePayloadOperator + resource: AffinityRulePayloadResource + tags: AffinityRulePayloadTags + within: None | str | Unset = UNSET + + + + + + def to_dict(self) -> dict[str, Any]: + from ..models.affinity_rule_payload_tags import AffinityRulePayloadTags # noqa: PLC0415 + operator = self.operator.value + + resource = self.resource.value + + tags = self.tags.to_dict() + + within: None | str | Unset + if isinstance(self.within, Unset): + within = UNSET + else: + within = self.within + + + field_dict: dict[str, Any] = {} + + field_dict.update({ + "operator": operator, + "resource": resource, + "tags": tags, + }) + if within is not UNSET: + field_dict["within"] = within + + return field_dict + + + + @classmethod + def from_dict(cls: type[T], src_dict: Mapping[str, Any]) -> T: + from ..models.affinity_rule_payload_tags import AffinityRulePayloadTags # noqa: PLC0415 + d = dict(src_dict) + operator = AffinityRulePayloadOperator(d.pop("operator")) + + + + + resource = AffinityRulePayloadResource(d.pop("resource")) + + + + + tags = AffinityRulePayloadTags.from_dict(d.pop("tags")) + + + + + def _parse_within(data: object) -> None | str | Unset: + if data is None: + return data + if isinstance(data, Unset): + return data + return cast(None | str | Unset, data) + + within = _parse_within(d.pop("within", UNSET)) + + + affinity_rule_payload = cls( + operator=operator, + resource=resource, + tags=tags, + within=within, + ) + + return affinity_rule_payload + diff --git a/clients/atlas-client/atlas_client/models/affinity_rule_payload_operator.py b/clients/atlas-client/atlas_client/models/affinity_rule_payload_operator.py new file mode 100644 index 000000000..a7b7e1b07 --- /dev/null +++ b/clients/atlas-client/atlas_client/models/affinity_rule_payload_operator.py @@ -0,0 +1,8 @@ +from enum import StrEnum + +class AffinityRulePayloadOperator(StrEnum): + HAS = "has" + HAS_NOT = "has_not" + + def __str__(self) -> str: + return str(self.value) diff --git a/clients/atlas-client/atlas_client/models/affinity_rule_payload_resource.py b/clients/atlas-client/atlas_client/models/affinity_rule_payload_resource.py new file mode 100644 index 000000000..553427220 --- /dev/null +++ b/clients/atlas-client/atlas_client/models/affinity_rule_payload_resource.py @@ -0,0 +1,8 @@ +from enum import StrEnum + +class AffinityRulePayloadResource(StrEnum): + METAL_SERVER = "metal_server" + VIRTUAL_MACHINE = "virtual_machine" + + def __str__(self) -> str: + return str(self.value) diff --git a/clients/atlas-client/atlas_client/models/affinity_rule_payload_tags.py b/clients/atlas-client/atlas_client/models/affinity_rule_payload_tags.py new file mode 100644 index 000000000..43b8ad2be --- /dev/null +++ b/clients/atlas-client/atlas_client/models/affinity_rule_payload_tags.py @@ -0,0 +1,66 @@ +from __future__ import annotations + +from collections.abc import Mapping +from typing import Any, TypeVar, BinaryIO, TextIO, TYPE_CHECKING, Generator + +from attrs import define as _attrs_define +from attrs import field as _attrs_field + +from ..types import UNSET, Unset + + + + + + + +T = TypeVar("T", bound="AffinityRulePayloadTags") + + + +@_attrs_define +class AffinityRulePayloadTags: + """ Tag pairs that must all be on one resource. + + """ + + additional_properties: dict[str, str] = _attrs_field(init=False, factory=dict) + + + + + + def to_dict(self) -> dict[str, Any]: + + field_dict: dict[str, Any] = {} + field_dict.update(self.additional_properties) + + return field_dict + + + + @classmethod + def from_dict(cls: type[T], src_dict: Mapping[str, Any]) -> T: + d = dict(src_dict) + affinity_rule_payload_tags = cls( + ) + + + affinity_rule_payload_tags.additional_properties = d + return affinity_rule_payload_tags + + @property + def additional_keys(self) -> list[str]: + return list(self.additional_properties.keys()) + + def __getitem__(self, key: str) -> str: + return self.additional_properties[key] + + def __setitem__(self, key: str, value: str) -> None: + self.additional_properties[key] = value + + def __delitem__(self, key: str) -> None: + del self.additional_properties[key] + + def __contains__(self, key: str) -> bool: + return key in self.additional_properties diff --git a/clients/atlas-client/atlas_client/models/affinity_unsatisfied_error.py b/clients/atlas-client/atlas_client/models/affinity_unsatisfied_error.py new file mode 100644 index 000000000..045b240d1 --- /dev/null +++ b/clients/atlas-client/atlas_client/models/affinity_unsatisfied_error.py @@ -0,0 +1,114 @@ +from __future__ import annotations + +from collections.abc import Mapping +from typing import Any, TypeVar, BinaryIO, TextIO, TYPE_CHECKING, Generator + +from attrs import define as _attrs_define +from attrs import field as _attrs_field + +from ..types import UNSET, Unset + +from typing import cast +from typing import Literal, cast + +if TYPE_CHECKING: + from ..models.api_error_field import ApiErrorField + + + + + +T = TypeVar("T", bound="AffinityUnsatisfiedError") + + + +@_attrs_define +class AffinityUnsatisfiedError: + """ No host with room meets the affinity rules of the VM. + + Attributes: + code (Literal['affinity_unsatisfied']): Stable machine-readable error code. + fields (list[ApiErrorField]): Invalid request fields, or an empty list. + message (str): Safe description of the failure. + """ + + code: Literal['affinity_unsatisfied'] + fields: list[ApiErrorField] + message: str + additional_properties: dict[str, Any] = _attrs_field(init=False, factory=dict) + + + + + + def to_dict(self) -> dict[str, Any]: + from ..models.api_error_field import ApiErrorField # noqa: PLC0415 + code = self.code + + fields = [] + for fields_item_data in self.fields: + fields_item = fields_item_data.to_dict() + fields.append(fields_item) + + + + message = self.message + + + field_dict: dict[str, Any] = {} + field_dict.update(self.additional_properties) + field_dict.update({ + "code": code, + "fields": fields, + "message": message, + }) + + return field_dict + + + + @classmethod + def from_dict(cls: type[T], src_dict: Mapping[str, Any]) -> T: + from ..models.api_error_field import ApiErrorField # noqa: PLC0415 + d = dict(src_dict) + code = cast(Literal['affinity_unsatisfied'] , d.pop("code")) + if code != 'affinity_unsatisfied': + raise ValueError(f"code must match const 'affinity_unsatisfied', got '{code}'") + + fields = [] + _fields = d.pop("fields") + for fields_item_data in (_fields): + fields_item = ApiErrorField.from_dict(fields_item_data) + + + + fields.append(fields_item) + + + message = d.pop("message") + + affinity_unsatisfied_error = cls( + code=code, + fields=fields, + message=message, + ) + + + affinity_unsatisfied_error.additional_properties = d + return affinity_unsatisfied_error + + @property + def additional_keys(self) -> list[str]: + return list(self.additional_properties.keys()) + + def __getitem__(self, key: str) -> Any: + return self.additional_properties[key] + + def __setitem__(self, key: str, value: Any) -> None: + self.additional_properties[key] = value + + def __delitem__(self, key: str) -> None: + del self.additional_properties[key] + + def __contains__(self, key: str) -> bool: + return key in self.additional_properties diff --git a/clients/atlas-client/atlas_client/models/capacity_unavailable_response.py b/clients/atlas-client/atlas_client/models/capacity_unavailable_response.py index d72cea094..b17b55acc 100644 --- a/clients/atlas-client/atlas_client/models/capacity_unavailable_response.py +++ b/clients/atlas-client/atlas_client/models/capacity_unavailable_response.py @@ -11,6 +11,7 @@ from typing import cast if TYPE_CHECKING: + from ..models.affinity_unsatisfied_error import AffinityUnsatisfiedError from ..models.out_of_capacity_error import OutOfCapacityError from ..models.placement_busy_error import PlacementBusyError @@ -30,10 +31,10 @@ class CapacityUnavailableResponse: caller can retry a busy placement at once and escalate a full one. Attributes: - error (OutOfCapacityError | PlacementBusyError): Capacity failure details. + error (AffinityUnsatisfiedError | OutOfCapacityError | PlacementBusyError): Capacity failure details. """ - error: OutOfCapacityError | PlacementBusyError + error: AffinityUnsatisfiedError | OutOfCapacityError | PlacementBusyError additional_properties: dict[str, Any] = _attrs_field(init=False, factory=dict) @@ -41,11 +42,14 @@ class CapacityUnavailableResponse: def to_dict(self) -> dict[str, Any]: + from ..models.affinity_unsatisfied_error import AffinityUnsatisfiedError # noqa: PLC0415 from ..models.out_of_capacity_error import OutOfCapacityError # noqa: PLC0415 from ..models.placement_busy_error import PlacementBusyError # noqa: PLC0415 error: dict[str, Any] if isinstance(self.error, OutOfCapacityError): error = self.error.to_dict() + elif isinstance(self.error, PlacementBusyError): + error = self.error.to_dict() else: error = self.error.to_dict() @@ -63,10 +67,11 @@ def to_dict(self) -> dict[str, Any]: @classmethod def from_dict(cls: type[T], src_dict: Mapping[str, Any]) -> T: + from ..models.affinity_unsatisfied_error import AffinityUnsatisfiedError # noqa: PLC0415 from ..models.out_of_capacity_error import OutOfCapacityError # noqa: PLC0415 from ..models.placement_busy_error import PlacementBusyError # noqa: PLC0415 d = dict(src_dict) - def _parse_error(data: object) -> OutOfCapacityError | PlacementBusyError: + def _parse_error(data: object) -> AffinityUnsatisfiedError | OutOfCapacityError | PlacementBusyError: try: if not isinstance(data, dict): raise TypeError() @@ -77,13 +82,23 @@ def _parse_error(data: object) -> OutOfCapacityError | PlacementBusyError: return error_type_0 except (TypeError, ValueError, AttributeError, KeyError): pass + try: + if not isinstance(data, dict): + raise TypeError() + error_type_1 = PlacementBusyError.from_dict(data) + + + + return error_type_1 + except (TypeError, ValueError, AttributeError, KeyError): + pass if not isinstance(data, dict): raise TypeError() - error_type_1 = PlacementBusyError.from_dict(data) + error_type_2 = AffinityUnsatisfiedError.from_dict(data) - return error_type_1 + return error_type_2 error = _parse_error(d.pop("error")) diff --git a/clients/atlas-client/atlas_client/models/create_virtual_machine_payload.py b/clients/atlas-client/atlas_client/models/create_virtual_machine_payload.py index 237f329ab..43c5b3078 100644 --- a/clients/atlas-client/atlas_client/models/create_virtual_machine_payload.py +++ b/clients/atlas-client/atlas_client/models/create_virtual_machine_payload.py @@ -12,7 +12,11 @@ from typing import cast if TYPE_CHECKING: + from ..models.affinity_all_of_payload import AffinityAllOfPayload + from ..models.affinity_any_of_payload import AffinityAnyOfPayload + from ..models.affinity_rule_payload import AffinityRulePayload from ..models.create_virtual_machine_payload_metadata import CreateVirtualMachinePayloadMetadata + from ..models.create_virtual_machine_payload_tags import CreateVirtualMachinePayloadTags from ..models.firewall_payload import FirewallPayload @@ -41,6 +45,9 @@ class CreateVirtualMachinePayload: is_privileged (bool | Unset): Whether the guest can reach every tenant through the mesh. Default: False. is_termination_protected (bool | Unset): Whether deletion is blocked. Default: False. metadata (CreateVirtualMachinePayloadMetadata | Unset): Custom guest metadata. + placement_rules (list[AffinityAllOfPayload | AffinityAnyOfPayload | AffinityRulePayload] | Unset): Rules that + limit the Metal Servers for the virtual machine. Every listed rule or group must hold. Placement uses only the + hosts that meet them. private_network_throughput_mibps (int | Unset): Private network throughput limit in MiB/s. Zero removes the limit. Default: 0. public_ipv4 (None | str | Unset): Reserved public IPv4 allocation ID, or null. @@ -49,6 +56,7 @@ class CreateVirtualMachinePayload: Default: 0. sleep_after_idle_seconds (int | Unset): Idle time before automatic stop. Zero disables it. Default: 0. ssh_keys (list[str] | Unset): Authorized SSH public keys. + tags (CreateVirtualMachinePayloadTags | Unset): Resource tags as key-value pairs. user_data (str | Unset): Cloud-init user data supplied to the guest. Default: ''. """ @@ -64,12 +72,14 @@ class CreateVirtualMachinePayload: is_privileged: bool | Unset = False is_termination_protected: bool | Unset = False metadata: CreateVirtualMachinePayloadMetadata | Unset = UNSET + placement_rules: list[AffinityAllOfPayload | AffinityAnyOfPayload | AffinityRulePayload] | Unset = UNSET private_network_throughput_mibps: int | Unset = 0 public_ipv4: None | str | Unset = UNSET public_ipv6: None | str | Unset = UNSET public_network_throughput_mibps: int | Unset = 0 sleep_after_idle_seconds: int | Unset = 0 ssh_keys: list[str] | Unset = UNSET + tags: CreateVirtualMachinePayloadTags | Unset = UNSET user_data: str | Unset = '' @@ -77,7 +87,11 @@ class CreateVirtualMachinePayload: def to_dict(self) -> dict[str, Any]: + from ..models.affinity_all_of_payload import AffinityAllOfPayload # noqa: PLC0415 + from ..models.affinity_any_of_payload import AffinityAnyOfPayload # noqa: PLC0415 + from ..models.affinity_rule_payload import AffinityRulePayload # noqa: PLC0415 from ..models.create_virtual_machine_payload_metadata import CreateVirtualMachinePayloadMetadata # noqa: PLC0415 + from ..models.create_virtual_machine_payload_tags import CreateVirtualMachinePayloadTags # noqa: PLC0415 from ..models.firewall_payload import FirewallPayload # noqa: PLC0415 cpu_millicores = self.cpu_millicores @@ -107,6 +121,22 @@ def to_dict(self) -> dict[str, Any]: if not isinstance(self.metadata, Unset): metadata = self.metadata.to_dict() + placement_rules: list[dict[str, Any]] | Unset = UNSET + if not isinstance(self.placement_rules, Unset): + placement_rules = [] + for placement_rules_item_data in self.placement_rules: + placement_rules_item: dict[str, Any] + if isinstance(placement_rules_item_data, AffinityRulePayload): + placement_rules_item = placement_rules_item_data.to_dict() + elif isinstance(placement_rules_item_data, AffinityAnyOfPayload): + placement_rules_item = placement_rules_item_data.to_dict() + else: + placement_rules_item = placement_rules_item_data.to_dict() + + placement_rules.append(placement_rules_item) + + + private_network_throughput_mibps = self.private_network_throughput_mibps public_ipv4: None | str | Unset @@ -131,6 +161,10 @@ def to_dict(self) -> dict[str, Any]: + tags: dict[str, Any] | Unset = UNSET + if not isinstance(self.tags, Unset): + tags = self.tags.to_dict() + user_data = self.user_data @@ -158,6 +192,8 @@ def to_dict(self) -> dict[str, Any]: field_dict["is_termination_protected"] = is_termination_protected if metadata is not UNSET: field_dict["metadata"] = metadata + if placement_rules is not UNSET: + field_dict["placement_rules"] = placement_rules if private_network_throughput_mibps is not UNSET: field_dict["private_network_throughput_mibps"] = private_network_throughput_mibps if public_ipv4 is not UNSET: @@ -170,6 +206,8 @@ def to_dict(self) -> dict[str, Any]: field_dict["sleep_after_idle_seconds"] = sleep_after_idle_seconds if ssh_keys is not UNSET: field_dict["ssh_keys"] = ssh_keys + if tags is not UNSET: + field_dict["tags"] = tags if user_data is not UNSET: field_dict["user_data"] = user_data @@ -179,7 +217,11 @@ def to_dict(self) -> dict[str, Any]: @classmethod def from_dict(cls: type[T], src_dict: Mapping[str, Any]) -> T: + from ..models.affinity_all_of_payload import AffinityAllOfPayload # noqa: PLC0415 + from ..models.affinity_any_of_payload import AffinityAnyOfPayload # noqa: PLC0415 + from ..models.affinity_rule_payload import AffinityRulePayload # noqa: PLC0415 from ..models.create_virtual_machine_payload_metadata import CreateVirtualMachinePayloadMetadata # noqa: PLC0415 + from ..models.create_virtual_machine_payload_tags import CreateVirtualMachinePayloadTags # noqa: PLC0415 from ..models.firewall_payload import FirewallPayload # noqa: PLC0415 d = dict(src_dict) cpu_millicores = d.pop("cpu_millicores") @@ -222,6 +264,45 @@ def from_dict(cls: type[T], src_dict: Mapping[str, Any]) -> T: + _placement_rules = d.pop("placement_rules", UNSET) + placement_rules: list[AffinityAllOfPayload | AffinityAnyOfPayload | AffinityRulePayload] | Unset = UNSET + if _placement_rules is not UNSET: + placement_rules = [] + for placement_rules_item_data in _placement_rules: + def _parse_placement_rules_item(data: object) -> AffinityAllOfPayload | AffinityAnyOfPayload | AffinityRulePayload: + try: + if not isinstance(data, dict): + raise TypeError() + placement_rules_item_type_0 = AffinityRulePayload.from_dict(data) + + + + return placement_rules_item_type_0 + except (TypeError, ValueError, AttributeError, KeyError): + pass + try: + if not isinstance(data, dict): + raise TypeError() + placement_rules_item_type_1 = AffinityAnyOfPayload.from_dict(data) + + + + return placement_rules_item_type_1 + except (TypeError, ValueError, AttributeError, KeyError): + pass + if not isinstance(data, dict): + raise TypeError() + placement_rules_item_type_2 = AffinityAllOfPayload.from_dict(data) + + + + return placement_rules_item_type_2 + + placement_rules_item = _parse_placement_rules_item(placement_rules_item_data) + + placement_rules.append(placement_rules_item) + + private_network_throughput_mibps = d.pop("private_network_throughput_mibps", UNSET) def _parse_public_ipv4(data: object) -> None | str | Unset: @@ -251,6 +332,16 @@ def _parse_public_ipv6(data: object) -> None | str | Unset: ssh_keys = cast(list[str], d.pop("ssh_keys", UNSET)) + _tags = d.pop("tags", UNSET) + tags: CreateVirtualMachinePayloadTags | Unset + if isinstance(_tags, Unset): + tags = UNSET + else: + tags = CreateVirtualMachinePayloadTags.from_dict(_tags) + + + + user_data = d.pop("user_data", UNSET) create_virtual_machine_payload = cls( @@ -266,12 +357,14 @@ def _parse_public_ipv6(data: object) -> None | str | Unset: is_privileged=is_privileged, is_termination_protected=is_termination_protected, metadata=metadata, + placement_rules=placement_rules, private_network_throughput_mibps=private_network_throughput_mibps, public_ipv4=public_ipv4, public_ipv6=public_ipv6, public_network_throughput_mibps=public_network_throughput_mibps, sleep_after_idle_seconds=sleep_after_idle_seconds, ssh_keys=ssh_keys, + tags=tags, user_data=user_data, ) diff --git a/clients/atlas-client/atlas_client/models/create_virtual_machine_payload_tags.py b/clients/atlas-client/atlas_client/models/create_virtual_machine_payload_tags.py new file mode 100644 index 000000000..8dd16d642 --- /dev/null +++ b/clients/atlas-client/atlas_client/models/create_virtual_machine_payload_tags.py @@ -0,0 +1,66 @@ +from __future__ import annotations + +from collections.abc import Mapping +from typing import Any, TypeVar, BinaryIO, TextIO, TYPE_CHECKING, Generator + +from attrs import define as _attrs_define +from attrs import field as _attrs_field + +from ..types import UNSET, Unset + + + + + + + +T = TypeVar("T", bound="CreateVirtualMachinePayloadTags") + + + +@_attrs_define +class CreateVirtualMachinePayloadTags: + """ Resource tags as key-value pairs. + + """ + + additional_properties: dict[str, str] = _attrs_field(init=False, factory=dict) + + + + + + def to_dict(self) -> dict[str, Any]: + + field_dict: dict[str, Any] = {} + field_dict.update(self.additional_properties) + + return field_dict + + + + @classmethod + def from_dict(cls: type[T], src_dict: Mapping[str, Any]) -> T: + d = dict(src_dict) + create_virtual_machine_payload_tags = cls( + ) + + + create_virtual_machine_payload_tags.additional_properties = d + return create_virtual_machine_payload_tags + + @property + def additional_keys(self) -> list[str]: + return list(self.additional_properties.keys()) + + def __getitem__(self, key: str) -> str: + return self.additional_properties[key] + + def __setitem__(self, key: str, value: str) -> None: + self.additional_properties[key] = value + + def __delitem__(self, key: str) -> None: + del self.additional_properties[key] + + def __contains__(self, key: str) -> bool: + return key in self.additional_properties diff --git a/clients/openapi/atlas-client.json b/clients/openapi/atlas-client.json index e07d07e6f..a09febde1 100644 --- a/clients/openapi/atlas-client.json +++ b/clients/openapi/atlas-client.json @@ -1,6 +1,157 @@ { "components": { "schemas": { + "AffinityAllOfPayload": { + "additionalProperties": false, + "description": "A group that holds when every one of its rules or groups holds.", + "properties": { + "all_of": { + "description": "Rules or groups. Every one must hold.", + "items": { + "oneOf": [ + { + "$ref": "#/components/schemas/AffinityRulePayload" + }, + { + "$ref": "#/components/schemas/AffinityAnyOfPayload" + }, + { + "$ref": "#/components/schemas/AffinityAllOfPayload" + } + ] + }, + "minItems": 1, + "title": "All Of", + "type": "array" + } + }, + "required": [ + "all_of" + ], + "title": "AffinityAllOfPayload", + "type": "object" + }, + "AffinityAnyOfPayload": { + "additionalProperties": false, + "description": "A group that holds when at least one of its rules or groups holds.", + "properties": { + "any_of": { + "description": "Rules or groups. One must hold.", + "items": { + "oneOf": [ + { + "$ref": "#/components/schemas/AffinityRulePayload" + }, + { + "$ref": "#/components/schemas/AffinityAnyOfPayload" + }, + { + "$ref": "#/components/schemas/AffinityAllOfPayload" + } + ] + }, + "minItems": 1, + "title": "Any Of", + "type": "array" + } + }, + "required": [ + "any_of" + ], + "title": "AffinityAnyOfPayload", + "type": "object" + }, + "AffinityRulePayload": { + "additionalProperties": false, + "description": "One rule on the tags of the candidate Metal Server, or of a VM that runs on it.", + "properties": { + "operator": { + "description": "`has` needs a resource with every tag pair. `has_not` rejects such a resource.", + "enum": [ + "has", + "has_not" + ], + "title": "Operator", + "type": "string" + }, + "resource": { + "description": "`metal_server` checks the candidate host. `virtual_machine` checks the VMs on the candidate host.", + "enum": [ + "metal_server", + "virtual_machine" + ], + "title": "Resource", + "type": "string" + }, + "tags": { + "additionalProperties": { + "maxLength": 1000, + "type": "string" + }, + "description": "Tag pairs that must all be on one resource.", + "maxProperties": 32, + "minProperties": 1, + "propertyNames": { + "maxLength": 128 + }, + "title": "Tags", + "type": "object" + }, + "within": { + "anyOf": [ + { + "maxLength": 128, + "minLength": 1, + "type": "string" + }, + { + "type": "null" + } + ], + "default": null, + "description": "A host tag key, such as `rack`. A `virtual_machine` rule then reads the VMs on every host that has the same value for this key as the candidate host. A host without the key fails the rule.", + "title": "Within" + } + }, + "required": [ + "resource", + "operator", + "tags" + ], + "title": "AffinityRulePayload", + "type": "object" + }, + "AffinityUnsatisfiedError": { + "description": "No host with room meets the affinity rules of the VM.", + "properties": { + "code": { + "const": "affinity_unsatisfied", + "description": "Stable machine-readable error code.", + "title": "Code", + "type": "string" + }, + "fields": { + "description": "Invalid request fields, or an empty list.", + "items": { + "$ref": "#/components/schemas/ApiErrorField" + }, + "title": "Fields", + "type": "array" + }, + "message": { + "description": "Safe description of the failure.", + "title": "Message", + "type": "string" + } + }, + "required": [ + "code", + "message", + "fields" + ], + "title": "AffinityUnsatisfiedError", + "type": "object" + }, "ApiErrorDetail": { "description": "The stable error details returned by Atlas.", "properties": { @@ -73,6 +224,7 @@ "description": "Capacity failure details.", "discriminator": { "mapping": { + "affinity_unsatisfied": "#/components/schemas/AffinityUnsatisfiedError", "out_of_capacity": "#/components/schemas/OutOfCapacityError", "placement_busy": "#/components/schemas/PlacementBusyError" }, @@ -84,6 +236,9 @@ }, { "$ref": "#/components/schemas/PlacementBusyError" + }, + { + "$ref": "#/components/schemas/AffinityUnsatisfiedError" } ], "title": "Error" @@ -269,6 +424,24 @@ "title": "Metadata", "type": "object" }, + "placement_rules": { + "description": "Rules that limit the Metal Servers for the virtual machine. Every listed rule or group must hold. Placement uses only the hosts that meet them.", + "items": { + "oneOf": [ + { + "$ref": "#/components/schemas/AffinityRulePayload" + }, + { + "$ref": "#/components/schemas/AffinityAnyOfPayload" + }, + { + "$ref": "#/components/schemas/AffinityAllOfPayload" + } + ] + }, + "title": "Placement Rules", + "type": "array" + }, "private_network_throughput_mibps": { "default": 0, "description": "Private network throughput limit in MiB/s. Zero removes the limit.", @@ -325,6 +498,19 @@ "title": "Ssh Keys", "type": "array" }, + "tags": { + "additionalProperties": { + "maxLength": 1000, + "type": "string" + }, + "description": "Resource tags as key-value pairs.", + "maxProperties": 32, + "propertyNames": { + "maxLength": 128 + }, + "title": "Tags", + "type": "object" + }, "user_data": { "default": "", "description": "Cloud-init user data supplied to the guest.", @@ -3421,7 +3607,7 @@ ] }, "post": { - "description": "Creates a tenant VM from an image and requests the specified compute, disk, network, and guest configuration. Only tenant 0 can set `is_privileged`, which lets the VM reach every tenant through the mesh.\n\nSet `is_termination_protected` to refuse deletion of the new VM. The termination protection route changes it later.", + "description": "Creates a tenant VM from an image and requests the specified compute, disk, network, and guest configuration. Only tenant 0 can set `is_privileged`, which lets the VM reach every tenant through the mesh.\n\nSet `is_termination_protected` to refuse deletion of the new VM. The termination protection route changes it later.\n\nUse `tags` to label the VM, for example, `{\"role\": \"cargo-server\"}`. `placement_rules` limits the Metal Servers for the VM by host tags and by the tags of other VMs on the host. Atlas Settings selects whether a VM that no host with room can satisfy fails with `affinity_unsatisfied` or goes to any host.", "operationId": "create_virtual_machine", "parameters": [ { @@ -3515,7 +3701,7 @@ } } }, - "description": "The virtual machine was not placed. `error.code` is `out_of_capacity` when no host can hold it, which needs more capacity in the region, or `placement_busy` when every candidate host was held by another placement, which only needs a retry. A busy response carries `Retry-After` in seconds.", + "description": "The virtual machine was not placed. `error.code` is `out_of_capacity` when no host can hold it, which needs more capacity in the region, or `placement_busy` when every candidate host was held by another placement, which only needs a retry. A busy response carries `Retry-After` in seconds. `affinity_unsatisfied` means that no host with room meets the affinity rules of the VM.", "headers": { "Retry-After": { "description": "Seconds to wait before retrying a `placement_busy` response.", diff --git a/docs/compute/placement.md b/docs/compute/placement.md index 950893075..4cf37abf3 100644 --- a/docs/compute/placement.md +++ b/docs/compute/placement.md @@ -39,6 +39,118 @@ For example, if host A already runs 2 of the tenant's VMs and host B runs 1, hos To add a strategy, subclass `PlacementStrategy` and register it. The [VM module specification](../../atlas/vm/SPEC.md#placement) shows how. +## Affinity rules + +Affinity rules limit the hosts that can hold a VM. A rule reads the tags of a host (Metal Server), or the tags of the VMs that run on the host. + +Atlas validates the rules in a create request and stores them in the `placement_rules` field of the Virtual Machine record. Placement applies them when it creates the VM, when it migrates the VM to a host that it chooses, and when a resize moves the VM to another host. A resize that stays on the current host does not check them. A migration to a host that an operator names does not check them. + +Any tenant can set rules on its own VMs. + +### Rule types + +```text +placement_rules = [node, ...] every node must hold (AND); absent or [] = no rules +node = rule | any_of | all_of +any_of = {"any_of": [node, ...]} at least one node holds (OR) +all_of = {"all_of": [node, ...]} every node holds (AND) +rule = {"resource": resource, "operator": operator, "tags": {key: value, ...}, "within": key} +resource = "metal_server" | "virtual_machine" +operator = "has" | "has_not" +``` + +| Field | Values | Meaning | +| --- | --- | --- | +| `resource` | `metal_server`, `virtual_machine` | `metal_server` reads the tags of the host. `virtual_machine` reads the tags of each VM on the host. | +| `operator` | `has`, `has_not` | `has` needs a resource with every pair in `tags`. `has_not` is the opposite of `has`. | +| `tags` | 1 to 32 key and value pairs | Every pair must be on the same resource. Atlas trims each key and value. A key has 1 to 128 characters. A value has at most 1000 characters. | +| `within` | A host tag key, such as `rack`. Optional. | Only for `virtual_machine`. The rule reads the VMs on every host that has the same value for this key as the candidate host. Without `within`, it reads only the VMs on the candidate host. | +| `any_of` | 1 or more nodes | The group holds when at least one node holds. | +| `all_of` | 1 or more nodes | The group holds when every node holds. | + +A rule asks for one of these conditions on a candidate host: + +| `resource` | `operator` | The host matches when | +| --- | --- | --- | +| `metal_server` | `has` | The host has every pair. | +| `metal_server` | `has_not` | The host is missing at least one pair. | +| `virtual_machine` | `has` | At least one VM on the host has every pair. | +| `virtual_machine` | `has_not` | No VM on the host has every pair. | + +For `virtual_machine`, one rule with two pairs needs one VM that has both pairs. An `all_of` group of two rules accepts two different VMs. + +A `virtual_machine` rule counts only the VMs of the same tenant. It counts a VM on its host, also when the VM is a draft, and on the destination host of its active migration. It does not count a VM that is being terminated, or the VM that placement moves. + +With `within`, the hosts that share the candidate host's value for the key are its **group**. For example, with `"within": "rack"` and a candidate host tagged `rack: a`, the group is every host tagged `rack: a`, whatever its status. `has` then needs a VM with every pair somewhere in the group, and `has_not` needs no such VM in the group. A candidate host without the key fails the rule for both operators, because Atlas cannot tell which group it is in. + +### Examples + +This create request tags a Cargo Server VM and asks for a host that has no other Cargo Server VM: + +```json +{ + "tags": {"role": "cargo-server"}, + "placement_rules": [ + {"resource": "virtual_machine", "operator": "has_not", "tags": {"role": "cargo-server"}} + ] +} +``` + +This request asks for a storage-optimised host in rack `a`, or for any memory-optimised host: + +```json +{ + "placement_rules": [ + {"any_of": [ + {"resource": "metal_server", "operator": "has", "tags": {"type": "storage-optimised", "rack": "a"}}, + {"resource": "metal_server", "operator": "has", "tags": {"type": "memory-optimised"}} + ]} + ] +} +``` + +This request keeps a database replica out of every rack that already has one: + +```json +{ + "tags": {"role": "db-replica"}, + "placement_rules": [ + {"resource": "virtual_machine", "operator": "has_not", "tags": {"role": "db-replica"}, "within": "rack"} + ] +} +``` + +If a replica runs on host 1 in rack `a`, host 3 in rack `a` fails the rule, although host 3 runs no replica itself. Hosts in rack `b` pass. A host without a `rack` tag fails. + +### Validation + +Atlas rejects the request with `400` when one of these is true: + +- A rule has an unknown field, for example `weight`, or a rule has no `resource`, `operator`, or `tags`. +- A `metal_server` rule has `within`, or `within` is empty or longer than 128 characters. +- `resource` or `operator` has an unknown value. +- `tags` is empty, a tag breaks the limits above, or two keys are the same after Atlas trims them. +- An `any_of` or `all_of` group is empty, or a group object has more than one key. +- The request has more than 16 rules in total, including the rules in groups. +- Groups nest more than 5 deep. + +### How placement applies the rules + +The rules are a filter before the strategy. When at least one host that meets the rules has room, placement removes the other hosts from its host pool. Then the strategy ranks the remaining hosts as usual. + +The host pool can be up to 1 second old, so two placements at the same time can see the same hosts. For this reason, Atlas checks the rules again while it holds the host lock, against the committed VMs on the host. If two Cargo Server VMs with `has_not {"role": "cargo-server"}` arrive together, the second one reads the draft of the first one. Placement then treats that host like a locked host and tries again with a fresh host pool, so the second VM goes to another host. + +A rule with `within` reads the VMs on every host of the group. Thus, placement also locks each other host of the group before it checks the rule, and keeps those locks until the VM record commits. It does not wait for these locks, so two placements that need the same group cannot deadlock. If another placement holds one of them, placement treats the candidate host as locked and tries again. + +The **Affinity Matching** setting in Atlas Settings, on the VM Scheduler tab, selects what happens when no host that meets the rules has room: + +| Value | No host that meets the rules has room | Every host that meets the rules is locked | +| --- | --- | --- | +| Enforced (default) | The request fails with `affinity_unsatisfied` (`503`). Atlas does not add a host. | `placement_busy` | +| Preferred | Placement ignores the rules and uses every host. | `placement_busy` | + +A locked host is not proof that no host meets the rules. Thus, Preferred matching does not ignore the rules when the hosts are only locked. + ## Under contention Concurrent requests rank hosts the same way, so they would all wait on the same top host. If another request holds that host's lock, Atlas tries the remaining hosts in random order. @@ -65,8 +177,10 @@ Offline and live simulators compare strategies outside the request path. Their r - [Placement context](../../atlas/vm/core/placement/context.py) loads samples, subtracts reservations, and rechecks capacity under the host lock. - [Placement strategy base](../../atlas/vm/core/placement/strategies/base.py) ranks candidate hosts and defines capacity and busy results. +- [Affinity rules](../../atlas/vm/core/placement/affinity.py) parse, validate, and serialize the rule tree. - [Transaction setup](../../atlas/vm/core/placement/transaction.py) enables READ COMMITTED. - [VM creation](../../atlas/vm/core/vm_service.py) commits the draft after placement. - [Placement tests](../../atlas/vm/core/placement/test_context.py) check capacity and lock behavior. +- [Affinity tests](../../atlas/vm/core/placement/test_affinity.py) check the rule shape and each rejected input. ::: diff --git a/docs/interfaces/tenant-api.md b/docs/interfaces/tenant-api.md index 50f6ef13c..b07415188 100644 --- a/docs/interfaces/tenant-api.md +++ b/docs/interfaces/tenant-api.md @@ -50,6 +50,7 @@ A create or resize that cannot place the VM returns `503`: | --- | --- | | `out_of_capacity` | Retry later when capacity is available. | | `placement_busy` | Retry after the response's `Retry-After` interval. | +| `affinity_unsatisfied` | No host with room meets the affinity rules of the VM. Change the rules, or add a host that meets them. | See [host selection](../compute/placement.md) for capacity rules.