Commit bec5b473 authored by Dimitrios Gogos's avatar Dimitrios Gogos
Browse files

feat: add GPU and storage management to application resources and update related schemas

parent feba9f09
Loading
Loading
Loading
Loading
Loading
+1 −1
Original line number Diff line number Diff line
@@ -23,4 +23,4 @@ repos:
    - id: mypy
      files: ^src/open_exposure_gateway/|^tests/
      args: [--strict, --ignore-missing-imports, --cache-dir, .cache/mypy]
      additional_dependencies: ["fastapi[standard]>=0.135.1", "pydantic>=2.0", "pydantic-settings>=2.0", "httpx>=0.27", "pytest>=9.0.2", "sqlalchemy>=2.0.48", "pytest-asyncio>=0.24", "testcontainers>=4.0.0"] #TODO, always add necessary
      additional_dependencies: ["fastapi[standard]>=0.135.1", "pydantic>=2.0", "pydantic-settings>=2.0", "httpx>=0.27", "pytest>=9.0.2", "sqlalchemy>=2.0.48", "pytest-asyncio>=0.24", "testcontainers>=4.0.0", "types-PyYAML>=6.0.1"]
+1 −0
Original line number Diff line number Diff line
@@ -36,6 +36,7 @@ dev = [
    "ruff>=0.15.6",
    "schemathesis>=4.22",
    "testcontainers>=4.0.0",
    "types-PyYAML>=6.0.1",
]

[tool.setuptools.packages.find]
+59 −1
Original line number Diff line number Diff line
@@ -3,7 +3,7 @@ from enum import StrEnum
from typing import Any, Literal, Optional
from uuid import UUID

from pydantic import BaseModel, ConfigDict, Field
from pydantic import BaseModel, ConfigDict, Field, field_validator


class AppInstanceStatus(StrEnum):
@@ -77,12 +77,32 @@ class ComponentSpecItem(BaseModel):
    networkInterfaces: list[NetworkInterface]


class GpuInfo(BaseModel):
    model_config = ConfigDict(extra="forbid")

    gpuMemory: int = Field(ge=0, le=16384)
    numGPU: int = Field(ge=0, le=16)


class AdditionalStorageItem(BaseModel):
    model_config = ConfigDict(extra="forbid")

    name: Optional[str] = Field(default=None, max_length=64)
    storageSize: str = Field(max_length=32, pattern=r"^\d+(GB|MB)$")
    mountPoint: str = Field(max_length=64)


AdditionalStorage = list[AdditionalStorageItem]


class VmResources(BaseModel):
    model_config = ConfigDict(extra="forbid")

    infraKind: Literal["virtualMachine"]
    numCPU: int = Field(ge=1, le=256)
    memory: int = Field(ge=1, le=32768)
    additionalStorages: Optional[AdditionalStorage] = Field(default=None, max_length=50)
    gpu: Optional[GpuInfo] = None


class ContainerResources(BaseModel):
@@ -91,6 +111,8 @@ class ContainerResources(BaseModel):
    infraKind: Literal["container"]
    numCPU: str = Field(pattern=r"^\d+((\.\d{1,3})|(m))?$")
    memory: int = Field(ge=1, le=16384)
    storage: Optional[AdditionalStorage] = Field(default=None, max_length=50)
    gpu: Optional[GpuInfo] = None


class DockerComposeResources(BaseModel):
@@ -99,6 +121,8 @@ class DockerComposeResources(BaseModel):
    infraKind: Literal["dockerCompose"]
    numCPU: int = Field(ge=1, le=256)
    memory: int = Field(ge=1, le=16384)
    storage: Optional[AdditionalStorage] = Field(default=None, max_length=50)
    gpu: Optional[GpuInfo] = None


class CpuPoolTopology(BaseModel):
@@ -142,6 +166,30 @@ class ApplicationResources(BaseModel):
    gpuPool: Optional[GpuPool] = None


class K8sPrimaryNetwork(BaseModel):
    model_config = ConfigDict(extra="forbid")

    provider: Optional[str] = Field(default=None, max_length=64)
    version: Optional[str] = Field(default=None, max_length=64)


class K8sAdditionalNetwork(BaseModel):
    model_config = ConfigDict(extra="forbid")

    name: Optional[str] = Field(default=None, max_length=64)
    interfaceType: Optional[Literal["netdevice", "vfio-pci", "interface"]] = None


class K8sNetworking(BaseModel):
    model_config = ConfigDict(extra="forbid")

    primaryNetwork: K8sPrimaryNetwork
    additionalNetworks: Optional[list[K8sAdditionalNetwork]] = Field(default=None, max_length=100)


K8sAddon = Literal["monitoring", "ingress"]


class KubernetesResources(BaseModel):
    model_config = ConfigDict(extra="forbid")

@@ -149,6 +197,16 @@ class KubernetesResources(BaseModel):
    applicationResources: ApplicationResources
    isStandalone: bool
    additionalStorage: Optional[str] = Field(default=None, max_length=32, pattern=r"^\d+(GB|MB)$")
    version: Optional[str] = Field(default=None, max_length=64)
    networking: Optional[K8sNetworking] = None
    addons: Optional[list[K8sAddon]] = Field(default=None, max_length=2)

    @field_validator("addons")
    @classmethod
    def _addons_unique(cls, v: Optional[list[str]]) -> Optional[list[str]]:
        if v is not None and len(set(v)) != len(v):
            raise ValueError("addons must be unique")
        return v


RequiredResources = VmResources | ContainerResources | DockerComposeResources | KubernetesResources
+114 −15
Original line number Diff line number Diff line
@@ -19,12 +19,18 @@ from open_exposure_gateway.api.camara.edge_application_management.vwip.schemas i
    SubmittedApp,
    VmResources,
)
from open_exposure_gateway.api.camara.edge_application_management.vwip.schemas import (
    AdditionalStorageItem as CamaraAdditionalStorageItem,
)
from open_exposure_gateway.api.camara.edge_application_management.vwip.schemas import (
    ApplicationResources as CamaraApplicationResources,
)
from open_exposure_gateway.api.camara.edge_application_management.vwip.schemas import (
    AppRepo as CamaraAppRepo,
)
from open_exposure_gateway.api.camara.edge_application_management.vwip.schemas import (
    GpuInfo as CamaraGpuInfo,
)
from open_exposure_gateway.api.camara.edge_application_management.vwip.schemas import (
    NetworkInterface as CamaraNetworkInterface,
)
@@ -39,6 +45,8 @@ from open_exposure_gateway.domain.edge_application_management import (
    CpuPool,
    CpuPoolTopology,
    GpuPool,
    GpuRequest,
    K8sClusterConfig,
    NetworkInterface,
    RequiredResources,
    SRMAccelerator,
@@ -61,6 +69,7 @@ from open_exposure_gateway.domain.edge_application_management import (
    SRMTerminatePayload,
    SRMTopologyConstraints,
    SRMZone,
    StorageRequest,
)
from open_exposure_gateway.domain.models import AppInstance, AppInstanceState

@@ -125,12 +134,65 @@ def _parse_storage_mb(value: str) -> int:
        raise ValueError(f"Cannot parse storage value: {value!r}")
    amount, unit = float(match.group(1)), match.group(2).upper()
    if unit == "TB":
        return int(amount * 1024 * 1024)
        return int(amount * 1000 * 1000)
    if unit == "GB":
        return int(amount * 1024)
        return int(amount * 1000)
    return int(amount)


def _format_storage_size(size_mb: int) -> str:
    """Inverse of `_parse_storage_mb` — whole decimal GB where it divides evenly."""
    if size_mb % 1000 == 0:
        return f"{size_mb // 1000}GB"
    return f"{size_mb}MB"


def _build_camara_gpu(compute: Optional[SRMComputeResources]) -> Optional[CamaraGpuInfo]:
    """SRM accelerator -> CAMARA `GpuInfo` (megabytes on both sides)."""
    if compute is None or compute.accelerator is None:
        return None
    acc = compute.accelerator
    if acc.type != "gpu":
        return None
    return CamaraGpuInfo(gpuMemory=acc.memory_mb, numGPU=acc.units)


def _build_camara_storage(
    compute: Optional[SRMComputeResources],
) -> Optional[list[CamaraAdditionalStorageItem]]:
    """SRM storage -> CAMARA `AdditionalStorage` (megabytes on both sides)."""
    if compute is None or not compute.storage:
        return None
    items = [
        CamaraAdditionalStorageItem(
            name=v.name,
            storageSize=_format_storage_size(v.size_mb),
            mountPoint=v.mount_point,
        )
        for v in compute.storage
        if v.mount_point
    ]
    return items or None


def _build_gpu_request(gpu: Optional[CamaraGpuInfo]) -> Optional[GpuRequest]:
    """CAMARA `GpuInfo` -> internal GPU request. Already megabytes; no conversion."""
    if gpu is None:
        return None
    return GpuRequest(num_gpu=gpu.numGPU, gpu_memory_mb=gpu.gpuMemory)


def _build_storage_requests(
    storages: Optional[list[CamaraAdditionalStorageItem]],
) -> list[StorageRequest]:
    """CAMARA `AdditionalStorage` -> internal storage requests."""
    if not storages:
        return []
    return [
        StorageRequest(name=s.name, size=s.storageSize, mount_point=s.mountPoint) for s in storages
    ]


def build_edge_cloud_zone(srm_zone: SRMZone) -> EdgeCloudZone:
    try:
        status = EdgeCloudZoneStatus(srm_zone.state)
@@ -154,9 +216,9 @@ def build_app_manifest(catalog: SRMCatalogPayload) -> AppManifest:
        None,
    )
    if unit is None:
        raise ValueError(
            f"no deployment unit in catalog entry {spec.ref!r} has an artifact_ref"
        )
        raise ValueError(f"no deployment unit in catalog entry {spec.ref!r} has an artifact_ref")
    artifact_ref = unit.artifact_ref
    assert artifact_ref is not None

    try:
        app_id: Optional[UUID] = UUID(spec.ref)
@@ -169,7 +231,7 @@ def build_app_manifest(catalog: SRMCatalogPayload) -> AppManifest:
    if repo_meta and repo_meta.type == "PRIVATEREPO":
        app_repo = CamaraAppRepo(
            type="PRIVATEREPO",
            imagePath=unit.artifact_ref,
            imagePath=artifact_ref,
            userName=repo_meta.user_ref,
            credentials=repo_meta.credentials,
            authType=repo_meta.auth_type,  # type: ignore[arg-type]
@@ -177,7 +239,7 @@ def build_app_manifest(catalog: SRMCatalogPayload) -> AppManifest:
    else:
        app_repo = CamaraAppRepo(
            type="PUBLICREPO",
            imagePath=unit.artifact_ref,
            imagePath=artifact_ref,
        )

    compute = unit.resource_requirements.compute
@@ -241,6 +303,8 @@ def build_app_manifest(catalog: SRMCatalogPayload) -> AppManifest:
            infraKind="container",
            numCPU=num_cpu_str,
            memory=compute.memory_mb if compute and compute.memory_mb else 0,
            storage=_build_camara_storage(compute),
            gpu=_build_camara_gpu(compute),
        )
    elif unit.runtime_kind in ("qcow2", "ova"):
        required_resources = VmResources(
@@ -249,6 +313,8 @@ def build_app_manifest(catalog: SRMCatalogPayload) -> AppManifest:
            if compute and compute.cpu_millicores
            else 1,
            memory=compute.memory_mb if compute and compute.memory_mb else 1,
            additionalStorages=_build_camara_storage(compute),
            gpu=_build_camara_gpu(compute),
        )
    elif unit.runtime_kind == "docker-compose":
        required_resources = DockerComposeResources(
@@ -257,6 +323,8 @@ def build_app_manifest(catalog: SRMCatalogPayload) -> AppManifest:
            if compute and compute.cpu_millicores
            else 1,
            memory=compute.memory_mb if compute and compute.memory_mb else 1,
            storage=_build_camara_storage(compute),
            gpu=_build_camara_gpu(compute),
        )
    else:
        required_resources = None
@@ -391,33 +459,48 @@ def build_app_registration_translation(
                ),
            )

        cluster_config = K8sClusterConfig(
            version=rr.version,
            networking=rr.networking.model_dump(exclude_none=True) if rr.networking else None,
            addons=list(rr.addons) if rr.addons else None,
        )
        required_resources = RequiredResources(
            infra_kind=rr.infraKind,
            is_standalone=rr.isStandalone or False,
            application_resources=ApplicationResources(cpu_pool=cpu_pool, gpu_pool=gpu_pool),
            additional_storage=rr.additionalStorage,
            k8s_cluster_config=(
                cluster_config if cluster_config.model_dump(exclude_none=True) else None
            ),
        )
    elif isinstance(manifest.requiredResources, (VmResources, DockerComposeResources)):
        vm_rr = manifest.requiredResources
        raw_storage = vm_rr.additionalStorages if isinstance(vm_rr, VmResources) else vm_rr.storage
        required_resources = RequiredResources(
            infra_kind=manifest.requiredResources.infraKind,
            infra_kind=vm_rr.infraKind,
            is_standalone=False,
            application_resources=ApplicationResources(
                cpu_pool=CpuPool(
                    num_cpu=float(manifest.requiredResources.numCPU),
                    memory=manifest.requiredResources.memory,
                    num_cpu=float(vm_rr.numCPU),
                    memory=vm_rr.memory,
                )
            ),
            gpu=_build_gpu_request(vm_rr.gpu),
            additional_storages=_build_storage_requests(raw_storage),
        )
    elif isinstance(manifest.requiredResources, ContainerResources):
        ctr_rr = manifest.requiredResources
        required_resources = RequiredResources(
            infra_kind=manifest.requiredResources.infraKind,
            infra_kind=ctr_rr.infraKind,
            is_standalone=False,
            application_resources=ApplicationResources(
                cpu_pool=CpuPool(
                    num_cpu=_parse_container_cpu_cores(manifest.requiredResources.numCPU),
                    memory=manifest.requiredResources.memory,
                    num_cpu=_parse_container_cpu_cores(ctr_rr.numCPU),
                    memory=ctr_rr.memory,
                )
            ),
            gpu=_build_gpu_request(ctr_rr.gpu),
            additional_storages=_build_storage_requests(ctr_rr.storage),
        )

    component_spec = [
@@ -506,9 +589,25 @@ def build_catalog_payload(translation: AppRegistrationTranslation) -> SRMCatalog
                    min_node_gpu_memory_mb=min_node_gpu_memory_mb,
                )

    if rr and rr.additional_storage:
    if rr and rr.gpu is not None and accelerator is None:
        accelerator = SRMAccelerator(
            type="gpu",
            units=rr.gpu.num_gpu,
            memory_mb=rr.gpu.gpu_memory_mb,
        )

    if rr and rr.additional_storages:
        storage = [
            SRMStorageVolume(
                name=s.name or "additional",
                size_mb=_parse_storage_mb(s.size),
                mount_point=s.mount_point,
            )
            for s in rr.additional_storages
        ]
    elif rr and rr.additional_storage:
        size_mb = _parse_storage_mb(rr.additional_storage)
        storage = [SRMStorageVolume(name="additional", size_mb=size_mb)]
        storage = [SRMStorageVolume(name="data", size_mb=size_mb, mount_point="/data")]

    compute = SRMComputeResources(
        cpu_millicores=cpu_millicores,
+20 −0
Original line number Diff line number Diff line
@@ -60,11 +60,31 @@ class ApplicationResources(BaseModel):
    gpu_pool: GpuPool | None = None


class GpuRequest(BaseModel):
    num_gpu: int
    gpu_memory_mb: int


class StorageRequest(BaseModel):
    name: str | None = None
    size: str
    mount_point: str


class K8sClusterConfig(BaseModel):
    version: str | None = None
    networking: Any | None = None
    addons: list[str] | None = None


class RequiredResources(BaseModel):
    infra_kind: str
    application_resources: ApplicationResources | None = None
    is_standalone: bool = False
    additional_storage: str | None = None
    gpu: GpuRequest | None = None
    additional_storages: list[StorageRequest] = []
    k8s_cluster_config: K8sClusterConfig | None = None


class AppRegistrationTranslation(BaseModel):
Loading