Loading manifests/automationservice.yaml +5 −0 Viewed Changes for manifests/automationservice.yaml: 5 added lines, 0 removed lines. Original line number Diff line number Diff line Loading @@ -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"] Loading proto/automation.proto +8 −7 Viewed Changes for proto/automation.proto: 8 added lines, 7 removed lines. Original line number Diff line number Diff line Loading @@ -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 ) {} Loading @@ -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 Loading src/analytics/backend/service/AnalyzerHandlers.py +49 −0 Viewed Changes for src/analytics/backend/service/AnalyzerHandlers.py: 49 added lines, 0 removed lines. Original line number Diff line number Diff line Loading @@ -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): """ Loading Loading @@ -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 src/analytics/backend/service/Streamer.py +5 −3 Viewed Changes for src/analytics/backend/service/Streamer.py: 5 added lines, 3 removed lines. Original line number Diff line number Diff line Loading @@ -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 Loading Loading @@ -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: Loading src/automation/service/database/models/_Base.py 0 → 100644 +25 −0 Viewed 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
manifests/automationservice.yaml +5 −0 Viewed Changes for manifests/automationservice.yaml: 5 added lines, 0 removed lines. Original line number Diff line number Diff line Loading @@ -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"] Loading
proto/automation.proto +8 −7 Viewed Changes for proto/automation.proto: 8 added lines, 7 removed lines. Original line number Diff line number Diff line Loading @@ -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 ) {} Loading @@ -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 Loading
src/analytics/backend/service/AnalyzerHandlers.py +49 −0 Viewed Changes for src/analytics/backend/service/AnalyzerHandlers.py: 49 added lines, 0 removed lines. Original line number Diff line number Diff line Loading @@ -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): """ Loading Loading @@ -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
src/analytics/backend/service/Streamer.py +5 −3 Viewed Changes for src/analytics/backend/service/Streamer.py: 5 added lines, 3 removed lines. Original line number Diff line number Diff line Loading @@ -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 Loading Loading @@ -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: Loading
src/automation/service/database/models/_Base.py 0 → 100644 +25 −0 Viewed 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)