Commit 3cdbe036 authored by Lluis Gifre Renom's avatar Lluis Gifre Renom
Browse files

Context:

- corrected report of config rules updated
- corrected update notifications for Device
- removed unneeded log messages
- migrated events for Link entity
parent 5f50df51
Loading
Loading
Loading
Loading
+11 −10
Original line number Diff line number Diff line
@@ -165,28 +165,29 @@ class ContextServiceServicerImpl(ContextServiceServicer, ContextPolicyServiceSer

    @safe_and_metered_rpc_method(METRICS_POOL, LOGGER)
    def ListLinkIds(self, request : Empty, context : grpc.ServicerContext) -> LinkIdList:
        return link_list_ids(self.db_engine)
        return LinkIdList(link_ids=link_list_ids(self.db_engine))

    @safe_and_metered_rpc_method(METRICS_POOL, LOGGER)
    def ListLinks(self, request : Empty, context : grpc.ServicerContext) -> LinkList:
        return link_list_objs(self.db_engine)
        return LinkList(links=link_list_objs(self.db_engine))

    @safe_and_metered_rpc_method(METRICS_POOL, LOGGER)
    def GetLink(self, request : LinkId, context : grpc.ServicerContext) -> Link:
        return link_get(self.db_engine, request)
        return Link(**link_get(self.db_engine, request))

    @safe_and_metered_rpc_method(METRICS_POOL, LOGGER)
    def SetLink(self, request : Link, context : grpc.ServicerContext) -> LinkId:
        link_id,updated = link_set(self.db_engine, request) # pylint: disable=unused-variable
        #event_type = EventTypeEnum.EVENTTYPE_UPDATE if updated else EventTypeEnum.EVENTTYPE_CREATE
        #notify_event(self.messagebroker, TOPIC_LINK, event_type, {'link_id': link_id})
        return link_id
        link_id,updated = link_set(self.db_engine, request)
        event_type = EventTypeEnum.EVENTTYPE_UPDATE if updated else EventTypeEnum.EVENTTYPE_CREATE
        notify_event(self.messagebroker, TOPIC_LINK, event_type, {'link_id': link_id})
        return LinkId(**link_id)

    @safe_and_metered_rpc_method(METRICS_POOL, LOGGER)
    def RemoveLink(self, request : LinkId, context : grpc.ServicerContext) -> Empty:
        deleted = link_delete(self.db_engine, request) # pylint: disable=unused-variable
        #if deleted:
        #    notify_event(self.messagebroker, TOPIC_LINK, EventTypeEnum.EVENTTYPE_REMOVE, {'link_id': request})
        link_id,deleted = link_delete(self.db_engine, request)
        if deleted:
            event_type = EventTypeEnum.EVENTTYPE_REMOVE
            notify_event(self.messagebroker, TOPIC_LINK, event_type, {'link_id': link_id})
        return Empty()

    @safe_and_metered_rpc_method(METRICS_POOL, LOGGER)
+3 −5
Original line number Diff line number Diff line
@@ -59,7 +59,7 @@ def upsert_config_rules(
    if slice_uuid   is not None: stmt = stmt.where(ConfigRuleModel.slice_uuid   == slice_uuid  )
    session.execute(stmt)

    updated = False
    configrule_updates = []
    if len(config_rules) > 0:
        stmt = insert(ConfigRuleModel).values(config_rules)
        #stmt = stmt.on_conflict_do_update(
@@ -69,11 +69,9 @@ def upsert_config_rules(
        #    )
        #)
        stmt = stmt.returning(ConfigRuleModel.created_at, ConfigRuleModel.updated_at)
        config_rule_updates = session.execute(stmt).fetchall()
        LOGGER.warning('config_rule_updates = {:s}'.format(str(config_rule_updates)))
        # TODO: updated = ...
        configrule_updates = session.execute(stmt).fetchall()

    return updated
    return configrule_updates

#Union_SpecificConfigRule = Union[
#    ConfigRuleCustomModel, ConfigRuleAclModel
+3 −2
Original line number Diff line number Diff line
@@ -148,13 +148,14 @@ def device_set(db_engine : Engine, request : Device) -> Tuple[Dict, bool]:
        )
        stmt = stmt.returning(EndPointModel.created_at, EndPointModel.updated_at)
        endpoint_updates = session.execute(stmt).fetchall()
        LOGGER.warning('endpoint_updates = {:s}'.format(str(endpoint_updates)))
        updated = updated or any([(updated_at > created_at) for created_at,updated_at in endpoint_updates])

        session.execute(insert(TopologyDeviceModel).values(related_topologies).on_conflict_do_nothing(
            index_elements=[TopologyDeviceModel.topology_uuid, TopologyDeviceModel.device_uuid]
        ))

        configrules_updated = upsert_config_rules(session, config_rules, device_uuid=device_uuid)
        configrule_updates = upsert_config_rules(session, config_rules, device_uuid=device_uuid)
        updated = updated or any([(updated_at > created_at) for created_at,updated_at in configrule_updates])

        return updated

+34 −21
Original line number Diff line number Diff line
@@ -12,12 +12,13 @@
# See the License for the specific language governing permissions and
# limitations under the License.

import datetime, logging
from sqlalchemy.dialects.postgresql import insert
from sqlalchemy.engine import Engine
from sqlalchemy.orm import Session, sessionmaker
from sqlalchemy_cockroachdb import run_transaction
from typing import Dict, List, Optional, Set, Tuple
from common.proto.context_pb2 import Link, LinkId, LinkIdList, LinkList
from common.proto.context_pb2 import Link, LinkId
from common.method_wrappers.ServiceExceptions import NotFoundException
from common.tools.object_factory.Link import json_link_id
from .models.LinkModel import LinkModel, LinkEndPointModel
@@ -25,21 +26,21 @@ from .models.TopologyModel import TopologyLinkModel
from .uuids.EndPoint import endpoint_get_uuid
from .uuids.Link import link_get_uuid

def link_list_ids(db_engine : Engine) -> LinkIdList:
LOGGER = logging.getLogger(__name__)

def link_list_ids(db_engine : Engine) -> List[Dict]:
    def callback(session : Session) -> List[Dict]:
        obj_list : List[LinkModel] = session.query(LinkModel).all()
        #.options(selectinload(LinkModel.topology)).filter_by(context_uuid=context_uuid).one_or_none()
        return [obj.dump_id() for obj in obj_list]
    return LinkIdList(link_ids=run_transaction(sessionmaker(bind=db_engine), callback))
    return run_transaction(sessionmaker(bind=db_engine), callback)

def link_list_objs(db_engine : Engine) -> LinkList:
def link_list_objs(db_engine : Engine) -> List[Dict]:
    def callback(session : Session) -> List[Dict]:
        obj_list : List[LinkModel] = session.query(LinkModel).all()
        #.options(selectinload(LinkModel.topology)).filter_by(context_uuid=context_uuid).one_or_none()
        return [obj.dump() for obj in obj_list]
    return LinkList(links=run_transaction(sessionmaker(bind=db_engine), callback))
    return run_transaction(sessionmaker(bind=db_engine), callback)

def link_get(db_engine : Engine, request : LinkId) -> Link:
def link_get(db_engine : Engine, request : LinkId) -> Dict:
    link_uuid = link_get_uuid(request, allow_random=False)
    def callback(session : Session) -> Optional[Dict]:
        obj : Optional[LinkModel] = session.query(LinkModel).filter_by(link_uuid=link_uuid).one_or_none()
@@ -50,14 +51,16 @@ def link_get(db_engine : Engine, request : LinkId) -> Link:
        raise NotFoundException('Link', raw_link_uuid, extra_details=[
            'link_uuid generated was: {:s}'.format(link_uuid)
        ])
    return Link(**obj)
    return obj

def link_set(db_engine : Engine, request : Link) -> Tuple[LinkId, bool]:
def link_set(db_engine : Engine, request : Link) -> Tuple[Dict, bool]:
    raw_link_uuid = request.link_id.link_uuid.uuid
    raw_link_name = request.name
    link_name = raw_link_uuid if len(raw_link_name) == 0 else raw_link_name
    link_uuid = link_get_uuid(request.link_id, link_name=link_name, allow_random=True)

    now = datetime.datetime.utcnow()

    topology_uuids : Set[str] = set()
    related_topologies : List[Dict] = list()
    link_endpoints_data : List[Dict] = list()
@@ -80,16 +83,24 @@ def link_set(db_engine : Engine, request : Link) -> Tuple[LinkId, bool]:
    link_data = [{
        'link_uuid' : link_uuid,
        'link_name' : link_name,
        'created_at': now,
        'updated_at': now,
    }]

    def callback(session : Session) -> None:
    def callback(session : Session) -> bool:
        stmt = insert(LinkModel).values(link_data)
        stmt = stmt.on_conflict_do_update(
            index_elements=[LinkModel.link_uuid],
            set_=dict(link_name = stmt.excluded.link_name)
            set_=dict(
                link_name  = stmt.excluded.link_name,
                updated_at = stmt.excluded.updated_at,
            )
        session.execute(stmt)
        )
        stmt = stmt.returning(LinkModel.created_at, LinkModel.updated_at)
        created_at,updated_at = session.execute(stmt).fetchone()
        updated = updated_at > created_at

        # TODO: manage add/remove of endpoints; manage changes in relations with topology
        stmt = insert(LinkEndPointModel).values(link_endpoints_data)
        stmt = stmt.on_conflict_do_nothing(
            index_elements=[LinkEndPointModel.link_uuid, LinkEndPointModel.endpoint_uuid]
@@ -100,13 +111,15 @@ def link_set(db_engine : Engine, request : Link) -> Tuple[LinkId, bool]:
            index_elements=[TopologyLinkModel.topology_uuid, TopologyLinkModel.link_uuid]
        ))

    run_transaction(sessionmaker(bind=db_engine), callback)
    updated = False # TODO: improve and check if created/updated
    return LinkId(**json_link_id(link_uuid)),updated
        return updated

def link_delete(db_engine : Engine, request : LinkId) -> bool:
    updated = run_transaction(sessionmaker(bind=db_engine), callback)
    return json_link_id(link_uuid),updated

def link_delete(db_engine : Engine, request : LinkId) -> Tuple[Dict, bool]:
    link_uuid = link_get_uuid(request, allow_random=False)
    def callback(session : Session) -> bool:
        num_deleted = session.query(LinkModel).filter_by(link_uuid=link_uuid).delete()
        return num_deleted > 0
    return run_transaction(sessionmaker(bind=db_engine), callback)
    deleted = run_transaction(sessionmaker(bind=db_engine), callback)
    return json_link_id(link_uuid),deleted
+3 −1
Original line number Diff line number Diff line
@@ -12,7 +12,7 @@
# See the License for the specific language governing permissions and
# limitations under the License.

from sqlalchemy import Column, ForeignKey, String
from sqlalchemy import Column, DateTime, ForeignKey, String
from sqlalchemy.dialects.postgresql import UUID
from sqlalchemy.orm import relationship
from typing import Dict
@@ -23,6 +23,8 @@ class LinkModel(_Base):

    link_uuid  = Column(UUID(as_uuid=False), primary_key=True)
    link_name  = Column(String, nullable=False)
    created_at = Column(DateTime)
    updated_at = Column(DateTime)

    #topology_links = relationship('TopologyLinkModel', back_populates='link')
    link_endpoints = relationship('LinkEndPointModel') # lazy='joined', back_populates='link'
Loading