Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
20 commits
Select commit Hold shift + click to select a range
225050c
feat(atlas): Metal Server tags
the-bokya Oct 1, 2026
02b4940
feat(atlas): Add an affinity field to VM creation
the-bokya Oct 1, 2026
160c76c
feat(atlas): Add the affinity rules model
the-bokya Oct 1, 2026
3ce3d17
feat(atlas): Accept affinity rules and tags on VM creation
the-bokya Oct 1, 2026
abdf9d3
test(atlas): Cover affinity rules and VM creation tags
the-bokya Oct 1, 2026
18e9219
docs(atlas): Describe the affinity rule types
the-bokya Oct 1, 2026
02f9199
feat(atlas): Filter hosts by affinity rules
the-bokya Oct 1, 2026
c9648a0
test(atlas): Cover affinity host filtering
the-bokya Oct 1, 2026
37fd51d
feat(atlas): Apply affinity rules during placement
the-bokya Oct 1, 2026
152011b
feat(atlas): Report affinity_unsatisfied from the VM API
the-bokya Oct 1, 2026
df65d0a
test(atlas): Cover affinity rules in placement
the-bokya Oct 1, 2026
4600a42
docs(atlas): Describe how placement applies affinity rules
the-bokya Oct 1, 2026
16a4d0d
fix(atlas): Ignore terminating VMs in affinity rules
the-bokya Oct 1, 2026
2cd934b
feat(atlas): Add within to affinity rules
the-bokya Oct 1, 2026
dd3e38a
feat(atlas): Accept within in VM affinity rules
the-bokya Oct 1, 2026
deaed80
test(atlas): Cover within and terminating VMs in affinity rules
the-bokya Oct 1, 2026
a0c0d0e
docs(atlas): Describe within and the VMs that affinity rules count
the-bokya Oct 1, 2026
02083cf
refactor(atlas): Rename affinity_rules to placement_rules
the-bokya Oct 1, 2026
89fe32c
feat(atlas): Allow placement rules for every tenant
the-bokya Oct 1, 2026
46267bf
Merge branch 'develop' into feat/affinity-rules
the-bokya Oct 1, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
87 changes: 85 additions & 2 deletions atlas/api/models.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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 (
Expand Down Expand Up @@ -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"),
]

Expand Down Expand Up @@ -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."""

Expand Down Expand Up @@ -451,13 +521,24 @@ 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:
if self.public_ipv4 and not self.ipv4_internet_access:
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(
Expand All @@ -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]),
)


Expand Down
5 changes: 4 additions & 1 deletion atlas/api/routes/virtual_machines.py
Original file line number Diff line number Diff line change
Expand Up @@ -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": {
Expand All @@ -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)
Expand Down
36 changes: 36 additions & 0 deletions atlas/api/tests/test_virtual_machines.py
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down Expand Up @@ -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)
Expand Down
12 changes: 11 additions & 1 deletion atlas/atlas/doctype/atlas_settings/atlas_settings.json
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down Expand Up @@ -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.",
Expand Down Expand Up @@ -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",
Expand Down
1 change: 1 addition & 0 deletions atlas/atlas/doctype/atlas_settings/atlas_settings.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
17 changes: 16 additions & 1 deletion atlas/metal_server/doctype/metal_server/metal_server.json
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,8 @@
"wireguard_public_key",
"section_break_metadata",
"provider_metadata",
"tags_section",
"tags",
"hidden_data_section",
"column_break_xhax",
"metald_tls_certificate",
Expand Down Expand Up @@ -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,
Expand All @@ -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",
Expand Down
6 changes: 5 additions & 1 deletion atlas/metal_server/doctype/metal_server/metal_server.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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"]
Expand All @@ -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
Expand Down Expand Up @@ -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:
Expand Down
18 changes: 18 additions & 0 deletions atlas/metal_server/doctype/metal_server/test_metal_server.py
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand Down
6 changes: 6 additions & 0 deletions atlas/vm/SPEC.md
Original file line number Diff line number Diff line change
Expand Up @@ -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. |
Expand All @@ -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.
Expand All @@ -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).
Expand Down
Loading
Loading