Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
16 commits
Select commit Hold shift + click to select a range
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
1 change: 1 addition & 0 deletions .vitepress/sidebar.mts
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@ const handbook: DefaultTheme.SidebarItem[] = [
{ text: 'Regional settings', link: '/docs/region/configuration' },
{ text: 'Provision a host', link: '/docs/region/' },
{ text: 'Metal daemon', link: '/docs/region/metald' },
{ text: 'Atlas access to hosts', link: '/docs/region/host-access' },
{ text: 'Host sync', link: '/docs/region/host-sync' },
{ text: 'Add a provider', link: '/docs/region/provider-guide' },
],
Expand Down
2 changes: 1 addition & 1 deletion atlas/atlas/SPEC.md
Original file line number Diff line number Diff line change
Expand Up @@ -46,7 +46,7 @@ Domain code reaches a provider only through `ServerProvider`. Implementations ar

## SSH tasks

- `target` is a Dynamic Link. The task connects to the target's `ssh_host`.
- `target` is a Dynamic Link. The task connects to the target's `ssh_host` through its `get_ssh_proxy_command()`. A VM uses its host as the proxy.
- A task stores no credentials. It uses the Atlas host identity.
- `script` and `environment` are plain text. Use `SSHRunner` directly for secret data.

Expand Down
20 changes: 15 additions & 5 deletions atlas/atlas/core/artifacts.py
Original file line number Diff line number Diff line change
Expand Up @@ -82,11 +82,21 @@ def get_linked_files() -> set[str]:


def get_download_url(file_name: str) -> str:
"""Return the URL a host uses to download one published File.
"""Return the URL a Metal host or a tenant-0 VM uses to download one published File."""
return get_file_url(file_name, get_internal_base_url())

`atlas_base_url` in the site configuration names an address that a host can
reach, which the site's own URL is not during local development.
"""

def get_internal_base_url() -> str:
"""Return `atlas_internal_url`, which hosts and tenant-0 VMs reach, or the public URL."""
return frappe.conf.atlas_internal_url or get_public_base_url()


def get_public_base_url() -> str:
"""Return `atlas_base_url`, which names a public address when the site URL is not one."""
return frappe.conf.atlas_base_url or frappe.utils.get_url(allow_header_override=False)


def get_file_url(file_name: str, base_url: str) -> str:
"""Return the URL of one File below a base URL."""
file_url = frappe.db.get_value("File", file_name, "file_url")
base_url = frappe.conf.atlas_base_url or frappe.utils.get_url(allow_header_override=False)
return f"{base_url.rstrip('/')}{file_url}"
43 changes: 37 additions & 6 deletions atlas/atlas/core/mesh_address.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@

import ipaddress
from typing import TYPE_CHECKING, cast
from uuid import UUID

import frappe
from frappe import _
Expand All @@ -14,31 +15,61 @@
MESH_PREFIX = 0xFDAA
MESH_NETWORK = ipaddress.IPv6Network((MESH_PREFIX << 112, 16))
MAXIMUM_VIRTUAL_MACHINE_NUMBER = 0xFFFFFFFFFFFFFFFF
# Atlas takes the last tenant-0 VM number, so no VM can use its address.
ATLAS_VIRTUAL_MACHINE_NUMBER = MAXIMUM_VIRTUAL_MACHINE_NUMBER
WIREGUARD_PREFIX = 0xFDAB
# The prefix and the region take the first 32 bits, so the low 96 bits of the UUID are the peer part.
WIREGUARD_PEER_MASK = (1 << 96) - 1


def get_virtual_machine_mesh_address(virtual_machine: Document | frappe._dict) -> str:
"""Return the mesh address from stable Atlas request metadata."""
settings = cast("AtlasSettings", frappe.get_single("Atlas Settings"))
region_id = settings.region_id
def validate_region_id(region_id: int) -> None:
"""Refuse a region ID that does not fit in one IPv6 field."""
if not 0 <= region_id <= 0xFFFF:
frappe.throw(_("Atlas Settings region ID must be a 16-bit unsigned integer."))


def get_virtual_machine_mesh_address(
virtual_machine: Document | frappe._dict, region_id: int | None = None
) -> str:
"""Return the mesh address from stable Atlas request metadata. A caller in a loop passes the region."""
if region_id is None:
region_id = cast("AtlasSettings", frappe.get_single("Atlas Settings")).region_id
validate_region_id(region_id)

virtual_machine_name = cast(str, virtual_machine.name)
virtual_machine_number = int(virtual_machine_name.rsplit("-", 1)[-1])
if virtual_machine_number > MAXIMUM_VIRTUAL_MACHINE_NUMBER:
frappe.throw(
_("Virtual Machine number {0} is too large for a mesh address.").format(virtual_machine_number)
)
if virtual_machine.tenant_id == 0 and virtual_machine_number == ATLAS_VIRTUAL_MACHINE_NUMBER:
frappe.throw(
_("Virtual Machine number {0} is the Atlas mesh address.").format(virtual_machine_number)
)

address = (
(MESH_PREFIX << 112) | (region_id << 96) | (virtual_machine.tenant_id << 64) | virtual_machine_number
)
return str(ipaddress.IPv6Address(address))


def get_atlas_mesh_address(region_id: int) -> str:
"""Return the tenant-0 mesh address of Atlas in one region."""
validate_region_id(region_id)

return str(ipaddress.IPv6Address((MESH_PREFIX << 112) | (region_id << 96) | ATLAS_VIRTUAL_MACHINE_NUMBER))


def get_region_mesh_address_prefix(region_id: int) -> str:
"""Return the leading hextets for VM mesh addresses in one region."""
if not 0 <= region_id <= 0xFFFF:
frappe.throw(_("Atlas Settings region ID must be a 16-bit unsigned integer."))
validate_region_id(region_id)

return f"{MESH_PREFIX:x}:{region_id:x}"


def get_wireguard_ip_address(peer_id: UUID, region_id: int) -> str:
"""Return the fdab::/16 wg0 address of one host or Atlas peer."""
validate_region_id(region_id)

address = (WIREGUARD_PREFIX << 112) | (region_id << 96) | (peer_id.int & WIREGUARD_PEER_MASK)
return str(ipaddress.IPv6Address(address))
46 changes: 33 additions & 13 deletions atlas/atlas/core/server_providers/aws/provider.py
Original file line number Diff line number Diff line change
Expand Up @@ -117,6 +117,38 @@ def ensure_server(self, request: ServerCreateRequest) -> ProviderServer:
"""Return the named AWS instance, and create it when necessary."""
return self.servers.ensure(request)

@override
def import_server(self, server: "MetalServer") -> None:
"""Match an existing instance to the catalog. Its storage volume must exist already."""
instance = self.servers.fetch(server.provider_server_id)
if not instance.get("ImageId"):
raise AwsError(f"AWS instance {server.provider_server_id} has no image")
# Public addresses attach to the primary interface, so it must be in the Atlas subnet.
if instance.get("SubnetId") != self.configuration.subnet_id:
raise AwsError(
f"AWS instance {server.provider_server_id} is in subnet {instance.get('SubnetId')}, "
f"not the Atlas subnet {self.configuration.subnet_id}"
)
instance_type = instance.get("InstanceType")
if not isinstance(instance_type, str) or not frappe.db.exists("Metal Server Size", instance_type):
raise AwsError(
f"No Metal Server Size matches {instance_type}. Sync the Metal Server Size catalog first."
)
server.server_size = instance_type
# The catalog keeps only the newest image of each version, so match the version.
images = self.client.call("ec2", "describe_images", ImageIds=[instance["ImageId"]]).get("Images", [])
versions = self.catalog.get_server_images(images)
if not versions or not frappe.db.exists("Metal Server Image", versions[0].name):
raise AwsError(
f"No Metal Server Image matches image {instance['ImageId']}. Sync the Metal Server Image catalog first."
)
server.server_image = versions[0].name
tags = {
tag.get("Key"): tag.get("Value") for tag in instance.get("Tags") or [] if isinstance(tag, Mapping)
}
server.title = tags.get("Name") or server.provider_server_id
self.apply_provider_server(server, self.servers.to_provider_server(instance))

@override
def prepare_server(self, server: "MetalServer") -> None:
"""Prepare the AWS instance before Secure Shell access."""
Expand All @@ -143,16 +175,6 @@ def configure_server_network(self, server: "MetalServer") -> None:
)
self.wait_for_private_address(server)

@override
def metald_listen_address(self, server: "MetalServer") -> str:
"""Return the primary interface address behind the internet gateway."""
metadata = frappe.parse_json(server.provider_metadata or "{}")
instance = metadata.get("instance") if isinstance(metadata, Mapping) else None
address = instance.get("PrivateIpAddress") if isinstance(instance, Mapping) else None
if not isinstance(address, str) or not address:
raise AwsError("Atlas server has no AWS primary private IPv4 address")
return address

@override
def storage_pool_device(self, server: "MetalServer") -> str:
"""Return the stable device path of the EBS volume for the storage pool."""
Expand Down Expand Up @@ -303,9 +325,7 @@ def uplink_interface(self, server: "MetalServer") -> str:
"""Return the guest device name that carries the AWS default route."""
from atlas.atlas.core.ssh import SSHRunner

result = SSHRunner(server.public_ipv4_address).run_command(
"ip -4 -o route show default", timeout_seconds=15
)
result = SSHRunner(server.ssh_host).run_command("ip -4 -o route show default", timeout_seconds=15)
fields = result.output.split()
if result.exit_code != 0 or "dev" not in fields:
raise AwsError(f"Atlas server {server.name} has no default route device")
Expand Down
57 changes: 37 additions & 20 deletions atlas/atlas/core/server_providers/aws/test_provider.py
Original file line number Diff line number Diff line change
Expand Up @@ -247,18 +247,6 @@ def test_configure_server_network_fails_without_a_mesh_mac_address(self) -> None
with self.assertRaises(AwsError):
provider.configure_server_network(server)

def test_metald_binds_the_primary_interface_not_the_mesh_interface(self) -> None:
provider = self.provider()
server = self.server()
server.provider_metadata = json.dumps(
{
"instance": {"PrivateIpAddress": "10.1.8.189", "PublicIpAddress": "56.155.92.65"},
"mesh_interface": {"PrivateIpAddress": "10.1.0.240"},
}
)

self.assertEqual(provider.metald_listen_address(server), "10.1.8.189")

def test_a_public_address_attaches_to_the_primary_interface(self) -> None:
provider = self.provider()
provider.ip_addresses = Mock()
Expand Down Expand Up @@ -294,14 +282,6 @@ def test_a_public_address_needs_the_primary_interface(self) -> None:
with self.assertRaisesRegex(AwsError, "primary network interface"):
provider.attach_public_ip_address("eipalloc-1", "203.0.113.9", server)

def test_metald_needs_the_primary_private_address(self) -> None:
provider = self.provider()
server = self.server()
server.provider_metadata = json.dumps({"instance": {"PublicIpAddress": "56.155.92.65"}})

with self.assertRaisesRegex(AwsError, "primary private IPv4 address"):
provider.metald_listen_address(server)

def test_the_storage_pool_device_names_the_attached_volume(self) -> None:
provider = self.provider()
server = self.server()
Expand Down Expand Up @@ -376,6 +356,42 @@ def test_uplink_interface_fails_without_a_default_route(self) -> None:
with patch("atlas.atlas.core.ssh.SSHRunner", return_value=runner), self.assertRaises(AwsError):
provider.uplink_interface(self.server())

def test_an_import_outside_the_atlas_subnet_is_refused(self) -> None:
"""Public addresses attach to the primary interface, so it must be in the Atlas subnet."""
provider = self.provider()
provider.configuration = SimpleNamespace(subnet_id="subnet-atlas")
provider.servers.fetch.return_value = {"ImageId": "ami-1", "SubnetId": "subnet-other"}

with self.assertRaisesRegex(AwsError, "not the Atlas subnet subnet-atlas"):
provider.import_server(self.server("i-1"))

def test_an_import_matches_an_older_image_of_a_catalog_version(self) -> None:
"""The catalog keeps only the newest image, and an instance can run an older one."""
provider = self.provider()
provider.configuration = SimpleNamespace(subnet_id="subnet-atlas")
provider.servers.fetch.return_value = {
"ImageId": "ami-old",
"SubnetId": "subnet-atlas",
"InstanceType": "c7i.xlarge",
"Tags": [{"Key": "Name", "Value": "osa-host"}],
}
provider.client.call.return_value = {
"Images": [
{
"ImageId": "ami-old",
"Name": "ubuntu/images/hvm-ssd-gp3/ubuntu-noble-24.04-amd64-server-20260904",
"CreationDate": "2026-09-04T11:45:55.000Z",
}
]
}
provider.apply_provider_server = Mock()
server = self.server("i-1")

with patch("atlas.atlas.core.server_providers.aws.provider.frappe.db.exists", return_value=True):
provider.import_server(server)

self.assertEqual((server.server_image, server.title), ("Ubuntu_24.04", "osa-host"))

def provider(self) -> AwsProvider:
provider = object.__new__(AwsProvider)
provider.settings = SimpleNamespace(
Expand All @@ -401,6 +417,7 @@ def server(provider_server_id: str | None = None) -> SimpleNamespace:
provider_server_id=provider_server_id,
provider_metadata="{}",
public_ipv4_address="203.0.113.1",
ssh_host="203.0.113.1",
private_ipv4_address=None,
public_network_interface=None,
private_network_interface=None,
Expand Down
2 changes: 1 addition & 1 deletion atlas/atlas/core/server_providers/aws/test_volumes.py
Original file line number Diff line number Diff line change
Expand Up @@ -73,7 +73,7 @@ def test_grow_runs_the_script_for_the_storage_pool(self) -> None:
volumes = self.volumes(
modifications=[{"StartTime": datetime.now(UTC), "ModificationState": "optimizing"}]
)
volumes.provider.storage_pool_device.return_value = "/dev/disk/by-id/pool"
volumes.provider.get_storage_pool_device.return_value = "/dev/disk/by-id/pool"
volumes.provider.poll.side_effect = lambda operation, **_: operation()
server = self.server()
task = SimpleNamespace(name="SSH-1", result=SimpleNamespace(is_success=True))
Expand Down
2 changes: 1 addition & 1 deletion atlas/atlas/core/server_providers/aws/volumes.py
Original file line number Diff line number Diff line change
Expand Up @@ -85,7 +85,7 @@ def grow(self, server: "MetalServer", kind: str) -> None:
)
environment = {"TARGET": kind}
if kind == "storage":
environment["STORAGE_POOL_DEVICE"] = self.provider.storage_pool_device(server)
environment["STORAGE_POOL_DEVICE"] = self.provider.get_storage_pool_device(server)
task = SSHTask.create_for_script_file(
target_type=server.doctype,
target=server.name,
Expand Down
31 changes: 19 additions & 12 deletions atlas/atlas/core/server_providers/base.py
Original file line number Diff line number Diff line change
Expand Up @@ -171,22 +171,17 @@ def configure_server_network(self, server: "MetalServer") -> None:
"""Configure the provider network after Secure Shell access is ready."""
...

def metald_listen_address(self, server: "MetalServer") -> str:
"""Return the configured metald address."""
address = (
server.public_ipv4_address
if self.settings.use_public_ip_for_metald
else server.private_ipv4_address
)
if not address:
raise self.error_class(f"Atlas server {server.name} has no address for metald")
return address

@abstractmethod
def storage_pool_device(self, server: "MetalServer") -> str:
"""Return the raw block device for the virtual machine storage pool."""
...

def get_storage_pool_device(self, server: "MetalServer") -> str:
"""Return the device that registration or import stored, else the provider device."""
metadata = frappe.parse_json(server.provider_metadata or "{}")
device = metadata.get("storage_pool_device") if isinstance(metadata, Mapping) else None
return device or self.storage_pool_device(server)

@abstractmethod
def set_power_state(self, provider_server_id: str, action: ServerPowerAction) -> None:
"""Apply one power action to a provider server."""
Expand All @@ -197,6 +192,18 @@ def delete_server(self, provider_server_id: str, provider_metadata: Mapping[str,
"""Delete a provider server and its owned resources if they exist."""
...

def import_server(self, server: "MetalServer") -> None:
"""Fill a Metal Server from a provider server that Atlas did not create."""
raise UnsupportedProviderOperation("server import")

def find_catalog_record(self, doctype: str, matches: Callable[[Mapping], bool], label: str) -> str:
"""Return the catalog record whose provider metadata matches an imported server."""
for name, metadata in frappe.get_all(doctype, fields=["name", "provider_metadata"], as_list=True):
values = frappe.parse_json(metadata or "{}")
if isinstance(values, Mapping) and matches(values):
return name
raise self.error_class(f"No {doctype} matches {label}. Sync the {doctype} catalog first.")

def reserve_public_ip_address(self, version: int) -> ReservedIPAddress:
"""Reserve one public IPv4 address or IPv6 block."""
raise UnsupportedProviderOperation(f"public IPv{version} address reservation")
Expand Down Expand Up @@ -233,7 +240,7 @@ def wait_for_private_address(self, server: "MetalServer") -> None:
def get_private_network_mac_address() -> str | None:
"""Return the interface MAC once the server has its private network address."""
try:
result = SSHRunner(server.public_ipv4_address).run_command(
result = SSHRunner(server.ssh_host).run_command(
f"ip -4 -o addr show dev {device} scope global && cat /sys/class/net/{device}/address",
timeout_seconds=15,
)
Expand Down
10 changes: 3 additions & 7 deletions atlas/atlas/core/server_providers/generic/provider.py
Original file line number Diff line number Diff line change
Expand Up @@ -47,7 +47,7 @@ def validate_settings(self) -> None:
@override
def validate_server(self, server: "MetalServer") -> None:
"""Require the storage pool device that host registration stores."""
self.storage_pool_device(server)
self.get_storage_pool_device(server)

@override
def validate_credentials(self) -> bool:
Expand Down Expand Up @@ -86,12 +86,8 @@ def configure_server_network(self, server: "MetalServer") -> None:

@override
def storage_pool_device(self, server: "MetalServer") -> str:
"""Return the storage pool device that host registration chose."""
metadata = frappe.parse_json(server.provider_metadata or "{}")
device = metadata.get("storage_pool_device") if isinstance(metadata, dict) else None
if not device:
raise GenericError(f"Metal Server {server.name} has no registered storage pool device")
return device
"""Refuse a host without the storage pool device that registration stores."""
raise GenericError(f"Metal Server {server.name} has no registered storage pool device")

@override
def set_power_state(self, provider_server_id: str, action: ServerPowerAction) -> None:
Expand Down
4 changes: 3 additions & 1 deletion atlas/atlas/core/server_providers/generic/test_provider.py
Original file line number Diff line number Diff line change
Expand Up @@ -44,7 +44,9 @@ def test_network_setup_only_checks_the_private_address(self) -> None:
run_setup_script.assert_not_called()

def test_storage_pool_device_is_the_registered_device(self) -> None:
self.assertEqual(GenericProvider(self.settings()).storage_pool_device(self.server()), "/dev/nvme1n1")
self.assertEqual(
GenericProvider(self.settings()).get_storage_pool_device(self.server()), "/dev/nvme1n1"
)

def test_public_address_attaches_as_itself(self) -> None:
provider = GenericProvider(self.settings())
Expand Down
Loading
Loading