From 1ce001e0e25cafc2adbb42cccfd6d98cce89e37d Mon Sep 17 00:00:00 2001 From: stentoumis Date: Wed, 2 Sep 2026 14:20:55 +0300 Subject: [PATCH 1/3] feat: get zones internal endpoint --- src/srm/adapters/database/repos/topology.py | 16 ++- src/srm/api/dependencies.py | 16 +++ src/srm/api/rest/router.py | 42 ++++++- src/srm/api/rest/schemas.py | 7 ++ src/srm/application/use_cases/topology.py | 21 ++++ src/srm/config.py | 3 +- src/srm/domain/ports/database/topology.py | 10 ++ tests/api/rest/test_zone_endpoint.py | 120 ++++++++++++++++++++ tests/integration/test_repositories.py | 33 ++++++ tests/test_config.py | 18 +++ tests/unit/fakes.py | 18 ++- 11 files changed, 300 insertions(+), 4 deletions(-) create mode 100644 src/srm/application/use_cases/topology.py create mode 100644 tests/api/rest/test_zone_endpoint.py diff --git a/src/srm/adapters/database/repos/topology.py b/src/srm/adapters/database/repos/topology.py index d6e1a10..98eff7c 100644 --- a/src/srm/adapters/database/repos/topology.py +++ b/src/srm/adapters/database/repos/topology.py @@ -59,6 +59,14 @@ class SqlZoneRepository(ZoneRepository): return ZoneMapper.to_domain(row) if row is not None else None async def list_active_resource_zones(self) -> list[Zone]: + return await self.list_resource_zones(state=ZoneState.ACTIVE) + + async def list_resource_zones( + self, + *, + state: ZoneState | None = None, + region: str | None = None, + ) -> list[Zone]: stmt = ( select(ZoneRow) .options( @@ -66,9 +74,15 @@ class SqlZoneRepository(ZoneRepository): .selectinload(DomainRow.capabilities) .selectinload(CapabilityRow.control_path_bindings) ) - .where(ZoneRow.kind == ZoneKind.RESOURCE, ZoneRow.state == ZoneState.ACTIVE) + .where(ZoneRow.kind == ZoneKind.RESOURCE) .order_by(ZoneRow.ref, ZoneRow.id) ) + if state is not None: + stmt = stmt.where(ZoneRow.state == state) + if region is not None: + stmt = stmt.where( + ZoneRow.zone_metadata["location"]["region"].astext == region, + ) rows = (await self._session.scalars(stmt)).all() return [ZoneMapper.to_domain(row) for row in rows] diff --git a/src/srm/api/dependencies.py b/src/srm/api/dependencies.py index 0f52613..6319b7c 100644 --- a/src/srm/api/dependencies.py +++ b/src/srm/api/dependencies.py @@ -4,12 +4,15 @@ from fastapi import Depends, Request from sqlalchemy.ext.asyncio import AsyncSession from srm.adapters.database.repos.catalog import SqlServiceSpecificationRepository +from srm.adapters.database.repos.topology import SqlZoneRepository from srm.app_state import AppState from srm.application.use_cases.catalog import ( CreateServiceSpecificationUseCase, DeleteServiceSpecificationUseCase, GetServiceSpecificationUseCase, ) +from srm.application.use_cases.topology import ListZonesUseCase +from srm.config import get_settings def get_app_state(request: Request) -> AppState: @@ -47,6 +50,14 @@ def get_delete_service_specification_use_case( return DeleteServiceSpecificationUseCase(SqlServiceSpecificationRepository(session)) +def get_list_zones_use_case(session: SessionDep) -> ListZonesUseCase: + return ListZonesUseCase(SqlZoneRepository(session)) + + +def get_zone_provider_name() -> str: + return get_settings().zone_provider_name + + CreateServiceSpecificationUseCaseDep = Annotated[ CreateServiceSpecificationUseCase, Depends(get_create_service_specification_use_case), @@ -59,3 +70,8 @@ DeleteServiceSpecificationUseCaseDep = Annotated[ DeleteServiceSpecificationUseCase, Depends(get_delete_service_specification_use_case), ] +ListZonesUseCaseDep = Annotated[ + ListZonesUseCase, + Depends(get_list_zones_use_case), +] +ZoneProviderNameDep = Annotated[str, Depends(get_zone_provider_name)] diff --git a/src/srm/api/rest/router.py b/src/srm/api/rest/router.py index b26c30b..64874f5 100644 --- a/src/srm/api/rest/router.py +++ b/src/srm/api/rest/router.py @@ -1,13 +1,16 @@ +from typing import Annotated from uuid import UUID import structlog -from fastapi import APIRouter, HTTPException, status +from fastapi import APIRouter, Header, HTTPException, Response, status from srm.adapters.errors import DuplicateServiceSpecificationError, ServiceSpecificationInUseError from srm.api.dependencies import ( CreateServiceSpecificationUseCaseDep, DeleteServiceSpecificationUseCaseDep, GetServiceSpecificationUseCaseDep, + ListZonesUseCaseDep, + ZoneProviderNameDep, ) from srm.api.rest.schemas import ( CreateServiceSpecificationRequest, @@ -17,7 +20,9 @@ from srm.api.rest.schemas import ( ServiceCapabilityRequirementResponseSchema, ServiceDeploymentUnitCreateSchema, ServiceSpecificationCreateSchema, + ZoneResponse, ) +from srm.application.use_cases.topology import ListZonesCommand from srm.domain.models.canonical_parameters.parameters import ( CapabilityParameters, CapabilityTarget, @@ -30,6 +35,7 @@ from srm.domain.models.catalog import ( ServiceDeploymentUnit, ServiceSpecification, ) +from srm.domain.models.topology import Zone, ZoneState internal = APIRouter(prefix="/internal", tags=["Internal"]) logger: structlog.BoundLogger = structlog.get_logger(__name__) @@ -134,6 +140,22 @@ def _build_get_response(saved: ServiceSpecification) -> GetServiceSpecificationR ) +def _zone_metadata(zone: Zone, provider_name: str) -> dict[str, object]: + metadata = dict(zone.metadata) + if metadata.get("provider") is None: + metadata["provider"] = provider_name + return metadata + + +def _build_zone_response(zone: Zone, provider_name: str) -> ZoneResponse: + return ZoneResponse( + id=zone.id, + name=zone.name, + state=zone.state.value, + metadata=_zone_metadata(zone, provider_name), + ) + + @internal.post( "/catalog/service-specifications", status_code=status.HTTP_201_CREATED, @@ -201,6 +223,24 @@ async def get_service_specification( return _build_get_response(specification) +@internal.get("/zones") +async def list_zones( + use_case: ListZonesUseCaseDep, + provider_name: ZoneProviderNameDep, + response: Response, + state: ZoneState | None = None, + region: str | None = None, + x_correlator: Annotated[str | None, Header()] = None, +) -> list[ZoneResponse]: + logger.info("list_zones_requested", state=state, region=region) + + zones = await use_case.execute(ListZonesCommand(state=state, region=region)) + if x_correlator is not None: + response.headers["x-correlator"] = x_correlator + + return [_build_zone_response(zone, provider_name) for zone in zones] + + @internal.delete( "/catalog/service-specifications/{id}", status_code=status.HTTP_204_NO_CONTENT, diff --git a/src/srm/api/rest/schemas.py b/src/srm/api/rest/schemas.py index 12c4882..01c8048 100644 --- a/src/srm/api/rest/schemas.py +++ b/src/srm/api/rest/schemas.py @@ -87,5 +87,12 @@ class GetServiceSpecificationResponse(DomainModel): service_capability_requirements: list[ServiceCapabilityRequirementResponseSchema] +class ZoneResponse(DomainModel): + id: UUID + name: str + state: str + metadata: dict[str, Any] + + class ErrorResponse(DomainModel): detail: str diff --git a/src/srm/application/use_cases/topology.py b/src/srm/application/use_cases/topology.py new file mode 100644 index 0000000..96000c6 --- /dev/null +++ b/src/srm/application/use_cases/topology.py @@ -0,0 +1,21 @@ +from dataclasses import dataclass + +from srm.domain.models.topology import Zone, ZoneState +from srm.domain.ports.database.topology import ZoneRepository + + +@dataclass(frozen=True, slots=True) +class ListZonesCommand: + state: ZoneState | None = None + region: str | None = None + + +class ListZonesUseCase: + def __init__(self, zones: ZoneRepository) -> None: + self._zones = zones + + async def execute(self, command: ListZonesCommand) -> list[Zone]: + return await self._zones.list_resource_zones( + state=command.state, + region=command.region, + ) diff --git a/src/srm/config.py b/src/srm/config.py index 20a4d87..c2150f6 100644 --- a/src/srm/config.py +++ b/src/srm/config.py @@ -19,7 +19,7 @@ Usage: from functools import lru_cache from typing import Literal -from pydantic import BaseModel +from pydantic import BaseModel, Field from pydantic_settings import BaseSettings, SettingsConfigDict @@ -45,6 +45,7 @@ class Settings(BaseSettings): postgres_settings: PostgreSQLSettings nats_settings: NatsSettings + zone_provider_name: str = Field(default="oop-provider", min_length=1) @lru_cache() diff --git a/src/srm/domain/ports/database/topology.py b/src/srm/domain/ports/database/topology.py index af8fa61..9aeec6a 100644 --- a/src/srm/domain/ports/database/topology.py +++ b/src/srm/domain/ports/database/topology.py @@ -8,6 +8,7 @@ from srm.domain.models.topology import ( ControlPathBinding, Domain, Zone, + ZoneState, ) @@ -20,6 +21,15 @@ class ZoneRepository(ABC): async def list_active_resource_zones(self) -> list[Zone]: pass + @abstractmethod + async def list_resource_zones( + self, + *, + state: ZoneState | None = None, + region: str | None = None, + ) -> list[Zone]: + pass + @abstractmethod async def create(self, zone: Zone) -> Zone: pass diff --git a/tests/api/rest/test_zone_endpoint.py b/tests/api/rest/test_zone_endpoint.py new file mode 100644 index 0000000..e5c63d4 --- /dev/null +++ b/tests/api/rest/test_zone_endpoint.py @@ -0,0 +1,120 @@ +from uuid import UUID, uuid4 + +from fastapi import FastAPI +from httpx import ASGITransport, AsyncClient, Response + +from srm.api.dependencies import get_list_zones_use_case, get_zone_provider_name +from srm.api.rest.router import internal +from srm.application.use_cases.topology import ListZonesCommand +from srm.domain.models.topology import Domain, DomainKind, DomainState, Zone, ZoneKind, ZoneState + +ZONE_ID = UUID("642f6105-7015-4af1-a4d1-e1ecb8437abc") + + +def _domain(kind: DomainKind) -> Domain: + return Domain( + id=uuid4(), + zone_id=ZONE_ID, + ref=f"{kind.value}-internal", + name=f"{kind.value.title()} Domain", + kind=kind, + state=DomainState.ACTIVE, + metadata={"summary": kind.value}, + ) + + +def _zone(*, metadata: dict[str, object] | None = None) -> Zone: + return Zone( + id=ZONE_ID, + ref="zone-internal", + name="Berlin Edge 1", + kind=ZoneKind.RESOURCE, + state=ZoneState.ACTIVE, + metadata=( + metadata + if metadata is not None + else {"provider": "acme", "location": {"region": "eu-central-1"}} + ), + domains=[ + _domain(DomainKind.COMPUTE), + _domain(DomainKind.NETWORK), + _domain(DomainKind.OBSERVABILITY), + ], + ) + + +def _app(zones: list[Zone], captured: list[ListZonesCommand] | None = None) -> FastAPI: + class _UseCase: + async def execute(self, command: ListZonesCommand) -> list[Zone]: + if captured is not None: + captured.append(command) + return zones + + app = FastAPI() + app.include_router(internal) + + async def override_use_case() -> _UseCase: + return _UseCase() + + async def override_provider_name() -> str: + return "test-provider" + + app.dependency_overrides[get_list_zones_use_case] = override_use_case + app.dependency_overrides[get_zone_provider_name] = override_provider_name + return app + + +async def _get(app: FastAPI, **headers: str) -> Response: + async with AsyncClient(transport=ASGITransport(app=app), base_url="http://test") as client: + return await client.get("/internal/zones", headers=headers) + + +async def test_get_zones_returns_oeg_compatible_shape() -> None: + response = await _get(_app([_zone()])) + + assert response.status_code == 200 + assert response.json() == [ + { + "id": str(ZONE_ID), + "name": "Berlin Edge 1", + "state": "active", + "metadata": { + "provider": "acme", + "location": {"region": "eu-central-1"}, + }, + } + ] + + +async def test_get_zones_passes_filters_to_the_use_case() -> None: + captured: list[ListZonesCommand] = [] + + async with AsyncClient( + transport=ASGITransport(app=_app([], captured)), base_url="http://test" + ) as client: + response = await client.get( + "/internal/zones", + params={"state": "active", "region": "athens"}, + ) + + assert response.status_code == 200 + assert captured == [ListZonesCommand(state=ZoneState.ACTIVE, region="athens")] + + +async def test_get_zones_rejects_unknown_state() -> None: + async with AsyncClient(transport=ASGITransport(app=_app([])), base_url="http://test") as client: + response = await client.get("/internal/zones", params={"state": "degraded"}) + + assert response.status_code == 422 + + +async def test_get_zones_echoes_correlator_when_supplied() -> None: + response = await _get(_app([]), **{"x-correlator": "corr-42"}) + + assert response.headers["x-correlator"] == "corr-42" + + +async def test_get_zones_guarantees_provider_metadata() -> None: + response = await _get(_app([_zone(metadata={})])) + + assert response.json()[0]["metadata"]["provider"] == "test-provider" diff --git a/tests/integration/test_repositories.py b/tests/integration/test_repositories.py index 40658e6..ef33079 100644 --- a/tests/integration/test_repositories.py +++ b/tests/integration/test_repositories.py @@ -103,6 +103,7 @@ def _zone( ref: str = "zone-1", kind: ZoneKind = ZoneKind.RESOURCE, state: ZoneState = ZoneState.ACTIVE, + metadata: dict[str, object] | None = None, domains: list[Domain] | None = None, ) -> Zone: return Zone( @@ -112,6 +113,7 @@ def _zone( name="Zone", kind=kind, state=state, + metadata=metadata or {}, domains=domains or [], created_at=_now(), updated_at=_now(), @@ -423,6 +425,37 @@ async def test_create_assigns_id_and_get_by_id_round_trips( assert reloaded == saved +async def test_zone_repo_lists_resource_zones_with_state_and_region_filters( + db_session: AsyncSession, +) -> None: + repo = SqlZoneRepository(db_session) + matching = await repo.create( + _zone( + ref="athens-active", + state=ZoneState.ACTIVE, + metadata={"provider": "acme", "location": {"region": "athens"}}, + ) + ) + await repo.create( + _zone( + ref="athens-offline", + state=ZoneState.OFFLINE, + metadata={"provider": "acme", "location": {"region": "athens"}}, + ) + ) + await repo.create( + _zone( + ref="berlin-active", + state=ZoneState.ACTIVE, + metadata={"provider": "acme", "location": {"region": "berlin"}}, + ) + ) + + found = await repo.list_resource_zones(state=ZoneState.ACTIVE, region="athens") + + assert [zone.id for zone in found] == [matching.id] + + def test_service_order_type_values_match_persistence_model() -> None: assert {order_type.value for order_type in ServiceOrderType} == { "deploy_service", diff --git a/tests/test_config.py b/tests/test_config.py index b1f8a58..e53b267 100644 --- a/tests/test_config.py +++ b/tests/test_config.py @@ -33,6 +33,24 @@ def test_loads_all_fields_from_env(monkeypatch: pytest.MonkeyPatch) -> None: assert s.nats_settings.url == "nats://localhost:4222" assert s.nats_settings.connect_timeout == 10 assert s.nats_settings.max_reconnect_attempts == 3 + assert s.zone_provider_name == "oop-provider" + + +def test_zone_provider_name_can_be_overridden(monkeypatch: pytest.MonkeyPatch) -> None: + set_valid_env(monkeypatch) + monkeypatch.setenv("ZONE_PROVIDER_NAME", "acme-edge") + + s = Settings(_env_file=None) # type: ignore[call-arg] + + assert s.zone_provider_name == "acme-edge" + + +def test_zone_provider_name_must_not_be_empty(monkeypatch: pytest.MonkeyPatch) -> None: + set_valid_env(monkeypatch) + monkeypatch.setenv("ZONE_PROVIDER_NAME", "") + + with pytest.raises(ValidationError): + Settings(_env_file=None) # type: ignore[call-arg] @pytest.mark.parametrize( diff --git a/tests/unit/fakes.py b/tests/unit/fakes.py index b004860..034d1d8 100644 --- a/tests/unit/fakes.py +++ b/tests/unit/fakes.py @@ -91,10 +91,26 @@ class InMemoryZoneRepository(ZoneRepository): return found.model_copy(deep=True) if found is not None else None async def list_active_resource_zones(self) -> list[Zone]: + return await self.list_resource_zones(state=ZoneState.ACTIVE) + + async def list_resource_zones( + self, + *, + state: ZoneState | None = None, + region: str | None = None, + ) -> list[Zone]: return [ zone.model_copy(deep=True) for zone in self.rows.values() - if zone.kind == ZoneKind.RESOURCE and zone.state == ZoneState.ACTIVE + if zone.kind == ZoneKind.RESOURCE + and (state is None or zone.state == state) + and ( + region is None + or ( + isinstance(zone.metadata.get("location"), dict) + and zone.metadata["location"].get("region") == region + ) + ) ] async def create(self, zone: Zone) -> Zone: -- GitLab From 21994b4b24b422913a2ab067a46ffb8080ec5daf Mon Sep 17 00:00:00 2001 From: dgogos Date: Thu, 3 Sep 2026 17:18:16 +0300 Subject: [PATCH 2/3] fix: add unit test for ListZonesUseCase execution --- tests/unit/test_list_zones_use_case.py | 19 +++++++++++++++++++ 1 file changed, 19 insertions(+) create mode 100644 tests/unit/test_list_zones_use_case.py diff --git a/tests/unit/test_list_zones_use_case.py b/tests/unit/test_list_zones_use_case.py new file mode 100644 index 0000000..f4b036c --- /dev/null +++ b/tests/unit/test_list_zones_use_case.py @@ -0,0 +1,19 @@ +from unittest.mock import AsyncMock + +from srm.application.use_cases.topology import ListZonesCommand, ListZonesUseCase +from srm.domain.models.topology import Zone, ZoneKind, ZoneState + + +def _zone() -> Zone: + return Zone(ref="zone-1", name="Zone", kind=ZoneKind.RESOURCE, state=ZoneState.ACTIVE) + + +async def test_execute_forwards_filters_to_the_repo_and_returns_its_result() -> None: + zones = AsyncMock() + zones.list_resource_zones.return_value = [_zone()] + use_case = ListZonesUseCase(zones) + + result = await use_case.execute(ListZonesCommand(state=ZoneState.ACTIVE, region="athens")) + + zones.list_resource_zones.assert_awaited_once_with(state=ZoneState.ACTIVE, region="athens") + assert result is zones.list_resource_zones.return_value -- GitLab From 8fb1cad6e0eded85ee57ee89ceaad6cfcc0a7849 Mon Sep 17 00:00:00 2001 From: stentoumis Date: Fri, 4 Sep 2026 11:19:16 +0300 Subject: [PATCH 3/3] feat: get zones now also return domains --- src/srm/api/rest/router.py | 14 +++++++++++++- src/srm/api/rest/schemas.py | 9 +++++++++ tests/api/rest/test_zone_endpoint.py | 16 ++++++++++++++++ 3 files changed, 38 insertions(+), 1 deletion(-) diff --git a/src/srm/api/rest/router.py b/src/srm/api/rest/router.py index 64874f5..7f0208f 100644 --- a/src/srm/api/rest/router.py +++ b/src/srm/api/rest/router.py @@ -20,6 +20,7 @@ from srm.api.rest.schemas import ( ServiceCapabilityRequirementResponseSchema, ServiceDeploymentUnitCreateSchema, ServiceSpecificationCreateSchema, + ZoneDomainSummaryResponse, ZoneResponse, ) from srm.application.use_cases.topology import ListZonesCommand @@ -35,7 +36,7 @@ from srm.domain.models.catalog import ( ServiceDeploymentUnit, ServiceSpecification, ) -from srm.domain.models.topology import Zone, ZoneState +from srm.domain.models.topology import DomainKind, Zone, ZoneState internal = APIRouter(prefix="/internal", tags=["Internal"]) logger: structlog.BoundLogger = structlog.get_logger(__name__) @@ -153,6 +154,17 @@ def _build_zone_response(zone: Zone, provider_name: str) -> ZoneResponse: name=zone.name, state=zone.state.value, metadata=_zone_metadata(zone, provider_name), + domains=[ + ZoneDomainSummaryResponse( + id=domain.id, + name=domain.name, + kind=domain.kind.value, + state=domain.state.value, + metadata=domain.metadata, + ) + for domain in zone.domains + if domain.kind in {DomainKind.COMPUTE, DomainKind.NETWORK} + ], ) diff --git a/src/srm/api/rest/schemas.py b/src/srm/api/rest/schemas.py index 01c8048..09a7610 100644 --- a/src/srm/api/rest/schemas.py +++ b/src/srm/api/rest/schemas.py @@ -87,11 +87,20 @@ class GetServiceSpecificationResponse(DomainModel): service_capability_requirements: list[ServiceCapabilityRequirementResponseSchema] +class ZoneDomainSummaryResponse(DomainModel): + id: UUID + name: str + kind: str + state: str + metadata: dict[str, Any] = Field(default_factory=dict) + + class ZoneResponse(DomainModel): id: UUID name: str state: str metadata: dict[str, Any] + domains: list[ZoneDomainSummaryResponse] = Field(default_factory=list) class ErrorResponse(DomainModel): diff --git a/tests/api/rest/test_zone_endpoint.py b/tests/api/rest/test_zone_endpoint.py index e5c63d4..c6a2db2 100644 --- a/tests/api/rest/test_zone_endpoint.py +++ b/tests/api/rest/test_zone_endpoint.py @@ -82,6 +82,22 @@ async def test_get_zones_returns_oeg_compatible_shape() -> None: "provider": "acme", "location": {"region": "eu-central-1"}, }, + "domains": [ + { + "id": response.json()[0]["domains"][0]["id"], + "name": "Compute Domain", + "kind": "compute", + "state": "active", + "metadata": {"summary": "compute"}, + }, + { + "id": response.json()[0]["domains"][1]["id"], + "name": "Network Domain", + "kind": "network", + "state": "active", + "metadata": {"summary": "network"}, + }, + ], } ] -- GitLab