feat: refactor automation for supporting zsm plugin handler logic.

Closes issue #267 (closed)

Edited by Georgios P. Katsikas

Merge request reports

Loading
+5 −0
Changes for manifests/automationservice.yaml: 5 added lines, 0 removed lines.
Original line number Diff line number Diff line
@@ -40,6 +40,11 @@ spec:
          env:
            - name: LOG_LEVEL
              value: "INFO"
            - name: CRDB_DATABASE
              value: "tfs_automation"
          envFrom:
            - secretRef:
                name: crdb-data
          startupProbe:
            exec:
              command: ["/bin/grpc_health_probe", "-addr=:30200"]
+8 −7
Changes for proto/automation.proto: 8 added lines, 7 removed lines.
Original line number Diff line number Diff line
@@ -17,11 +17,11 @@ package automation;

import "context.proto";
import "policy.proto";
import "analytics_frontend.proto";

// Automation service RPCs
service AutomationService {
  rpc ZSMCreate       (ZSMCreateRequest ) returns (ZSMService     ) {}
  rpc ZSMUpdate       (ZSMCreateUpdate  ) returns (ZSMService     ) {}
  rpc ZSMDelete       (ZSMServiceID     ) returns (ZSMServiceState) {}
  rpc ZSMGetById      (ZSMServiceID     ) returns (ZSMService     ) {}
  rpc ZSMGetByService (context.ServiceId) returns (ZSMService     ) {}
@@ -37,14 +37,15 @@ enum ZSMServiceStateEnum {
  ZSM_REMOVED = 5;       // ZSM loop is removed
}

message ZSMCreateRequest {
  context.ServiceId serviceId = 1;
  policy.PolicyRuleList policyList = 2;
enum ZSMTypeEnum {
  ZSMTYPE_UNKNOWN = 0;
}

message ZSMCreateUpdate {
  context.Uuid ZSMServiceID = 1;
  policy.PolicyRuleList policyList = 2;
message ZSMCreateRequest {
  context.ServiceId target_service_id = 1;
  context.ServiceId telemetry_service_id = 2;
  analytics_frontend.Analyzer analyzer = 3;
  policy.PolicyRuleService policy = 4;
}

// A unique identifier per ZSM service
+49 −0
Changes for src/analytics/backend/service/AnalyzerHandlers.py: 49 added lines, 0 removed lines.
Original line number Diff line number Diff line
@@ -15,18 +15,28 @@
import logging
from enum import Enum
import pandas as pd
from collections import defaultdict

logger = logging.getLogger(__name__)


class Handlers(Enum):
    AGGREGATION_HANDLER = "AggregationHandler"
    AGGREGATION_HANDLER_THREE_TO_ONE = "AggregationHandlerThreeToOne"
    UNSUPPORTED_HANDLER = "UnsupportedHandler"

    @classmethod
    def is_valid_handler(cls, handler_name):
        return handler_name in cls._value2member_map_

def select_handler(handler_name):
    if handler_name == "AggregationHandler":
        return aggregation_handler
    elif handler_name == "AggregationHandlerThreeToOne":
        return aggregation_handler_three_to_one
    else:
        return "UnsupportedHandler"

# This method is top-level and should not be part of the class due to serialization issues.
def threshold_handler(key, aggregated_df, thresholds):
    """
@@ -134,3 +144,42 @@ def aggregation_handler(
            return results
        else:
            return []

def find(data , type , value):
    return next((item for item in data if item[type] == value), None)


def aggregation_handler_three_to_one(
        batch_type_name, key, batch, input_kpi_list, output_kpi_list, thresholds
):

    # Group and sum
    # Track sum and count
    sum_dict = defaultdict(int)
    count_dict = defaultdict(int)

    for item in batch:
        kpi_id = item["kpi_id"]
        if kpi_id in input_kpi_list:
            sum_dict[kpi_id] += item["kpi_value"]
            count_dict[kpi_id] += 1

    # Compute average
    avg_dict = {kpi_id: sum_dict[kpi_id] / count_dict[kpi_id] for kpi_id in sum_dict}

    total_kpi_metric = 0
    for kpi_id, total_value in avg_dict.items():
        total_kpi_metric += total_value

    result = {
        "kpi_id": output_kpi_list[0],
        "avg": total_kpi_metric,
        "THRESHOLD_RAISE": bool(total_kpi_metric > 2600),
        "THRESHOLD_FALL": bool(total_kpi_metric < 699)
    }
    results = []

    results.append(result)
    logger.warning(f"result : {result}.")

    return results
+5 −3
Changes for src/analytics/backend/service/Streamer.py: 5 added lines, 3 removed lines.
Original line number Diff line number Diff line
@@ -19,7 +19,7 @@ import logging

from confluent_kafka                            import KafkaException, KafkaError
from common.tools.kafka.Variables               import KafkaTopic
from analytics.backend.service.AnalyzerHandlers import Handlers, aggregation_handler
from analytics.backend.service.AnalyzerHandlers import Handlers, aggregation_handler, aggregation_handler_three_to_one , select_handler
from analytics.backend.service.AnalyzerHelper   import AnalyzerHelper


@@ -114,11 +114,13 @@ class DaskStreamer(threading.Thread):
        if Handlers.is_valid_handler(self.thresholds["task_type"]):
            if self.client is not None and self.client.status == 'running':
                try:
                    future = self.client.submit(aggregation_handler, "batch size", self.key,
                    future = self.client.submit(select_handler(self.thresholds["task_type"]), "batch size",
                                                self.key,
                                                self.batch, self.input_kpis, self.output_kpis, self.thresholds)
                    future.add_done_callback(lambda fut: self.produce_result(fut.result(), KafkaTopic.ALARMS.value))
                except Exception as e:
                    logger.error(f"Failed to submit task to Dask client or unable to process future. See error for detail: {e}")
                    logger.error(
                        f"Failed to submit task to Dask client or unable to process future. See error for detail: {e}")
            else:
                logger.warning("Dask client is not running. Skipping processing.")
        else:
+25 −0
Changes for src/automation/service/database/models/_Base.py: 25 added lines, 0 removed lines.
Original line number Diff line number Diff line
# Copyright 2022-2025 ETSI SDG TeraFlowSDN (TFS) (https://tfs.etsi.org/)
#
# 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 sqlalchemy
from typing import Any, List
from sqlalchemy.orm import Session, sessionmaker, declarative_base
from sqlalchemy.sql import text
from sqlalchemy_cockroachdb import run_transaction

_Base = declarative_base()

def rebuild_database(db_engine : sqlalchemy.engine.Engine, drop_if_exists : bool = False):
    if drop_if_exists: _Base.metadata.drop_all(db_engine)
    _Base.metadata.create_all(db_engine)
Loading
Loading