From 6620b9a63ee402a0e78289a9bd4c82885fe9094a Mon Sep 17 00:00:00 2001 From: Rafael Date: Wed, 27 May 2026 17:02:19 +0300 Subject: [PATCH] Adding CAMARA API EFN for accelleran rApp adapter --- .../oran/adapters/accelleran_ric/client.py | 160 ++++++++++ .../energyFootprintNotification/schemas.py | 59 ++++ src/sunrise6g_opensdk/oran/core/common.py | 13 +- src/sunrise6g_opensdk/oran/core/schemas.py | 9 + tests/oran/test_efn_subscription.py | 294 ++++++++++++++++++ 5 files changed, 534 insertions(+), 1 deletion(-) create mode 100644 src/sunrise6g_opensdk/oran/adapters/accelleran_ric/client.py create mode 100644 src/sunrise6g_opensdk/oran/adapters/accelleran_ric/energyFootprintNotification/schemas.py create mode 100644 tests/oran/test_efn_subscription.py diff --git a/src/sunrise6g_opensdk/oran/adapters/accelleran_ric/client.py b/src/sunrise6g_opensdk/oran/adapters/accelleran_ric/client.py new file mode 100644 index 0000000..42d7419 --- /dev/null +++ b/src/sunrise6g_opensdk/oran/adapters/accelleran_ric/client.py @@ -0,0 +1,160 @@ +#!/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] diff --git a/src/sunrise6g_opensdk/oran/adapters/accelleran_ric/energyFootprintNotification/schemas.py b/src/sunrise6g_opensdk/oran/adapters/accelleran_ric/energyFootprintNotification/schemas.py new file mode 100644 index 0000000..383cc80 --- /dev/null +++ b/src/sunrise6g_opensdk/oran/adapters/accelleran_ric/energyFootprintNotification/schemas.py @@ -0,0 +1,59 @@ +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" + ) diff --git a/src/sunrise6g_opensdk/oran/core/common.py b/src/sunrise6g_opensdk/oran/core/common.py index aacd7ff..b51aa56 100644 --- a/src/sunrise6g_opensdk/oran/core/common.py +++ b/src/sunrise6g_opensdk/oran/core/common.py @@ -1,5 +1,4 @@ # -*- 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) diff --git a/src/sunrise6g_opensdk/oran/core/schemas.py b/src/sunrise6g_opensdk/oran/core/schemas.py index 9175343..742fc64 100644 --- a/src/sunrise6g_opensdk/oran/core/schemas.py +++ b/src/sunrise6g_opensdk/oran/core/schemas.py @@ -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" diff --git a/tests/oran/test_efn_subscription.py b/tests/oran/test_efn_subscription.py new file mode 100644 index 0000000..5994f24 --- /dev/null +++ b/tests/oran/test_efn_subscription.py @@ -0,0 +1,294 @@ +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" -- GitLab