Commit 6620b9a6 authored by Rafael Pires's avatar Rafael Pires
Browse files

Adding CAMARA API EFN for accelleran rApp adapter

parent b7dc64c1
Loading
Loading
Loading
Loading
Loading
+160 −0
Original line number Diff line number Diff line
#!/bin/python3

###
import datetime
import logging
from typing import Dict, Optional
from uuid import UUID

from pydantic import AnyUrl

from sunrise6g_opensdk import logger

from ...adapters.accelleran_ric.energyFootprintNotification.schemas import (
    EFN_ReportCreationRequest,
    EnergyFootprintNotificationConfig,
)
from ...adapters.accelleran_ric.energyFootprintNotification.schemas import (
    SubscriptionDetail as EFN_sd,
)
from ...adapters.accelleran_ric.energyFootprintNotification.schemas import (
    SubscriptionEventType as EFN_set,
)
from ...adapters.accelleran_ric.energyFootprintNotification.schemas import (
    SubscriptionRequest as EFN_sr,
)
from ...adapters.accelleran_ric.energyFootprintNotification.schemas import (
    TimePeriod,
)
from ...core import common as oran_common
from ...core.base_oran_client import BaseOranClient
from ...core.common import OranHttpError, requires_capability
from ...core.schemas import Device, Protocol, SinkCredential

###
log = logger.get_logger(__name__)


class OranManager(BaseOranClient):
    """This client implements the BaseOranClient and translate CAMARA apis (see https://github.com/camaraproject/DeviceReachabilityStatus/blob/r1.2/code/API_definitions/device-reachability-status-subscriptions.yaml) to Accelleran r/xApp APIs"""

    capabilities = ["oran-energy-footprint-notification"]

    def __init__(self, base_url: str, scs_as_id: str):
        """
        This client implements the BaseOranClient and translates the
        CAMARA APIs into specific HTTP requests understandable by the Accelleran EFN rApp.
        """
        try:
            self.base_url = base_url
            self.scs_as_id = scs_as_id
            log.info(
                f"initialized Accellleran rApp client with base url: {self.base_url} \
                and scs {self.scs_as_id}"
            )
        except Exception as e:
            log.error(f"Failed to init Accelleran rApp client\n->{e}")
            raise e

    @requires_capability("oran-energy-footprint-notification")
    def load_efn_config(self, config: Dict) -> EFN_ReportCreationRequest:
        """Helper function to load the details of the energy footprint notification subscription from a dictionary. This is useful for testing purposes, as it allows to easily create a subscription request from a dictionary."""
        try:
            device: Device = self._get_device_from_config(
                config
            )  # There may be different types of Device Identifiers

            subscription_detail = EFN_sd(device=device)
            efn_config = EnergyFootprintNotificationConfig(
                subscriptionDetail=subscription_detail,
                subscriptionExpireTime=datetime(config.get("subscriptionExpireTime")),
                subscriptionMaxEvents=int(config.get("subscriptionMaxEvents")),
                initialEvent=bool(config.get("initialEvent")),
            )
            efn_request = EFN_ReportCreationRequest(
                service=config["service"],
                subscription_request=EFN_sr(
                    sink=config["sink"],
                    sink_credential=config.get("sinkCredential"),
                    protocol=config["protocol"],
                    config=efn_config,
                    types=config["types"],
                ),
                time_period=config.get("timePeriod"),
            )
        except Exception as e:
            logging.error(f"Error loading EFN config\n->{e}")
            raise e
        return efn_request

    @requires_capability("oran-energy-footprint-notification")
    def calculate_energy_footprint_notification_report(
        self,
        service: list[UUID],
        sink: AnyUrl,
        config: EnergyFootprintNotificationConfig,
        protocol: Protocol,
        types: EFN_set,
        time_period: Optional[TimePeriod] = None,
        sink_credential: Optional[SinkCredential] = None,
    ) -> dict | str:
        """
        Implements POST for energy footprint notification request in path is /calculate-carbon-footprint creates ReportCreationRequest (service{array}, timeperiod:optional, subscriptionRequest, subscriptionID:optional), either returns json response or status code as string
        """
        subscription_request = EFN_sr(
            sink=sink,
            sink_credential=sink_credential,
            protocol=protocol,
            config=config,
            types=[types],
        )

        payload = EFN_ReportCreationRequest(
            service=service, subscription_request=subscription_request, time_period=time_period
        )
        try:
            request = oran_common.oran_subscription_post(
                self.base_url, self.scs_as_id, model_payload=payload
            )
            return request
        except OranHttpError as e:
            logging.error(
                f"POST Response error code\n->Status Code: {e.status_code}\n->Error Body: {e.body}"
            )
            return e.status_code

    @requires_capability("oran-energy-footprint-notification")
    def retrieve_energy_footprint_subscription(self, session_id) -> dict | str:
        """
        Implements GET request for energy footprint notification, either returns json response or status code as string
        """
        try:
            request = oran_common.oran_subscription_get(self.base_url, self.scs_as_id, session_id)
            return request
        except OranHttpError as e:
            logging.error(
                f"GET Response error code\n->Status Code: {e.status_code}\n->Error Body: {e.body}"
            )
            return e.status_code

    @requires_capability("oran-energy-footprint-notification")
    def delete_energy_footprint_subscription(self, session_id) -> dict | str:
        """
        Implements DELETE request for energy footprint notification, either returns json response or status code as string
        """
        try:
            request = oran_common.oran_subscription_delete(
                self.base_url, self.scs_as_id, session_id
            )
            return request
        except OranHttpError as e:
            logging.error(
                f"DELETE Response error code\n->Status Code: {e.status_code}\n->Error Body: {e.body}"
            )
            return e.status_code

    def _get_device_from_config(self, config):
        """
        Retrieves one device from the multiple identifiers possible
        """
        return config.get("device").values()[0]
+59 −0
Original line number Diff line number Diff line
from datetime import datetime
from enum import Enum
from typing import Annotated, List, Optional, Union

from pydantic import AnyUrl, BaseModel, Field, IPvAnyInterface

from ....core.schemas import Device, Protocol, SinkCredential


class SubscriptionEventType(str, Enum):
    ENERGY = "org.camaraproject.energy-footprint-notification.v1.energy"
    CARBON_FOOTPRINT = "org.camaraproject.energy-footprint-notification.v1.carbon-footprint"


class SubscriptionDetail(BaseModel):
    device: Device  # stated as an object but not defined


class EnergyFootprintNotificationConfig(BaseModel):
    subscription_detail: SubscriptionDetail = Field(
        ...,
        alias="subscriptionDetail",
        description="Target device for the reachbility subscription",
    )
    subscription_expire_time: Optional[datetime] = Field(None, alias="subscriptionExpireTime")
    subscription_max_events: Optional[int] = Field(None, alias="subscriptionMaxEvents", ge=1)
    initial_event: Optional[bool] = Field(
        False,
        alias="initialEvent",
        description="If true, it sends an event immmedately after subscription",
    )


class TimePeriod(BaseModel):
    startDate: datetime
    endDate: datetime


class SubscriptionRequest(BaseModel):
    sink: Union[IPvAnyInterface, AnyUrl]
    sink_credential: Optional[SinkCredential] = Field(
        None,
        description="Provides authentication and authorization information necessary to enable delivery of event target",
    )
    protocol: Protocol
    config: EnergyFootprintNotificationConfig
    types: Annotated[List[SubscriptionEventType], Field(min_length=1, max_length=1)]


class EFN_ReportCreationRequest(BaseModel):
    service: list[str]  # list[UUID]
    time_period: Optional[TimePeriod] = Field(
        None,
        description="start date is required, ending is optional. It must follow [RFC 3339](https://datatracker.ietf.org/doc/html/rfc3339#section-5.6) including timezone",
    )
    subscription_request: SubscriptionRequest
    request_id: Optional[str] = Field(
        None, description="Request Identifier. Returned by the API, used to update it"
    )
+12 −1
Original line number Diff line number Diff line
# -*- coding: utf-8 -*-

import requests
from pydantic import BaseModel

@@ -84,6 +83,12 @@ def oran_subscription_post(base_url: str, scs_as_id: str, model_payload: BaseMod
    return _make_request("POST", url, data=data)


def oran_subscription_get(url: str, scs_as_id: str, session_id: BaseModel) -> dict:
    url = oran_subscription_build(url, scs_as_id, session_id)
    print(url)
    return _make_request("GET", url)


def oran_subscription_build(base_url: str, scs_as_id: str, session_id: str = None):
    url = f"{base_url}/{scs_as_id}/subscriptions"
    if session_id is not None and len(session_id) > 0:
@@ -92,6 +97,12 @@ def oran_subscription_build(base_url: str, scs_as_id: str, session_id: str = Non
        return url


def oran_subscription_delete(base_url: str, scs_as_id: str, session_id: str):
    url = oran_subscription_build(base_url, scs_as_id, session_id)
    print(url)
    return _make_request("DELETE", url)


# Policy methods
def oran_policy_post(base_url: str, scs_as_id: str, model_payload: BaseModel) -> dict:
    data = model_payload.model_dump_json(exclude_none=True, by_alias=True)
+9 −0
Original line number Diff line number Diff line
@@ -648,3 +648,12 @@ class SessionInfo(BaseSessionInfo):
    ] = None
    qosStatus: QosStatus
    statusInfo: StatusInfo | None = None


class Protocol(Enum):
    HTTP = "HTTP"
    MQTT3 = "MQTT3"
    MQTT5 = "MQTT5"
    AMQP = "AMQP"
    NATS = "NATS"
    KAFKA = "KAFKA"
+294 −0
Original line number Diff line number Diff line
from datetime import datetime, timezone
from unittest.mock import MagicMock
from uuid import uuid4

import pytest

from sunrise6g_opensdk.oran.adapters.accelleran_ric.client import OranManager
from sunrise6g_opensdk.oran.adapters.accelleran_ric.energyFootprintNotification.schemas import (
    EnergyFootprintNotificationConfig,
    SubscriptionDetail,
    SubscriptionEventType,
    SubscriptionRequest,
    TimePeriod,
)
from sunrise6g_opensdk.oran.core import common as oran_common
from sunrise6g_opensdk.oran.core import schemas as oran_schemas

# --- --- --- --- --- --- ---
# Create pytest fixtures

BASE_URL = "http://localhost:9091"
SCS_AS_ID = "test-scs-id"
SESSION_ID = str(uuid4())


@pytest.fixture(autouse=True)
def mock_oran_common(monkeypatch):
    monkeypatch.setattr(oran_common, "oran_subscription_post", MagicMock())
    monkeypatch.setattr(oran_common, "oran_subscription_get", MagicMock())
    monkeypatch.setattr(oran_common, "oran_subscription_delete", MagicMock())
    yield


@pytest.fixture
def manager():
    return OranManager(base_url=BASE_URL, scs_as_id=SCS_AS_ID)


@pytest.fixture
def efn_config():
    device = oran_schemas.Device(
        networkAccessIdentifier=oran_schemas.NetworkAccessIdentifier("123")
    )
    detail = SubscriptionDetail(device=device)
    return EnergyFootprintNotificationConfig(subscriptionDetail=detail)


@pytest.fixture
def time_period():
    return TimePeriod(
        startDate=datetime(2024, 1, 1, tzinfo=timezone.utc),
        endDate=datetime(2024, 1, 2, tzinfo=timezone.utc),
    )


@pytest.fixture(autouse=True)
def reset_mocks():
    oran_common.oran_subscription_post.reset_mock(side_effect=True)
    oran_common.oran_subscription_get.reset_mock(side_effect=True)
    oran_common.oran_subscription_delete.reset_mock(side_effect=True)
    yield


# --- --- --- --- --- --- ---
# Test OranManager initialization


class TestOranManagerInit:
    def test_stores_base_url(self):
        assert OranManager(BASE_URL, SCS_AS_ID).base_url == BASE_URL

    def test_stores_scs_as_id(self):
        assert OranManager(BASE_URL, SCS_AS_ID).scs_as_id == SCS_AS_ID

    def test_capabilities_declared(self):
        assert "oran-energy-footprint-notification" in OranManager.capabilities


# --- --- --- --- --- --- ---
# Test HTTP Methods
# --- --- --- --- --- --- ---

# --- --- --- --- --- --- ---
# Test POST Method, calculate_energy_footprint_notification_report


class TestCalculateEnergyFootprintNotificationReport:
    def test_calls_oran_post(self, manager, efn_config):
        oran_common.oran_subscription_post.return_value = {"requestId": "abc"}
        manager.calculate_energy_footprint_notification_report(
            service=[str(uuid4())],
            sink="https://example.com/callback",
            config=efn_config,
            protocol="HTTP",
            types=SubscriptionEventType.ENERGY,
        )
        oran_common.oran_subscription_post.assert_called_once()

    def test_passes_base_url_and_scs(self, manager, efn_config):
        manager.calculate_energy_footprint_notification_report(
            service=[str(uuid4())],
            sink="https://example.com/callback",
            config=efn_config,
            protocol="HTTP",
            types=SubscriptionEventType.ENERGY,
        )
        args, _ = oran_common.oran_subscription_post.call_args
        assert args[0] == BASE_URL
        assert args[1] == SCS_AS_ID

    def test_returns_api_response_on_success(self, manager, efn_config):
        expected = {"requestId": "xyz-123"}
        oran_common.oran_subscription_post.return_value = expected
        result = manager.calculate_energy_footprint_notification_report(
            service=[str(uuid4())],
            sink="https://example.com/callback",
            config=efn_config,
            protocol="HTTP",
            types=SubscriptionEventType.ENERGY,
        )
        assert result == expected

    def test_returns_status_code_on_http_error(self, manager, efn_config):
        oran_common.oran_subscription_post.side_effect = oran_common.OranHttpError(
            "Unprocessable", 422
        )
        result = manager.calculate_energy_footprint_notification_report(
            service=[str(uuid4())],
            sink="https://example.com/callback",
            config=efn_config,
            protocol="HTTP",
            types=SubscriptionEventType.ENERGY,
        )
        assert result == 422

    def test_different_http_error_codes_are_returned(self, manager, efn_config):
        for code in [400, 401, 403, 404, 500]:
            oran_common.oran_subscription_post.side_effect = oran_common.OranHttpError(
                "error", code
            )
            result = manager.calculate_energy_footprint_notification_report(
                service=[str(uuid4())],
                sink="https://example.com/callback",
                config=efn_config,
                protocol="HTTP",
                types=SubscriptionEventType.ENERGY,
            )
            assert result == code

    def test_with_optional_time_period(self, manager, efn_config, time_period):
        manager.calculate_energy_footprint_notification_report(
            service=[str(uuid4())],
            sink="https://example.com/callback",
            config=efn_config,
            protocol="HTTP",
            types=SubscriptionEventType.ENERGY,
            time_period=time_period,
        )
        _, kwargs = oran_common.oran_subscription_post.call_args
        assert kwargs["model_payload"].time_period == time_period

    def test_no_time_period_by_default(self, manager, efn_config):
        manager.calculate_energy_footprint_notification_report(
            service=[str(uuid4())],
            sink="https://example.com/callback",
            config=efn_config,
            protocol="HTTP",
            types=SubscriptionEventType.ENERGY,
        )
        _, kwargs = oran_common.oran_subscription_post.call_args
        assert kwargs["model_payload"].time_period is None

    def test_with_multiple_service_ids(self, manager, efn_config):
        services = [str(uuid4()), str(uuid4()), str(uuid4())]
        manager.calculate_energy_footprint_notification_report(
            service=services,
            sink="https://example.com/callback",
            config=efn_config,
            protocol="HTTP",
            types=SubscriptionEventType.ENERGY,
        )
        _, kwargs = oran_common.oran_subscription_post.call_args
        assert kwargs["model_payload"].service == services

    def test_carbon_footprint_event_type(self, manager, efn_config):
        manager.calculate_energy_footprint_notification_report(
            service=[str(uuid4())],
            sink="https://example.com/callback",
            config=efn_config,
            protocol="HTTP",
            types=SubscriptionEventType.CARBON_FOOTPRINT,
        )
        _, kwargs = oran_common.oran_subscription_post.call_args
        assert (
            SubscriptionEventType.CARBON_FOOTPRINT
            in kwargs["model_payload"].subscription_request.types
        )


# --- --- --- --- --- --- ---
# Test GET Method, retrieve_energy_footprint_subscription


class TestRetrieveEnergyFootprintSubscription:
    def test_calls_oran_get_with_correct_args(self, manager):
        manager.retrieve_energy_footprint_subscription(SESSION_ID)
        oran_common.oran_subscription_get.assert_called_once_with(BASE_URL, SCS_AS_ID, SESSION_ID)

    def test_returns_subscription_data_on_success(self, manager):
        expected = {"subscriptionId": SESSION_ID, "status": "active"}
        oran_common.oran_subscription_get.return_value = expected
        assert manager.retrieve_energy_footprint_subscription(SESSION_ID) == expected

    def test_returns_status_code_on_http_error(self, manager):
        oran_common.oran_subscription_get.side_effect = oran_common.OranHttpError("Not found", 404)
        assert manager.retrieve_energy_footprint_subscription(SESSION_ID) == 404

    def test_different_http_error_codes_are_returned(self, manager):
        for code in [401, 403, 500]:
            oran_common.oran_subscription_get.side_effect = oran_common.OranHttpError("error", code)
            assert manager.retrieve_energy_footprint_subscription(SESSION_ID) == code


# --- --- --- --- --- --- ----
# Test DELETE Method, delete_energy_footprint_subscription


class TestDeleteEnergyFootprintSubscription:
    def test_calls_oran_delete_with_correct_args(self, manager):
        manager.delete_energy_footprint_subscription(SESSION_ID)
        oran_common.oran_subscription_delete.assert_called_once_with(
            BASE_URL, SCS_AS_ID, SESSION_ID
        )

    def test_returns_response_on_success(self, manager):
        oran_common.oran_subscription_delete.return_value = {"status": "deleted"}
        assert manager.delete_energy_footprint_subscription(SESSION_ID) == {"status": "deleted"}

    def test_returns_status_code_on_http_error(self, manager):
        oran_common.oran_subscription_delete.side_effect = oran_common.OranHttpError(
            "Forbidden", 403
        )
        assert manager.delete_energy_footprint_subscription(SESSION_ID) == 403

    def test_different_http_error_codes_are_returned(self, manager):
        for code in [401, 404, 500]:
            oran_common.oran_subscription_delete.side_effect = oran_common.OranHttpError(
                "error", code
            )
            assert manager.delete_energy_footprint_subscription(SESSION_ID) == code


# --- --- --- --- --- --- ---
# Schema validation


class TestSchemas:
    def test_subscription_event_type_energy_value(self):
        assert (
            SubscriptionEventType.ENERGY
            == "org.camaraproject.energy-footprint-notification.v1.energy"
        )

    def test_subscription_event_type_carbon_value(self):
        assert (
            SubscriptionEventType.CARBON_FOOTPRINT
            == "org.camaraproject.energy-footprint-notification.v1.carbon-footprint"
        )

    def test_efn_config_defaults(self):
        detail = SubscriptionDetail(device=oran_schemas.Device(deviceId="123"))
        cfg = EnergyFootprintNotificationConfig(subscriptionDetail=detail)
        assert cfg.initial_event is False
        assert cfg.subscription_expire_time is None
        assert cfg.subscription_max_events is None

    def test_efn_config_max_events_minimum(self):
        detail = SubscriptionDetail(device=oran_schemas.Device(deviceId="456"))
        with pytest.raises(Exception):
            EnergyFootprintNotificationConfig(subscriptionDetail=detail, subscriptionMaxEvents=0)

    def test_subscription_request_requires_at_least_one_type(self, efn_config):
        with pytest.raises(Exception):
            SubscriptionRequest(
                sink="https://example.com/cb",
                protocol="HTTP",
                config=efn_config,
                types=[],
            )

    def test_oran_http_error_exposes_status_code_and_body(self):
        err = oran_common.OranHttpError("bad payload", 422, "bad payload")
        assert err.status_code == 422
        assert err.body == "bad payload"