diff --git a/config/cgm_worker/merger.properties b/config/cgm_worker/merger.properties index 933a2424..66458f45 100644 --- a/config/cgm_worker/merger.properties +++ b/config/cgm_worker/merger.properties @@ -10,4 +10,6 @@ OPDE_MODELS_ELK_INDEX = emfos-opde-models ENABLE_DYNAMIC_MERGE_SETTINGS = True MERGE_LOAD_FLOW_SETTINGS = CGM_DEFAULT MERGE_LOAD_FLOW_SETTINGS_PRIORITY = CGM_DEFAULT,CGM_RELAXED_1,CGM_RELAXED_2 -REMOVE_GENERATORS_FROM_SLACK_DISTRIBUTION = True \ No newline at end of file +REMOVE_GENERATORS_FROM_SLACK_DISTRIBUTION = True +QAS_EIC = QAS_EIC +QAS_MSG_TYPE = QAS_MSG_TYPE \ No newline at end of file diff --git a/config/task_generator/process_conf.json b/config/task_generator/process_conf.json index 1c267498..7b3add26 100644 --- a/config/task_generator/process_conf.json +++ b/config/task_generator/process_conf.json @@ -39,7 +39,9 @@ "scaling": "False", "upload_to_opdm": "True", "upload_to_minio": "True", - "send_merge_report": "True" + "send_merge_report": "True", + "force_outage_fix": "False", + "lvl8_reporting" : "True" } }, { @@ -73,7 +75,9 @@ "scaling": "False", "upload_to_opdm": "True", "upload_to_minio": "True", - "send_merge_report": "True" + "send_merge_report": "True", + "force_outage_fix": "False", + "lvl8_reporting" : "True" } }, { @@ -107,7 +111,9 @@ "scaling": "False", "upload_to_opdm": "True", "upload_to_minio": "True", - "send_merge_report": "True" + "send_merge_report": "True", + "force_outage_fix": "False", + "lvl8_reporting" : "True" } }, { @@ -141,7 +147,9 @@ "scaling": "False", "upload_to_opdm": "True", "upload_to_minio": "True", - "send_merge_report": "True" + "send_merge_report": "True", + "force_outage_fix": "False", + "lvl8_reporting" : "True" } }, { @@ -175,7 +183,9 @@ "scaling": "False", "upload_to_opdm": "True", "upload_to_minio": "True", - "send_merge_report": "True" + "send_merge_report": "True", + "force_outage_fix": "False", + "lvl8_reporting" : "True" } } ] @@ -227,7 +237,9 @@ "scaling": "False", "upload_to_opdm": "False", "upload_to_minio": "True", - "send_merge_report": "True" + "send_merge_report": "True", + "force_outage_fix": "False", + "lvl8_reporting" : "True" } }, { @@ -268,7 +280,9 @@ "scaling": "False", "upload_to_opdm": "False", "upload_to_minio": "True", - "send_merge_report": "True" + "send_merge_report": "True", + "force_outage_fix": "False", + "lvl8_reporting" : "True" } }, { @@ -309,7 +323,9 @@ "scaling": "False", "upload_to_opdm": "False", "upload_to_minio": "True", - "send_merge_report": "True" + "send_merge_report": "True", + "force_outage_fix": "False", + "lvl8_reporting" : "True" } }, { @@ -349,7 +365,9 @@ "scaling": "False", "upload_to_opdm": "False", "upload_to_minio": "True", - "send_merge_report": "True" + "send_merge_report": "True", + "force_outage_fix": "False", + "lvl8_reporting" : "True" } }, { @@ -390,7 +408,9 @@ "scaling": "False", "upload_to_opdm": "False", "upload_to_minio": "True", - "send_merge_report": "True" + "send_merge_report": "True", + "force_outage_fix": "False", + "lvl8_reporting" : "True" } }, { @@ -431,7 +451,9 @@ "scaling": "False", "upload_to_opdm": "False", "upload_to_minio": "True", - "send_merge_report": "True" + "send_merge_report": "True", + "force_outage_fix": "False", + "lvl8_reporting" : "True" } }, { @@ -472,7 +494,9 @@ "scaling": "False", "upload_to_opdm": "False", "upload_to_minio": "True", - "send_merge_report": "True" + "send_merge_report": "True", + "force_outage_fix": "False", + "lvl8_reporting" : "True" } }, @@ -517,7 +541,9 @@ "scaling": "False", "upload_to_opdm": "False", "upload_to_minio": "True", - "send_merge_report": "True" + "send_merge_report": "True", + "force_outage_fix": "False", + "lvl8_reporting" : "True" } }, @@ -562,7 +588,9 @@ "scaling": "False", "upload_to_opdm": "False", "upload_to_minio": "True", - "send_merge_report": "True" + "send_merge_report": "True", + "force_outage_fix": "False", + "lvl8_reporting" : "True" } } diff --git a/emf/model_merger/merge_functions.py b/emf/model_merger/merge_functions.py index 5a65dee1..772b6285 100644 --- a/emf/model_merger/merge_functions.py +++ b/emf/model_merger/merge_functions.py @@ -9,6 +9,7 @@ import triplets import uuid import config +import xml.etree.ElementTree as ET from emf.common.config_parser import parse_app_properties from emf.common.integrations import elastic from emf.model_merger import temporary @@ -16,6 +17,7 @@ from emf.common.helpers.loadflow import get_model_outages, get_network_elements from emf.common.helpers.opdm_objects import load_opdm_objects_to_triplets, filename_from_opdm_metadata + logger = logging.getLogger(__name__) parse_app_properties(caller_globals=globals(), path=config.paths.cgm_worker.post_processing) @@ -857,6 +859,93 @@ def run_post_merge_processing(input_models: list, return sv_data, ssh_data, opdm_object_meta +def lvl8_report_cgm(merge_report: dict): + + # Create root + qa_attribs = { + 'created': datetime.datetime.strptime(merge_report["@timestamp"], '%Y-%m-%dT%H:%M:%S.%f').strftime('%Y-%m-%dT%H:%M:%SZ'), + 'schemeVersion': "2.0", + 'serviceProvider': merge_report["merge_entity"], + 'xmlns': "http://entsoe.eu/checks" + } + qa_root = ET.Element("QAReport", attrib=qa_attribs) + + # Add RuleViolations + violations_list = [ + { + 'ruleId': "CGMConvergence", + 'validationLevel': "8", + 'severity': "WARNING", + 'Message': "Power flow could not be calculated for CGM with default settings." + }, + { + 'ruleId': "CGMConvergenceRelaxed", + 'validationLevel': "8", + 'severity': "ERROR", + 'Message': "Power flow could not be calculated for CGM with CGM_RELAXED_2 settings." + } + ] + # TODO:pick the correct setting based on retruned LF setting and convergance from model. Set model quality indicator based on violations + violations = [] + if merge_report["loadflow_status"] == 'CONVERGED': + if merge_report["loadflow_settings"] == 'CGM_DEFAULT': + logger.info(f"Merge successful with default settings included in lvl8 report") + quality_indicator_cgm = "Valid" + else: + violations.append(violations_list[0]) + quality_indicator_cgm = "Warning - non fatal inconsistencies" + else: + violations = violations.append(violations_list) + quality_indicator_cgm = "Invalid - inconsistent data" + + # Create + cgm_attribs = { + 'created': datetime.datetime.strptime(merge_report["@timestamp"], '%Y-%m-%dT%H:%M:%S.%f').strftime('%Y-%m-%dT%H:%M:%SZ'), + 'resource': merge_report['network_meta']['fullModel_ID'], # TODO get here correct content ID + 'scenarioTime': datetime.datetime.fromisoformat(merge_report["@scenario_timestamp"]).strftime('%Y-%m-%dT%H:%M:%SZ'), + 'version': str(merge_report["@version"]), + 'processType': merge_report["@time_horizon"], + 'qualityIndicator': quality_indicator_cgm + } + cgm = ET.SubElement(qa_root, "CGM", attrib=cgm_attribs) + + try: + for v in violations: + rv = ET.SubElement(cgm, "RuleViolation", { + 'ruleId': v['ruleId'], + 'validationLevel': v['validationLevel'], + 'severity': v['severity'] + }) + msg = ET.SubElement(rv, "Message") + msg.text = v['Message'] + except: + logger.info(f"No violations present in merge") + + # TODO:pick the TSOs from QA report. Missing parameters below for all IGMs + for i in merge_report['merge_included_entity'] + merge_report['replaced_entity']: + igm = ET.SubElement(cgm, "IGM", { + 'created': i["creation_timestamp"], + 'scenarioTime': datetime.datetime.fromisoformat(i['scenario_timestamp']).strftime('%Y-%m-%dT%H:%M:%SZ'), + 'tso': i['tso'], + 'version': str(i['version']), + 'processType': i['time_horizon'], + 'qualityIndicator': i['quality_indicator'], + }) + resource_igm = ET.SubElement(igm, "resource") + resource_igm.text = i['model_sv_id'] + + # Add EMFInformation + ET.SubElement(cgm, "EMFInformation", { + 'mergingEntity': merge_report["merge_entity"], + 'cgmType': merge_report["merge_type"] + }) + + # Generate final XML + qa_report_lvl8 = ET.tostring(qa_root, encoding='utf-8', xml_declaration=True) + + return qa_report_lvl8 + + if __name__ == "__main__": from emf.common.integrations.object_storage.models import get_latest_boundary, get_latest_models_and_download diff --git a/emf/model_merger/model_merger.py b/emf/model_merger/model_merger.py index 0386bf0c..bfe9a77d 100644 --- a/emf/model_merger/model_merger.py +++ b/emf/model_merger/model_merger.py @@ -11,7 +11,7 @@ from io import BytesIO from zipfile import ZipFile from emf.common.config_parser import parse_app_properties -from emf.common.integrations import opdm, minio_api, elastic +from emf.common.integrations import opdm, minio_api, elastic, edx from emf.common.integrations.object_storage.models import get_latest_boundary, get_latest_models_and_download from emf.common.integrations.object_storage.schedules import query_acnp_schedules, query_hvdc_schedules from emf.common.loadflow_tool import loadflow_settings @@ -28,7 +28,6 @@ from concurrent.futures import ThreadPoolExecutor from lxml import etree - logger = logging.getLogger(__name__) parse_app_properties(caller_globals=globals(), path=config.paths.cgm_worker.merger) executor = ThreadPoolExecutor(max_workers=20) @@ -68,6 +67,32 @@ class MergedModel: replacement_reason: List = field(default_factory=list) outages_updated: List = field(default_factory=list) outages_unmapped: List = field(default_factory=list) + merge_included_entity: List = field(default_factory=list) + + +@dataclass(init=False) +class ModelEntity: + data_source: str = "OPDM" + quality_indicator: str = "Valid" + tso: str = None + time_horizon: str = None + scenario_timestamp: str = None + model_sv_id: str = None + version: int = None + quality_indicator: str = "Valid" + creation_timestamp: str = None + file_name: str = None + + def __init__(self, data_source: str, quality_indicator: str, **kwargs): + self.data_source = data_source + self.quality_indicator = quality_indicator + self.tso = kwargs.get('pmd:TSO', 'unknown') + self.time_horizon = kwargs.get('pmd:timeHorizon', 'unknown') + self.scenario_timestamp = kwargs.get('pmd:scenarioDate', 'unknown') + self.model_sv_id = kwargs.get('pmd:fullModel_ID', 'unknown') + self.version = int(kwargs.get('pmd:version', 999)) + self.creation_timestamp = kwargs.get('pmd:creationDate', 'unknown') + self.file_name = kwargs.get('pmd:fileName', 'unknown') class HandlerMergeModels: @@ -82,7 +107,8 @@ def run_loadflow(merged_model): # Set starting point of lf settings priority list if json.loads(ENABLE_DYNAMIC_MERGE_SETTINGS.lower()): settings_list = [param.strip() for param in MERGE_LOAD_FLOW_SETTINGS_PRIORITY.split(",")] - settings_priority = next((i for i, value in enumerate(settings_list) if value == MERGE_LOAD_FLOW_SETTINGS), None) + settings_priority = next((i for i, value in enumerate(settings_list) if value == MERGE_LOAD_FLOW_SETTINGS), + None) settings_list = settings_list[settings_priority:] else: settings_list = [MERGE_LOAD_FLOW_SETTINGS] @@ -165,6 +191,7 @@ def handle(self, task_object: dict, properties: dict, **kwargs): model_merge_report_send_to_elk = task_properties["send_merge_report"] post_temp_fixes = task_properties['post_temp_fixes'] force_outage_fix = task_properties['force_outage_fix'] + lvl8_reporting = task_properties['lvl8_reporting'] # Collect valid models from ObjectStorage downloaded_models = get_latest_models_and_download(time_horizon=time_horizon, @@ -179,6 +206,8 @@ def handle(self, task_object: dict, properties: dict, **kwargs): excluded_models=excluded_models, filter_on='pmd:TSO') + merged_model.merge_included_entity = [ModelEntity(data_source='OPDM', quality_indicator='Valid', **model).__dict__ for model in models] + # Get additional models from ObjectStorage if local import is configured if local_import_models: additional_models = get_latest_models_and_download(time_horizon=time_horizon, @@ -188,8 +217,11 @@ def handle(self, task_object: dict, properties: dict, **kwargs): additional_models = merge_functions.filter_models(models=additional_models, included_models=local_import_models, filter_on='pmd:TSO') + merged_model.merge_included_entity.extend( + [ModelEntity(data_source='PDN', quality_indicator='Valid', **model).__dict__ for model in additional_models]) - missing_local_import = [tso for tso in local_import_models if tso not in [model['pmd:TSO'] for model in additional_models]] + missing_local_import = [tso for tso in local_import_models if + tso not in [model['pmd:TSO'] for model in additional_models]] merged_model.excluded.extend([{'tso': tso, 'reason': 'missing-pdn'} for tso in missing_local_import]) # Perform local replacement if configured @@ -201,10 +233,10 @@ def handle(self, task_object: dict, properties: dict, **kwargs): scenario_date=scenario_datetime, data_source='PDN') - logger.info(f"Local storage replacement model(s) found: {[model['pmd:fileName'] for model in replacement_models_local]}") - replaced_entities_local = [{'tso': model['pmd:TSO'], - 'time_horizon': model['pmd:timeHorizon'], - 'scenario_timestamp': model['pmd:scenarioDate']} for model in replacement_models_local] + logger.info( + f"Local storage replacement model(s) found: {[model['pmd:fileName'] for model in replacement_models_local]}") + replaced_entities_local = [ModelEntity(data_source='PDN', quality_indicator='Substituted', **model).__dict__ for model in + replacement_models_local] merged_model.replaced_entity.extend(replaced_entities_local) additional_models.extend(replacement_models_local) except Exception as error: @@ -233,11 +265,10 @@ def handle(self, task_object: dict, properties: dict, **kwargs): logger.info(f"Running replacement for missing models: {missing_models}") replacement_models = run_replacement(missing_models, time_horizon, scenario_datetime) if replacement_models: - logger.info(f"Replacement model(s) found: {[model['pmd:fileName'] for model in replacement_models]}") - replaced_entities = [{'tso': model['pmd:TSO'], - 'time_horizon': model['pmd:timeHorizon'], - 'scenario_timestamp': model['pmd:scenarioDate']} - for model in replacement_models] + logger.info( + f"Replacement model(s) found: {[model['pmd:fileName'] for model in replacement_models]}") + replaced_entities = [ModelEntity(data_source='OPDM', quality_indicator='Substituted', **model).__dict__ for model in + replacement_models] merged_model.replaced_entity.extend(replaced_entities) models.extend(replacement_models) merged_model.replaced = True @@ -267,7 +298,8 @@ def handle(self, task_object: dict, properties: dict, **kwargs): tso_list = [] if force_outage_fix: # force outage fix on all models if set tso_list = merged_model.included - elif merging_area == 'BA' and any(tso in ['LITGRID', 'AST', 'ELERING'] for tso in replaced_tso_list): # by default do it on Baltic merge replaced models + elif merging_area == 'BA' and any(tso in ['LITGRID', 'AST', 'ELERING'] for tso in + replaced_tso_list): # by default do it on Baltic merge replaced models tso_list = replaced_tso_list if tso_list: # if not set force and not replaced BA then nothing to fix merged_model = merge_functions.update_model_outages(merged_model=merged_model, @@ -280,7 +312,8 @@ def handle(self, task_object: dict, properties: dict, **kwargs): merged_model.network = handle_igm_ssh_vs_cgm_ssh_error(network_pre_instance=merged_model.network) # Ensure boundary point EquivalentInjection are set to zero for paired tie lines - merged_model.network = merge_functions.ensure_paired_equivalent_injection_compatibility(network=merged_model.network) + merged_model.network = merge_functions.ensure_paired_equivalent_injection_compatibility( + network=merged_model.network) # Ensure boundary line connectivity consistency for paired boundary lines merged_model.network = merge_functions.ensure_paired_boundary_line_connectivity(network=merged_model.network) @@ -288,7 +321,8 @@ def handle(self, task_object: dict, properties: dict, **kwargs): # TODO - run other LF if default fails # Run loadflow on merged model merged_model = self.run_loadflow(merged_model=merged_model) - logger.info(f"Loadflow status of main island: {merged_model.loadflow_status} [settings: {merged_model.loadflow_settings}]") + logger.info( + f"Loadflow status of main island: {merged_model.loadflow_status} [settings: {merged_model.loadflow_settings}]") # Perform scaling if model_scaling: @@ -310,7 +344,8 @@ def handle(self, task_object: dict, properties: dict, **kwargs): merged_model = scaler.scale_balance(model=merged_model, ac_schedules=ac_schedules, dc_schedules=dc_schedules, - lf_settings=getattr(loadflow_settings, merged_model.loadflow_settings)) + lf_settings=getattr(loadflow_settings, + merged_model.loadflow_settings)) except Exception as e: logger.error(e) merged_model.scaled = False @@ -359,6 +394,8 @@ def handle(self, task_object: dict, properties: dict, **kwargs): task_properties=task_properties ) + # for merge report need to get the final uuid. + merged_model.network_meta['fullModel_ID'] = opdm_object_meta['pmd:fullModel_ID'] # Package both input models and exported CGM profiles to in memory zip files serialized_data = export_to_cgmes_zip([ssh_data, sv_data]) post_p_end = datetime.datetime.now(datetime.UTC) @@ -429,7 +466,8 @@ def handle(self, task_object: dict, properties: dict, **kwargs): opdm_object_meta['pmd:content-reference'] = merged_model.content_reference response = elastic.Elastic.send_to_elastic(index=OPDE_MODELS_ELK_INDEX, json_message=opdm_object_meta) - # Send merge report and opdm object metadata to Elastic + # Send merge report and OPDM object metadata to Elastic + merge_report = None if model_merge_report_send_to_elk: logger.info(f"Sending merge report to Elastic") try: @@ -441,6 +479,20 @@ def handle(self, task_object: dict, properties: dict, **kwargs): except Exception as error: logger.error(f"Failed to create merge report: {error}") + # Send QAS level 8 report if configured + if lvl8_reporting and merge_report: + try: + lvl8_report = merge_functions.lvl8_report_cgm(merge_report=merge_report) + service_edx = edx.EDX() + message_id = service_edx.send_message(receiver_EIC=QAS_EIC, + business_type=QAS_MSG_TYPE, + content=lvl8_report) + logger.info(f"QAS-Level-8 report generated and sent with ID: {message_id}") + except Exception as error: + logger.error(f"Failed to send QAS-Level-8 report with error: {error}") + else: + logger.warning(f"QAS-Level-8 not generated because merge report unavailable or disabled by configuration") + # Append message headers with OPDM root metadata extracted_meta = {key: value for key, value in opdm_object_meta.items() if isinstance(value, str)} properties.headers = extracted_meta @@ -454,7 +506,6 @@ def handle(self, task_object: dict, properties: dict, **kwargs): if __name__ == "__main__": - sample_task = { "@context": "https://example.com/task_context.jsonld", "@type": "Task", @@ -483,7 +534,7 @@ def handle(self, task_object: dict, properties: dict, **kwargs): "job_period_start": "2024-05-24T22:00:00+00:00", "job_period_end": "2024-05-25T06:00:00+00:00", "task_properties": { - "timestamp_utc": "2025-07-06T12:30:00+00:00", + "timestamp_utc": "2025-07-31T12:30:00+00:00", "merge_type": "BA", "merging_entity": "BALTICRCC", "included": ["LITGRID", "AST", "ELERING", "PSE"], @@ -500,6 +551,7 @@ def handle(self, task_object: dict, properties: dict, **kwargs): "upload_to_minio": "True", "send_merge_report": "False", "force_outage_fix": "False", + "lvl8_reporting": "True" } }