Commit 400bc525 authored by Lluis Gifre Renom's avatar Lluis Gifre Renom
Browse files

Context:

- Implemented support for Slices
- missing: unitary tests
parent 17719702
Loading
Loading
Loading
Loading
+16 −0
Original line number Diff line number Diff line
@@ -21,6 +21,7 @@ from .DeviceModel import DeviceModel
from .EndPointModel import EndPointModel
from .LinkModel import LinkModel
from .ServiceModel import ServiceModel
from .SliceModel import SliceModel
from .TopologyModel import TopologyModel

LOGGER = logging.getLogger(__name__)
@@ -40,6 +41,21 @@ class ServiceEndPointModel(Model): # pylint: disable=abstract-method
    service_fk = ForeignKeyField(ServiceModel)
    endpoint_fk = ForeignKeyField(EndPointModel)

class SliceEndPointModel(Model): # pylint: disable=abstract-method
    pk = PrimaryKeyField()
    slice_fk = ForeignKeyField(SliceModel)
    endpoint_fk = ForeignKeyField(EndPointModel)

class SliceServiceModel(Model): # pylint: disable=abstract-method
    pk = PrimaryKeyField()
    slice_fk = ForeignKeyField(SliceModel)
    service_fk = ForeignKeyField(ServiceModel)

class SliceSubSliceModel(Model): # pylint: disable=abstract-method
    pk = PrimaryKeyField()
    slice_fk = ForeignKeyField(SliceModel)
    sub_slice_fk = ForeignKeyField(SliceModel)

class TopologyDeviceModel(Model): # pylint: disable=abstract-method
    pk = PrimaryKeyField()
    topology_fk = ForeignKeyField(TopologyModel)
+85 −0
Original line number Diff line number Diff line
# Copyright 2021-2023 H2020 TeraFlow (https://www.teraflow-h2020.eu/)
#
# Licensed under the Apache License, Version 2.0 (the "License");
# you may not use this file except in compliance with the License.
# You may obtain a copy of the License at
#
#      http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.

import functools, logging, operator
from enum import Enum
from typing import Dict, List
from common.orm.fields.EnumeratedField import EnumeratedField
from common.orm.fields.ForeignKeyField import ForeignKeyField
from common.orm.fields.PrimaryKeyField import PrimaryKeyField
from common.orm.fields.StringField import StringField
from common.orm.model.Model import Model
from common.orm.HighLevel import get_related_objects
from context.proto.context_pb2 import SliceStatusEnum
from .ConstraintModel import ConstraintsModel
from .ContextModel import ContextModel
from .Tools import grpc_to_enum

LOGGER = logging.getLogger(__name__)

class ORM_SliceStatusEnum(Enum):
    UNDEFINED = SliceStatusEnum.SLICESTATUS_UNDEFINED
    PLANNED   = SliceStatusEnum.SLICESTATUS_PLANNED
    INIT      = SliceStatusEnum.SLICESTATUS_INIT
    ACTIVE    = SliceStatusEnum.SLICESTATUS_ACTIVE
    DEINIT    = SliceStatusEnum.SLICESTATUS_DEINIT

grpc_to_enum__slice_status = functools.partial(
    grpc_to_enum, SliceStatusEnum, ORM_SliceStatusEnum)

class SliceModel(Model):
    pk = PrimaryKeyField()
    context_fk = ForeignKeyField(ContextModel)
    slice_uuid = StringField(required=True, allow_empty=False)
    slice_constraints_fk = ForeignKeyField(ConstraintsModel)
    slice_status = EnumeratedField(ORM_SliceStatusEnum, required=True)

    def dump_id(self) -> Dict:
        context_id = ContextModel(self.database, self.context_fk).dump_id()
        return {
            'context_id': context_id,
            'slice_uuid': {'uuid': self.slice_uuid},
        }

    def dump_endpoint_ids(self) -> List[Dict]:
        from .RelationModels import SliceEndPointModel # pylint: disable=import-outside-toplevel
        db_endpoints = get_related_objects(self, SliceEndPointModel, 'endpoint_fk')
        return [db_endpoint.dump_id() for db_endpoint in sorted(db_endpoints, key=operator.attrgetter('pk'))]

    def dump_constraints(self) -> List[Dict]:
        return ConstraintsModel(self.database, self.slice_constraints_fk).dump()

    def dump_service_ids(self) -> List[Dict]:
        from .RelationModels import SliceServiceModel # pylint: disable=import-outside-toplevel
        db_services = get_related_objects(self, SliceServiceModel, 'service_fk')
        return [db_service.dump_id() for db_service in sorted(db_services, key=operator.attrgetter('pk'))]

    def dump_subslice_ids(self) -> List[Dict]:
        from .RelationModels import SliceSubSliceModel # pylint: disable=import-outside-toplevel
        db_subslices = get_related_objects(self, SliceSubSliceModel, 'sub_slice_fk')
        return [db_subslice.dump_id() for db_subslice in sorted(db_subslices, key=operator.attrgetter('pk'))]

    def dump(   # pylint: disable=arguments-differ
            self, include_endpoint_ids=True, include_constraints=True, include_service_ids=True,
            include_subslice_ids=True
        ) -> Dict:
        result = {
            'slice_id': self.dump_id(),
            'slice_status': {'slice_status': self.slice_status.value},
        }
        if include_endpoint_ids: result['slice_endpoint_ids'] = self.dump_endpoint_ids()
        if include_constraints: result['slice_constraints'] = self.dump_constraints()
        if include_service_ids: result['slice_service_ids'] = self.dump_service_ids()
        if include_subslice_ids: result['sub_subslice_ids'] = self.dump_subslice_ids()
        return result
+3 −2
Original line number Diff line number Diff line
@@ -12,13 +12,14 @@
# See the License for the specific language governing permissions and
# limitations under the License.

TOPIC_CONNECTION = 'connection'
TOPIC_CONTEXT    = 'context'
TOPIC_TOPOLOGY   = 'topology'
TOPIC_DEVICE     = 'device'
TOPIC_LINK       = 'link'
TOPIC_SERVICE    = 'service'
TOPIC_CONNECTION = 'connection'
TOPIC_SLICE      = 'slice'

TOPICS = {TOPIC_CONTEXT, TOPIC_TOPOLOGY, TOPIC_DEVICE, TOPIC_LINK, TOPIC_SERVICE, TOPIC_CONNECTION}
TOPICS = {TOPIC_CONNECTION, TOPIC_CONTEXT, TOPIC_TOPOLOGY, TOPIC_DEVICE, TOPIC_LINK, TOPIC_SERVICE, TOPIC_SLICE}

CONSUME_TIMEOUT = 0.5 # seconds
+154 −7
Original line number Diff line number Diff line
@@ -24,8 +24,8 @@ from common.rpc_method_wrapper.ServiceExceptions import InvalidArgumentException
from context.proto.context_pb2 import (
    Connection, ConnectionEvent, ConnectionId, ConnectionIdList, ConnectionList, Context, ContextEvent, ContextId,
    ContextIdList, ContextList, Device, DeviceEvent, DeviceId, DeviceIdList, DeviceList, Empty, EventTypeEnum, Link,
    LinkEvent, LinkId, LinkIdList, LinkList, Service, ServiceEvent, ServiceId, ServiceIdList, ServiceList, Topology,
    TopologyEvent, TopologyId, TopologyIdList, TopologyList)
    LinkEvent, LinkId, LinkIdList, LinkList, Service, ServiceEvent, ServiceId, ServiceIdList, ServiceList, Slice, SliceEvent,
    SliceId, SliceIdList, SliceList, Topology, TopologyEvent, TopologyId, TopologyIdList, TopologyList)
from context.proto.context_pb2_grpc import ContextServiceServicer
from context.service.database.ConfigModel import ConfigModel, ConfigRuleModel, grpc_config_rules_to_raw, update_config
from context.service.database.ConnectionModel import ConnectionModel, PathHopModel, PathModel, set_path
@@ -37,12 +37,13 @@ from context.service.database.EndPointModel import EndPointModel, KpiSampleTypeM
from context.service.database.Events import notify_event
from context.service.database.LinkModel import LinkModel
from context.service.database.RelationModels import (
    ConnectionSubServiceModel, LinkEndPointModel, ServiceEndPointModel, TopologyDeviceModel, TopologyLinkModel)
    ConnectionSubServiceModel, LinkEndPointModel, ServiceEndPointModel, SliceEndPointModel, SliceServiceModel, SliceSubSliceModel, TopologyDeviceModel, TopologyLinkModel)
from context.service.database.ServiceModel import (
    ServiceModel, grpc_to_enum__service_status, grpc_to_enum__service_type)
from context.service.database.SliceModel import SliceModel, grpc_to_enum__slice_status
from context.service.database.TopologyModel import TopologyModel
from .Constants import (
    CONSUME_TIMEOUT, TOPIC_CONNECTION, TOPIC_CONTEXT, TOPIC_DEVICE, TOPIC_LINK, TOPIC_SERVICE, TOPIC_TOPOLOGY)
    CONSUME_TIMEOUT, TOPIC_CONNECTION, TOPIC_CONTEXT, TOPIC_DEVICE, TOPIC_LINK, TOPIC_SERVICE, TOPIC_SLICE, TOPIC_TOPOLOGY)

LOGGER = logging.getLogger(__name__)

@@ -54,6 +55,7 @@ METHOD_NAMES = [
    'ListDeviceIds',     'ListDevices',     'GetDevice',     'SetDevice',     'RemoveDevice',     'GetDeviceEvents',
    'ListLinkIds',       'ListLinks',       'GetLink',       'SetLink',       'RemoveLink',       'GetLinkEvents',
    'ListServiceIds',    'ListServices',    'GetService',    'SetService',    'RemoveService',    'GetServiceEvents',
    'ListSliceIds',      'ListSlices',      'GetSlice',      'SetSlice',      'RemoveSlice',      'GetSliceEvents',
]
METRICS = create_metrics(SERVICE_NAME, METHOD_NAMES)

@@ -183,7 +185,8 @@ class ContextServiceServicerImpl(ContextServiceServicer):
            topology_uuid = request.topology_id.topology_uuid.uuid
            str_topology_key = key_to_str([context_uuid, topology_uuid])
            result : Tuple[TopologyModel, bool] = update_or_create_object(
                self.database, TopologyModel, str_topology_key, {'context_fk': db_context, 'topology_uuid': topology_uuid})
                self.database, TopologyModel, str_topology_key, {
                    'context_fk': db_context, 'topology_uuid': topology_uuid})
            db_topology,updated = result

            for device_id in request.device_ids:
@@ -403,7 +406,8 @@ class ContextServiceServicerImpl(ContextServiceServicer):
                    str_topology_key = key_to_str([endpoint_topology_context_uuid, endpoint_topology_uuid])
                    db_topology : TopologyModel = get_object(self.database, TopologyModel, str_topology_key)
                    str_topology_device_key = key_to_str([str_topology_key, endpoint_device_uuid], separator='--')
                    get_object(self.database, TopologyDeviceModel, str_topology_device_key) # check device is in topology
                    # check device is in topology
                    get_object(self.database, TopologyDeviceModel, str_topology_device_key)
                    str_endpoint_key = key_to_str([str_endpoint_key, str_topology_key], separator=':')

                db_endpoint : EndPointModel = get_object(self.database, EndPointModel, str_endpoint_key)
@@ -491,7 +495,8 @@ class ContextServiceServicerImpl(ContextServiceServicer):
                    raise InvalidArgumentException(
                        'request.service_endpoint_ids[{:d}].topology_id.context_id.context_uuid.uuid'.format(i),
                        endpoint_topology_context_uuid,
                        ['should be == {:s}({:s})'.format('request.service_id.context_id.context_uuid.uuid', context_uuid)])
                        ['should be == {:s}({:s})'.format(
                            'request.service_id.context_id.context_uuid.uuid', context_uuid)])

            service_uuid = request.service_id.service_uuid.uuid
            str_service_key = key_to_str([context_uuid, service_uuid])
@@ -574,6 +579,148 @@ class ContextServiceServicerImpl(ContextServiceServicer):
            yield ServiceEvent(**json.loads(message.content))


    # ----- Slice ----------------------------------------------------------------------------------------------------

    @safe_and_metered_rpc_method(METRICS, LOGGER)
    def ListSliceIds(self, request: ContextId, context : grpc.ServicerContext) -> SliceIdList:
        with self.lock:
            db_context : ContextModel = get_object(self.database, ContextModel, request.context_uuid.uuid)
            db_slices : Set[SliceModel] = get_related_objects(db_context, SliceModel)
            db_slices = sorted(db_slices, key=operator.attrgetter('pk'))
            return SliceIdList(slice_ids=[db_slice.dump_id() for db_slice in db_slices])

    @safe_and_metered_rpc_method(METRICS, LOGGER)
    def ListSlices(self, request: ContextId, context : grpc.ServicerContext) -> SliceList:
        with self.lock:
            db_context : ContextModel = get_object(self.database, ContextModel, request.context_uuid.uuid)
            db_slices : Set[SliceModel] = get_related_objects(db_context, SliceModel)
            db_slices = sorted(db_slices, key=operator.attrgetter('pk'))
            return SliceList(slices=[db_slice.dump() for db_slice in db_slices])

    @safe_and_metered_rpc_method(METRICS, LOGGER)
    def GetSlice(self, request: SliceId, context : grpc.ServicerContext) -> Slice:
        with self.lock:
            str_key = key_to_str([request.context_id.context_uuid.uuid, request.slice_uuid.uuid])
            db_slice : SliceModel = get_object(self.database, SliceModel, str_key)
            return Slice(**db_slice.dump(
                include_endpoint_ids=True, include_constraints=True, include_service_ids=True,
                include_subslice_ids=True))

    @safe_and_metered_rpc_method(METRICS, LOGGER)
    def SetSlice(self, request: Slice, context : grpc.ServicerContext) -> SliceId:
        with self.lock:
            context_uuid = request.service_id.context_id.context_uuid.uuid
            db_context : ContextModel = get_object(self.database, ContextModel, context_uuid)

            for i,endpoint_id in enumerate(request.slice_endpoint_ids):
                endpoint_topology_context_uuid = endpoint_id.topology_id.context_id.context_uuid.uuid
                if len(endpoint_topology_context_uuid) > 0 and context_uuid != endpoint_topology_context_uuid:
                    raise InvalidArgumentException(
                        'request.slice_endpoint_ids[{:d}].topology_id.context_id.context_uuid.uuid'.format(i),
                        endpoint_topology_context_uuid,
                        ['should be == {:s}({:s})'.format(
                            'request.slice_id.context_id.context_uuid.uuid', context_uuid)])

            slice_uuid = request.slice_id.slice_uuid.uuid
            str_slice_key = key_to_str([context_uuid, slice_uuid])

            constraints_result = set_constraints(
                self.database, str_slice_key, 'constraints', request.slice_constraints)
            db_constraints = constraints_result[0][0]

            result : Tuple[SliceModel, bool] = update_or_create_object(self.database, SliceModel, str_slice_key, {
                'context_fk'          : db_context,
                'slice_uuid'          : slice_uuid,
                'slice_constraints_fk': db_constraints,
                'slice_status'        : grpc_to_enum__slice_status(request.slice_status.slice_status),
            })
            db_slice, updated = result

            for i,endpoint_id in enumerate(request.slice_endpoint_ids):
                endpoint_uuid                  = endpoint_id.endpoint_uuid.uuid
                endpoint_device_uuid           = endpoint_id.device_id.device_uuid.uuid
                endpoint_topology_uuid         = endpoint_id.topology_id.topology_uuid.uuid
                endpoint_topology_context_uuid = endpoint_id.topology_id.context_id.context_uuid.uuid

                str_endpoint_key = key_to_str([endpoint_device_uuid, endpoint_uuid])
                if len(endpoint_topology_context_uuid) > 0 and len(endpoint_topology_uuid) > 0:
                    str_topology_key = key_to_str([endpoint_topology_context_uuid, endpoint_topology_uuid])
                    str_endpoint_key = key_to_str([str_endpoint_key, str_topology_key], separator=':')

                db_endpoint : EndPointModel = get_object(self.database, EndPointModel, str_endpoint_key)

                str_slice_endpoint_key = key_to_str([slice_uuid, str_endpoint_key], separator='--')
                result : Tuple[SliceEndPointModel, bool] = get_or_create_object(
                    self.database, SliceEndPointModel, str_slice_endpoint_key, {
                        'slice_fk': db_slice, 'endpoint_fk': db_endpoint})
                #db_slice_endpoint, slice_endpoint_created = result

            for i,service_id in enumerate(request.slice_service_ids):
                service_uuid         = service_id.service_uuid.uuid
                service_context_uuid = service_id.context_id.context_uuid.uuid
                str_service_key = key_to_str([service_context_uuid, service_uuid])
                db_service : ServiceModel = get_object(self.database, ServiceModel, str_service_key)

                str_slice_service_key = key_to_str([str_slice_key, str_service_key], separator='--')
                result : Tuple[SliceServiceModel, bool] = get_or_create_object(
                    self.database, SliceServiceModel, str_slice_service_key, {
                        'slice_fk': db_slice, 'service_fk': db_service})
                #db_slice_service, slice_service_created = result

            for i,subslice_id in enumerate(request.slice_subslice_ids):
                subslice_uuid         = subslice_id.slice_uuid.uuid
                subslice_context_uuid = subslice_id.context_id.context_uuid.uuid
                str_subslice_key = key_to_str([subslice_context_uuid, subslice_uuid])
                db_subslice : SliceModel = get_object(self.database, SliceModel, str_subslice_key)

                str_slice_subslice_key = key_to_str([str_slice_key, str_subslice_key], separator='--')
                result : Tuple[SliceSubSliceModel, bool] = get_or_create_object(
                    self.database, SliceSubSliceModel, str_slice_subslice_key, {
                        'slice_fk': db_slice, 'sub_slice_fk': db_subslice})
                #db_slice_subslice, slice_subslice_created = result

            event_type = EventTypeEnum.EVENTTYPE_UPDATE if updated else EventTypeEnum.EVENTTYPE_CREATE
            dict_slice_id = db_slice.dump_id()
            notify_event(self.messagebroker, TOPIC_SLICE, event_type, {'slice_id': dict_slice_id})
            return SliceId(**dict_slice_id)

    @safe_and_metered_rpc_method(METRICS, LOGGER)
    def RemoveSlice(self, request: SliceId, context : grpc.ServicerContext) -> Empty:
        with self.lock:
            context_uuid = request.context_id.context_uuid.uuid
            slice_uuid = request.slice_uuid.uuid
            db_slice = SliceModel(self.database, key_to_str([context_uuid, slice_uuid]), auto_load=False)
            found = db_slice.load()
            if not found: return Empty()

            dict_slice_id = db_slice.dump_id()

            for db_slice_endpoint_pk,_ in db_slice.references(SliceEndPointModel):
                SliceEndPointModel(self.database, db_slice_endpoint_pk).delete()

            db_constraints = ConstraintsModel(self.database, db_slice.slice_constraints_fk)
            for db_constraint_pk,_ in db_constraints.references(ConstraintModel):
                ConstraintModel(self.database, db_constraint_pk).delete()

            for db_slice_service_pk,_ in db_slice.references(SliceServiceModel):
                SliceServiceModel(self.database, db_slice_service_pk).delete()

            for db_slice_subslice_pk,_ in db_slice.references(SliceSubSliceModel):
                SliceSubSliceModel(self.database, db_slice_subslice_pk).delete()

            db_slice.delete()
            db_constraints.delete()

            event_type = EventTypeEnum.EVENTTYPE_REMOVE
            notify_event(self.messagebroker, TOPIC_SLICE, event_type, {'slice_id': dict_slice_id})
            return Empty()

    @safe_and_metered_rpc_method(METRICS, LOGGER)
    def GetSliceEvents(self, request: Empty, context : grpc.ServicerContext) -> Iterator[SliceEvent]:
        for message in self.messagebroker.consume({TOPIC_SLICE}, consume_timeout=CONSUME_TIMEOUT):
            yield SliceEvent(**json.loads(message.content))


    # ----- Connection -------------------------------------------------------------------------------------------------

    @safe_and_metered_rpc_method(METRICS, LOGGER)