From 78b7868301ad8445f92c00d123b026f395cfeb7a Mon Sep 17 00:00:00 2001 From: VeikoAunapuu Date: Tue, 29 Jul 2025 11:40:00 +0300 Subject: [PATCH 01/11] Func for lvl8 --- emf/model_merger/merge_functions.py | 89 +++++++++++++++++++++++++++++ 1 file changed, 89 insertions(+) diff --git a/emf/model_merger/merge_functions.py b/emf/model_merger/merge_functions.py index 5a65dee1..4eca5f8e 100644 --- a/emf/model_merger/merge_functions.py +++ b/emf/model_merger/merge_functions.py @@ -856,6 +856,95 @@ def run_post_merge_processing(input_models: list, return sv_data, ssh_data, opdm_object_meta +def lvl8_report_cgm(merge_report): + + # Create root + qa_attribs = { + 'created': "2025-07-27T16:04:34Z", + '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': "resource" ,#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['included']: + igm = ET.SubElement(cgm, "IGM", { + 'created': datetime.datetime.strptime(merge_report["@timestamp"],'%Y-%m-%dT%H:%M:%S.%f').strftime('%Y-%m-%dT%H:%M:%SZ'), + 'scenarioTime': datetime.datetime.fromisoformat(merge_report["@scenario_timestamp"]).strftime('%Y-%m-%dT%H:%M:%SZ'), + 'tso': i, + 'version': 'version', + 'processType': 'processType', + 'qualityIndicator': 'qualityIndicator', + 'resource':"resource" + }) + resource_igm= ET.SubElement(igm, "resource") + resource_igm.text="resource" + + + # Add EMFInformation + ET.SubElement(cgm, "EMFInformation", { + 'mergingEntity': merge_report["merge_entity"], + 'cgmType': merge_report["merge_type"] + }) + + # Generate final XML + qa_report_lvl8 = minidom.parseString(ET.tostring(qa_root, encoding='utf-8', xml_declaration=True)) + logger.debug(qa_report_lvl8.toprettyxml(indent=" ")) + return qa_report_lvl8 if __name__ == "__main__": From b22dda4115d64c61a6cfe62d3b0c467e1f251365 Mon Sep 17 00:00:00 2001 From: VeikoAunapuu Date: Tue, 29 Jul 2025 11:42:27 +0300 Subject: [PATCH 02/11] Integration of lvl 8 --- emf/model_merger/model_merger.py | 10 +++++++++- 1 file changed, 9 insertions(+), 1 deletion(-) diff --git a/emf/model_merger/model_merger.py b/emf/model_merger/model_merger.py index 0386bf0c..0695bbc7 100644 --- a/emf/model_merger/model_merger.py +++ b/emf/model_merger/model_merger.py @@ -165,6 +165,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, @@ -440,7 +441,13 @@ def handle(self, task_object: dict, properties: dict, **kwargs): logger.error(f"Merge report sending to Elastic failed: {error}") except Exception as error: logger.error(f"Failed to create merge report: {error}") - + + #send lvl 8 report + if lvl8_reporting: + lvl8_report = merge_functions.lvl8_report_cgm(merge_report) + logger.debug(f"lvl8 report generated: {lvl8_report}") + #TODO Sending the report via EDX. + # 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 @@ -500,6 +507,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" } } From 9bfbd58535a9e84a748726f68d9a9f1863384a8e Mon Sep 17 00:00:00 2001 From: VeikoAunapuu Date: Tue, 29 Jul 2025 16:01:33 +0300 Subject: [PATCH 03/11] Update model_merger.py --- emf/model_merger/model_merger.py | 47 ++++++++++++++++++++++++++++---- 1 file changed, 41 insertions(+), 6 deletions(-) diff --git a/emf/model_merger/model_merger.py b/emf/model_merger/model_merger.py index 0695bbc7..1a68fcc4 100644 --- a/emf/model_merger/model_merger.py +++ b/emf/model_merger/model_merger.py @@ -68,7 +68,8 @@ class MergedModel: replacement_reason: List = field(default_factory=list) outages_updated: List = field(default_factory=list) outages_unmapped: List = field(default_factory=list) - + included_opdm: List = field(default_factory=list) + included_local: List = field(default_factory=list) class HandlerMergeModels: @@ -180,6 +181,19 @@ def handle(self, task_object: dict, properties: dict, **kwargs): excluded_models=excluded_models, filter_on='pmd:TSO') + merged_model.included_opdm=[ + { + 'tso': item['pmd:TSO'], + 'time_horizon': item['pmd:timeHorizon'], + 'scenario_timestamp': item['pmd:scenarioDate'], + 'pmd:fullModel_ID': item['pmd:fullModel_ID'], + 'pmd:version': item['pmd:version'], + 'qualityIndicator': 'valid', + 'pmd:creationDate': item['pmd:creationDate'], + } + for item 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, @@ -189,7 +203,18 @@ 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.included_local = [ + { + 'tso': item['pmd:TSO'], + 'time_horizon': item['pmd:timeHorizon'], + 'scenario_timestamp': item['pmd:scenarioDate'], + 'pmd:fullModel_ID': item['pmd:fullModel_ID'], + 'pmd:version': item['pmd:version'], + 'qualityIndicator': 'valid', + 'pmd:creationDate': item['pmd:creationDate'], + } + for item 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]) @@ -205,7 +230,13 @@ def handle(self, task_object: dict, properties: dict, **kwargs): 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] + 'scenario_timestamp': model['pmd:scenarioDate'], + 'pmd:fullModel_ID': model['pmd:fullModel_ID'], + 'pmd:version': model['pmd:version'], + 'data_source' : 'PDN', + 'qualityIndicator': 'replaced', + 'pmd:creationDate': model['pmd:creationDate'] + } for model in replacement_models_local] merged_model.replaced_entity.extend(replaced_entities_local) additional_models.extend(replacement_models_local) except Exception as error: @@ -236,9 +267,13 @@ def handle(self, task_object: dict, properties: dict, **kwargs): 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] + 'time_horizon': model['pmd:timeHorizon'], + 'scenario_timestamp': model['pmd:scenarioDate'], + 'pmd:fullModel_ID': model['pmd:fullModel_ID'], + 'pmd:version': model['pmd:version'], + 'data_source' : 'OPDM', + 'qualityIndicator': 'replaced', + 'pmd:creationDate': model['pmd:creationDate'] } for model in replacement_models] merged_model.replaced_entity.extend(replaced_entities) models.extend(replacement_models) merged_model.replaced = True From b980090f16eb76facae38daaa82c5e7534f65cb4 Mon Sep 17 00:00:00 2001 From: VeikoAunapuu Date: Tue, 29 Jul 2025 16:03:49 +0300 Subject: [PATCH 04/11] Update merge_functions.py --- emf/model_merger/merge_functions.py | 19 +++++++++---------- 1 file changed, 9 insertions(+), 10 deletions(-) diff --git a/emf/model_merger/merge_functions.py b/emf/model_merger/merge_functions.py index 4eca5f8e..01fbd646 100644 --- a/emf/model_merger/merge_functions.py +++ b/emf/model_merger/merge_functions.py @@ -860,7 +860,7 @@ def lvl8_report_cgm(merge_report): # Create root qa_attribs = { - 'created': "2025-07-27T16:04:34Z", + 'created': merge_report['pmd:creationDate'], 'schemeVersion': "2.0", 'serviceProvider': merge_report["merge_entity"], 'xmlns': "http://entsoe.eu/checks" @@ -921,18 +921,17 @@ def lvl8_report_cgm(merge_report): # TODO:pick the TSOs from QA report. Missing parameters below for all IGMs - for i in merge_report['included']: + for i in merge_report['included_opdm'] + merge_report['included_local'] + merge_report['replaced_entity']: igm = ET.SubElement(cgm, "IGM", { - 'created': datetime.datetime.strptime(merge_report["@timestamp"],'%Y-%m-%dT%H:%M:%S.%f').strftime('%Y-%m-%dT%H:%M:%SZ'), - 'scenarioTime': datetime.datetime.fromisoformat(merge_report["@scenario_timestamp"]).strftime('%Y-%m-%dT%H:%M:%SZ'), - 'tso': i, - 'version': 'version', - 'processType': 'processType', - 'qualityIndicator': 'qualityIndicator', - 'resource':"resource" + 'created': i["pmd:creationDate"], + 'scenarioTime': datetime.datetime.fromisoformat(i['scenario_timestamp']).strftime('%Y-%m-%dT%H:%M:%SZ'), + 'tso': i['tso'], + 'version': i['pmd:version'], + 'processType': i['time_horizon'], + 'qualityIndicator': i['qualityIndicator'], }) resource_igm= ET.SubElement(igm, "resource") - resource_igm.text="resource" + resource_igm.text=i['pmd:fullModel_ID'] # Add EMFInformation From 1222104391dffd18483a593c22a3555eaa03df9a Mon Sep 17 00:00:00 2001 From: VeikoAunapuu Date: Wed, 30 Jul 2025 11:33:23 +0300 Subject: [PATCH 05/11] Init_working_version --- config/cgm_worker/merger.properties | 4 +- config/task_generator/process_conf.json | 56 ++++++++++++++++------ emf/model_merger/merge_functions.py | 22 +++++---- emf/model_merger/model_merger.py | 63 +++++++++++++++---------- 4 files changed, 94 insertions(+), 51 deletions(-) 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 01fbd646..a6af7f7d 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) @@ -860,9 +862,9 @@ def lvl8_report_cgm(merge_report): # Create root qa_attribs = { - 'created': merge_report['pmd:creationDate'], + '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"], + 'serviceProvider': merge_report['network_meta']['fullModel_ID'], 'xmlns': "http://entsoe.eu/checks" } qa_root = ET.Element("QAReport", attrib=qa_attribs) @@ -898,7 +900,7 @@ def lvl8_report_cgm(merge_report): # 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': "resource" ,#TODO get here correct content ID + '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"], @@ -923,15 +925,15 @@ def lvl8_report_cgm(merge_report): # TODO:pick the TSOs from QA report. Missing parameters below for all IGMs for i in merge_report['included_opdm'] + merge_report['included_local'] + merge_report['replaced_entity']: igm = ET.SubElement(cgm, "IGM", { - 'created': i["pmd:creationDate"], - 'scenarioTime': datetime.datetime.fromisoformat(i['scenario_timestamp']).strftime('%Y-%m-%dT%H:%M:%SZ'), + 'created': i["@timestamp"], + 'scenarioTime': datetime.datetime.fromisoformat(i['@scenario_timestamp']).strftime('%Y-%m-%dT%H:%M:%SZ'), 'tso': i['tso'], - 'version': i['pmd:version'], - 'processType': i['time_horizon'], + 'version': i['@version'], + 'processType': i['@time_horizon'], 'qualityIndicator': i['qualityIndicator'], }) resource_igm= ET.SubElement(igm, "resource") - resource_igm.text=i['pmd:fullModel_ID'] + resource_igm.text=i['fullModel_ID'] # Add EMFInformation @@ -941,8 +943,8 @@ def lvl8_report_cgm(merge_report): }) # Generate final XML - qa_report_lvl8 = minidom.parseString(ET.tostring(qa_root, encoding='utf-8', xml_declaration=True)) - logger.debug(qa_report_lvl8.toprettyxml(indent=" ")) + qa_report_lvl8 = ET.tostring(qa_root, encoding='utf-8', xml_declaration=True) + return qa_report_lvl8 if __name__ == "__main__": diff --git a/emf/model_merger/model_merger.py b/emf/model_merger/model_merger.py index 1a68fcc4..3bc8cab8 100644 --- a/emf/model_merger/model_merger.py +++ b/emf/model_merger/model_merger.py @@ -184,16 +184,17 @@ def handle(self, task_object: dict, properties: dict, **kwargs): merged_model.included_opdm=[ { 'tso': item['pmd:TSO'], - 'time_horizon': item['pmd:timeHorizon'], - 'scenario_timestamp': item['pmd:scenarioDate'], - 'pmd:fullModel_ID': item['pmd:fullModel_ID'], - 'pmd:version': item['pmd:version'], + '@time_horizon': item['pmd:timeHorizon'], + '@scenario_timestamp': item['pmd:scenarioDate'], + 'fullModel_ID': item['pmd:fullModel_ID'], + '@version': item['pmd:version'], 'qualityIndicator': 'valid', - 'pmd:creationDate': item['pmd:creationDate'], + '@timestamp': item['pmd:creationDate'], } for item 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, @@ -206,15 +207,16 @@ def handle(self, task_object: dict, properties: dict, **kwargs): merged_model.included_local = [ { 'tso': item['pmd:TSO'], - 'time_horizon': item['pmd:timeHorizon'], - 'scenario_timestamp': item['pmd:scenarioDate'], - 'pmd:fullModel_ID': item['pmd:fullModel_ID'], - 'pmd:version': item['pmd:version'], + '@time_horizon': item['pmd:timeHorizon'], + '@scenario_timestamp': item['pmd:scenarioDate'], + 'fullModel_ID': item['pmd:fullModel_ID'], + '@version': item['pmd:version'], 'qualityIndicator': 'valid', - 'pmd:creationDate': item['pmd:creationDate'], + '@timestamp': item['pmd:creationDate'], } for item 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]) @@ -229,13 +231,13 @@ def handle(self, task_object: dict, properties: dict, **kwargs): 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'], - 'pmd:fullModel_ID': model['pmd:fullModel_ID'], - 'pmd:version': model['pmd:version'], - 'data_source' : 'PDN', + '@time_horizon': model['pmd:timeHorizon'], + '@scenario_timestamp': model['pmd:scenarioDate'], + 'fullModel_ID': model['pmd:fullModel_ID'], + '@version': model['pmd:version'], + '@data_source' : 'PDN', 'qualityIndicator': 'replaced', - 'pmd:creationDate': model['pmd:creationDate'] + '@timestamp': model['pmd:creationDate'] } for model in replacement_models_local] merged_model.replaced_entity.extend(replaced_entities_local) additional_models.extend(replacement_models_local) @@ -267,13 +269,13 @@ def handle(self, task_object: dict, properties: dict, **kwargs): 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'], - 'pmd:fullModel_ID': model['pmd:fullModel_ID'], - 'pmd:version': model['pmd:version'], - 'data_source' : 'OPDM', + '@time_horizon': model['pmd:timeHorizon'], + '@scenario_timestamp': model['pmd:scenarioDate'], + 'fullModel_ID': model['pmd:fullModel_ID'], + '@version': model['pmd:version'], + '@data_source' : 'OPDM', 'qualityIndicator': 'replaced', - 'pmd:creationDate': model['pmd:creationDate'] } for model in replacement_models] + '@timestamp': model['pmd:creationDate'] } for model in replacement_models] merged_model.replaced_entity.extend(replaced_entities) models.extend(replacement_models) merged_model.replaced = True @@ -395,6 +397,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) @@ -477,11 +481,18 @@ def handle(self, task_object: dict, properties: dict, **kwargs): except Exception as error: logger.error(f"Failed to create merge report: {error}") - #send lvl 8 report - if lvl8_reporting: + #send lvl 8 report + if lvl8_reporting: + try: lvl8_report = merge_functions.lvl8_report_cgm(merge_report) - logger.debug(f"lvl8 report generated: {lvl8_report}") - #TODO Sending the report via EDX. + service_edx = edx.EDX()#.create_client(server=EDX_SERVER, username=EDX_USERNAME, password=EDX_PASSWORD) + message_id = service_edx.send_message(receiver_EIC=QAS_EIC, + business_type=QAS_MSG_TYPE, + content=lvl8_report) + logger.info(f"lvl8 report generated and sent-ID: {message_id}")#: {lvl8_report}") + except Exception as error: + logger.error(f"Failed to send lvl8 report: {error}") + # Append message headers with OPDM root metadata extracted_meta = {key: value for key, value in opdm_object_meta.items() if isinstance(value, str)} From 94b0925fe353815b186b77af326d12f148d8a25d Mon Sep 17 00:00:00 2001 From: VeikoAunapuu Date: Wed, 30 Jul 2025 13:54:28 +0300 Subject: [PATCH 06/11] Added missing dependency --- emf/model_merger/model_merger.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/emf/model_merger/model_merger.py b/emf/model_merger/model_merger.py index 3bc8cab8..9f0b2e04 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 From 23ddeb3596419dca9038ea1307df7524c04930d9 Mon Sep 17 00:00:00 2001 From: VeikoAunapuu Date: Wed, 30 Jul 2025 14:58:40 +0300 Subject: [PATCH 07/11] fixes, and tested in QAS --- emf/model_merger/merge_functions.py | 2 +- emf/model_merger/model_merger.py | 8 ++++---- 2 files changed, 5 insertions(+), 5 deletions(-) diff --git a/emf/model_merger/merge_functions.py b/emf/model_merger/merge_functions.py index a6af7f7d..65f085b7 100644 --- a/emf/model_merger/merge_functions.py +++ b/emf/model_merger/merge_functions.py @@ -864,7 +864,7 @@ def lvl8_report_cgm(merge_report): 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['network_meta']['fullModel_ID'], + 'serviceProvider': merge_report["merge_entity"], 'xmlns': "http://entsoe.eu/checks" } qa_root = ET.Element("QAReport", attrib=qa_attribs) diff --git a/emf/model_merger/model_merger.py b/emf/model_merger/model_merger.py index 9f0b2e04..e9591609 100644 --- a/emf/model_merger/model_merger.py +++ b/emf/model_merger/model_merger.py @@ -188,7 +188,7 @@ def handle(self, task_object: dict, properties: dict, **kwargs): '@scenario_timestamp': item['pmd:scenarioDate'], 'fullModel_ID': item['pmd:fullModel_ID'], '@version': item['pmd:version'], - 'qualityIndicator': 'valid', + 'qualityIndicator': 'Valid', '@timestamp': item['pmd:creationDate'], } for item in models @@ -211,7 +211,7 @@ def handle(self, task_object: dict, properties: dict, **kwargs): '@scenario_timestamp': item['pmd:scenarioDate'], 'fullModel_ID': item['pmd:fullModel_ID'], '@version': item['pmd:version'], - 'qualityIndicator': 'valid', + 'qualityIndicator': 'Valid', '@timestamp': item['pmd:creationDate'], } for item in additional_models @@ -236,7 +236,7 @@ def handle(self, task_object: dict, properties: dict, **kwargs): 'fullModel_ID': model['pmd:fullModel_ID'], '@version': model['pmd:version'], '@data_source' : 'PDN', - 'qualityIndicator': 'replaced', + 'qualityIndicator': 'Substituted', '@timestamp': model['pmd:creationDate'] } for model in replacement_models_local] merged_model.replaced_entity.extend(replaced_entities_local) @@ -274,7 +274,7 @@ def handle(self, task_object: dict, properties: dict, **kwargs): 'fullModel_ID': model['pmd:fullModel_ID'], '@version': model['pmd:version'], '@data_source' : 'OPDM', - 'qualityIndicator': 'replaced', + 'qualityIndicator': 'Substituted', '@timestamp': model['pmd:creationDate'] } for model in replacement_models] merged_model.replaced_entity.extend(replaced_entities) models.extend(replacement_models) From 13098d96745cc9f70ce69f1dea9bdffe145e2f12 Mon Sep 17 00:00:00 2001 From: VeikoAunapuu Date: Wed, 30 Jul 2025 15:18:27 +0300 Subject: [PATCH 08/11] Added filename to models for better trace --- emf/model_merger/model_merger.py | 8 ++++++-- 1 file changed, 6 insertions(+), 2 deletions(-) diff --git a/emf/model_merger/model_merger.py b/emf/model_merger/model_merger.py index e9591609..c6e0d4bd 100644 --- a/emf/model_merger/model_merger.py +++ b/emf/model_merger/model_merger.py @@ -190,6 +190,7 @@ def handle(self, task_object: dict, properties: dict, **kwargs): '@version': item['pmd:version'], 'qualityIndicator': 'Valid', '@timestamp': item['pmd:creationDate'], + 'filename': item['pmd:fileName'] } for item in models ] @@ -213,6 +214,7 @@ def handle(self, task_object: dict, properties: dict, **kwargs): '@version': item['pmd:version'], 'qualityIndicator': 'Valid', '@timestamp': item['pmd:creationDate'], + 'filename': item['pmd:fileName'] } for item in additional_models ] @@ -237,7 +239,8 @@ def handle(self, task_object: dict, properties: dict, **kwargs): '@version': model['pmd:version'], '@data_source' : 'PDN', 'qualityIndicator': 'Substituted', - '@timestamp': model['pmd:creationDate'] + '@timestamp': model['pmd:creationDate'], + 'filename': model['pmd:fileName'] } for model in replacement_models_local] merged_model.replaced_entity.extend(replaced_entities_local) additional_models.extend(replacement_models_local) @@ -275,7 +278,8 @@ def handle(self, task_object: dict, properties: dict, **kwargs): '@version': model['pmd:version'], '@data_source' : 'OPDM', 'qualityIndicator': 'Substituted', - '@timestamp': model['pmd:creationDate'] } for model in replacement_models] + '@timestamp': model['pmd:creationDate'], + 'filename': model['pmd:fileName']} for model in replacement_models] merged_model.replaced_entity.extend(replaced_entities) models.extend(replacement_models) merged_model.replaced = True From c76c65bf0b1adffb0eeed03b7b779336266c900c Mon Sep 17 00:00:00 2001 From: "martynas.karobcikas" Date: Thu, 31 Jul 2025 11:53:04 +0300 Subject: [PATCH 09/11] temporary commit --- emf/model_merger/merge_functions.py | 7 ++-- emf/model_merger/model_merger.py | 51 ++++++++++++++--------------- 2 files changed, 29 insertions(+), 29 deletions(-) diff --git a/emf/model_merger/merge_functions.py b/emf/model_merger/merge_functions.py index 65f085b7..706abc25 100644 --- a/emf/model_merger/merge_functions.py +++ b/emf/model_merger/merge_functions.py @@ -858,7 +858,7 @@ def run_post_merge_processing(input_models: list, return sv_data, ssh_data, opdm_object_meta -def lvl8_report_cgm(merge_report): +def lvl8_report_cgm(merge_report: dict): # Create root qa_attribs = { @@ -932,8 +932,8 @@ def lvl8_report_cgm(merge_report): 'processType': i['@time_horizon'], 'qualityIndicator': i['qualityIndicator'], }) - resource_igm= ET.SubElement(igm, "resource") - resource_igm.text=i['fullModel_ID'] + resource_igm = ET.SubElement(igm, "resource") + resource_igm.text = i['fullModel_ID'] # Add EMFInformation @@ -947,6 +947,7 @@ def lvl8_report_cgm(merge_report): 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 c6e0d4bd..2b817eeb 100644 --- a/emf/model_merger/model_merger.py +++ b/emf/model_merger/model_merger.py @@ -181,21 +181,20 @@ def handle(self, task_object: dict, properties: dict, **kwargs): excluded_models=excluded_models, filter_on='pmd:TSO') - merged_model.included_opdm=[ + merged_model.merge_included_opdm = [ { 'tso': item['pmd:TSO'], - '@time_horizon': item['pmd:timeHorizon'], - '@scenario_timestamp': item['pmd:scenarioDate'], - 'fullModel_ID': item['pmd:fullModel_ID'], - '@version': item['pmd:version'], - 'qualityIndicator': 'Valid', - '@timestamp': item['pmd:creationDate'], - 'filename': item['pmd:fileName'] + 'time_horizon': item['pmd:timeHorizon'], + 'scenario_timestamp': item['pmd:scenarioDate'], + 'model_sv_id': item['pmd:fullModel_ID'], + 'version': item['pmd:version'], + 'quality_indicator': 'Valid', + 'creation_timestamp': item['pmd:creationDate'], + 'file_name': item['pmd:fileName'] } - for item in models + for item 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, @@ -205,16 +204,16 @@ 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.included_local = [ + merged_model.merge_included_local = [ { 'tso': item['pmd:TSO'], - '@time_horizon': item['pmd:timeHorizon'], - '@scenario_timestamp': item['pmd:scenarioDate'], - 'fullModel_ID': item['pmd:fullModel_ID'], - '@version': item['pmd:version'], - 'qualityIndicator': 'Valid', - '@timestamp': item['pmd:creationDate'], - 'filename': item['pmd:fileName'] + 'time_horizon': item['pmd:timeHorizon'], + 'scenario_timestamp': item['pmd:scenarioDate'], + 'model_sv_id': item['pmd:fullModel_ID'], + 'version': item['pmd:version'], + 'quality_indicator': 'Valid', + 'creation_timestamp': item['pmd:creationDate'], + 'file_name': item['pmd:fileName'] } for item in additional_models ] @@ -233,7 +232,7 @@ def handle(self, task_object: dict, properties: dict, **kwargs): 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'], + 'time_horizon': model['pmd:timeHorizon'], '@scenario_timestamp': model['pmd:scenarioDate'], 'fullModel_ID': model['pmd:fullModel_ID'], '@version': model['pmd:version'], @@ -485,17 +484,17 @@ def handle(self, task_object: dict, properties: dict, **kwargs): except Exception as error: logger.error(f"Failed to create merge report: {error}") - #send lvl 8 report + # Send QAS level 8 report if configured if lvl8_reporting: try: - lvl8_report = merge_functions.lvl8_report_cgm(merge_report) - service_edx = edx.EDX()#.create_client(server=EDX_SERVER, username=EDX_USERNAME, password=EDX_PASSWORD) + 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"lvl8 report generated and sent-ID: {message_id}")#: {lvl8_report}") + 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 lvl8 report: {error}") + logger.error(f"Failed to send QAS-Level-8 report with error: {error}") # Append message headers with OPDM root metadata From a693edcff11d5aa90a7831b0a8e9c56399a74036 Mon Sep 17 00:00:00 2001 From: "martynas.karobcikas" Date: Thu, 31 Jul 2025 12:45:07 +0300 Subject: [PATCH 10/11] reafactor: changed to use dataclasss for model entity --- emf/model_merger/merge_functions.py | 28 +++---- emf/model_merger/model_merger.py | 123 +++++++++++++--------------- 2 files changed, 71 insertions(+), 80 deletions(-) diff --git a/emf/model_merger/merge_functions.py b/emf/model_merger/merge_functions.py index 706abc25..d78f1977 100644 --- a/emf/model_merger/merge_functions.py +++ b/emf/model_merger/merge_functions.py @@ -858,11 +858,12 @@ 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'), + '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" @@ -884,8 +885,8 @@ def lvl8_report_cgm(merge_report: dict): '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=[] + # 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") @@ -899,8 +900,8 @@ def lvl8_report_cgm(merge_report: dict): # 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 + '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"], @@ -908,7 +909,6 @@ def lvl8_report_cgm(merge_report: dict): } cgm = ET.SubElement(qa_root, "CGM", attrib=cgm_attribs) - try: for v in violations: rv = ET.SubElement(cgm, "RuleViolation", { @@ -921,20 +921,18 @@ def lvl8_report_cgm(merge_report: dict): 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['included_opdm'] + merge_report['included_local'] + merge_report['replaced_entity']: + for i in merge_report['merge_included_entity'] + merge_report['replaced_entity']: igm = ET.SubElement(cgm, "IGM", { - 'created': i["@timestamp"], - 'scenarioTime': datetime.datetime.fromisoformat(i['@scenario_timestamp']).strftime('%Y-%m-%dT%H:%M:%SZ'), + 'created': i["creation_timestamp"], + 'scenarioTime': datetime.datetime.fromisoformat(i['scenario_timestamp']).strftime('%Y-%m-%dT%H:%M:%SZ'), 'tso': i['tso'], - 'version': i['@version'], - 'processType': i['@time_horizon'], - 'qualityIndicator': i['qualityIndicator'], + 'version': i['version'], + 'processType': i['time_horizon'], + 'qualityIndicator': i['quality_indicator'], }) resource_igm = ET.SubElement(igm, "resource") - resource_igm.text = i['fullModel_ID'] - + resource_igm.text = i['model_sv_id'] # Add EMFInformation ET.SubElement(cgm, "EMFInformation", { diff --git a/emf/model_merger/model_merger.py b/emf/model_merger/model_merger.py index 2b817eeb..b8c09d4d 100644 --- a/emf/model_merger/model_merger.py +++ b/emf/model_merger/model_merger.py @@ -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,8 +67,31 @@ class MergedModel: replacement_reason: List = field(default_factory=list) outages_updated: List = field(default_factory=list) outages_unmapped: List = field(default_factory=list) - included_opdm: List = field(default_factory=list) - included_local: List = field(default_factory=list) + merge_included_entity: List = field(default_factory=list) + + +@dataclass(init=False) +class ModelEntity: + data_source: str = "OPDM" + 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, **kwargs): + self.data_source = data_source + 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: @@ -83,7 +105,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] @@ -181,19 +204,7 @@ def handle(self, task_object: dict, properties: dict, **kwargs): excluded_models=excluded_models, filter_on='pmd:TSO') - merged_model.merge_included_opdm = [ - { - 'tso': item['pmd:TSO'], - 'time_horizon': item['pmd:timeHorizon'], - 'scenario_timestamp': item['pmd:scenarioDate'], - 'model_sv_id': item['pmd:fullModel_ID'], - 'version': item['pmd:version'], - 'quality_indicator': 'Valid', - 'creation_timestamp': item['pmd:creationDate'], - 'file_name': item['pmd:fileName'] - } - for item in models - ] + merged_model.merge_included_entity = [ModelEntity(data_source='OPDM', **model).__dict__ for model in models] # Get additional models from ObjectStorage if local import is configured if local_import_models: @@ -204,21 +215,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_local = [ - { - 'tso': item['pmd:TSO'], - 'time_horizon': item['pmd:timeHorizon'], - 'scenario_timestamp': item['pmd:scenarioDate'], - 'model_sv_id': item['pmd:fullModel_ID'], - 'version': item['pmd:version'], - 'quality_indicator': 'Valid', - 'creation_timestamp': item['pmd:creationDate'], - 'file_name': item['pmd:fileName'] - } - for item 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.merge_included_entity.extend( + [ModelEntity(data_source='PDN', **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]] merged_model.excluded.extend([{'tso': tso, 'reason': 'missing-pdn'} for tso in missing_local_import]) # Perform local replacement if configured @@ -230,17 +231,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'], - 'fullModel_ID': model['pmd:fullModel_ID'], - '@version': model['pmd:version'], - '@data_source' : 'PDN', - 'qualityIndicator': 'Substituted', - '@timestamp': model['pmd:creationDate'], - 'filename': model['pmd:fileName'] - } 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', **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: @@ -269,16 +263,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'], - 'fullModel_ID': model['pmd:fullModel_ID'], - '@version': model['pmd:version'], - '@data_source' : 'OPDM', - 'qualityIndicator': 'Substituted', - '@timestamp': model['pmd:creationDate'], - 'filename': model['pmd:fileName']} 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', **model).__dict__ for model in + replacement_models] merged_model.replaced_entity.extend(replaced_entities) models.extend(replacement_models) merged_model.replaced = True @@ -308,7 +296,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, @@ -321,7 +310,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) @@ -329,7 +319,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: @@ -351,7 +342,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 @@ -400,7 +392,7 @@ def handle(self, task_object: dict, properties: dict, **kwargs): task_properties=task_properties ) - #for merge report need to get the final uuid. + # 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]) @@ -472,7 +464,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: @@ -483,9 +476,9 @@ def handle(self, task_object: dict, properties: dict, **kwargs): logger.error(f"Merge report sending to Elastic failed: {error}") except Exception as error: logger.error(f"Failed to create merge report: {error}") - + # Send QAS level 8 report if configured - if lvl8_reporting: + if lvl8_reporting and merge_report: try: lvl8_report = merge_functions.lvl8_report_cgm(merge_report=merge_report) service_edx = edx.EDX() @@ -495,8 +488,9 @@ def handle(self, task_object: dict, properties: dict, **kwargs): 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 @@ -510,7 +504,6 @@ def handle(self, task_object: dict, properties: dict, **kwargs): if __name__ == "__main__": - sample_task = { "@context": "https://example.com/task_context.jsonld", "@type": "Task", @@ -539,7 +532,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"], @@ -556,7 +549,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" + "lvl8_reporting": "True" } } From c7b5593abe843400d227a4278d1f7a9d52475023 Mon Sep 17 00:00:00 2001 From: VeikoAunapuu Date: Thu, 31 Jul 2025 13:46:32 +0300 Subject: [PATCH 11/11] minor fixes --- emf/model_merger/merge_functions.py | 2 +- emf/model_merger/model_merger.py | 12 +++++++----- 2 files changed, 8 insertions(+), 6 deletions(-) diff --git a/emf/model_merger/merge_functions.py b/emf/model_merger/merge_functions.py index d78f1977..772b6285 100644 --- a/emf/model_merger/merge_functions.py +++ b/emf/model_merger/merge_functions.py @@ -927,7 +927,7 @@ def lvl8_report_cgm(merge_report: dict): 'created': i["creation_timestamp"], 'scenarioTime': datetime.datetime.fromisoformat(i['scenario_timestamp']).strftime('%Y-%m-%dT%H:%M:%SZ'), 'tso': i['tso'], - 'version': i['version'], + 'version': str(i['version']), 'processType': i['time_horizon'], 'qualityIndicator': i['quality_indicator'], }) diff --git a/emf/model_merger/model_merger.py b/emf/model_merger/model_merger.py index b8c09d4d..bfe9a77d 100644 --- a/emf/model_merger/model_merger.py +++ b/emf/model_merger/model_merger.py @@ -73,6 +73,7 @@ class MergedModel: @dataclass(init=False) class ModelEntity: data_source: str = "OPDM" + quality_indicator: str = "Valid" tso: str = None time_horizon: str = None scenario_timestamp: str = None @@ -82,8 +83,9 @@ class ModelEntity: creation_timestamp: str = None file_name: str = None - def __init__(self, data_source: str, **kwargs): + 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') @@ -204,7 +206,7 @@ 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', **model).__dict__ for model in models] + 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: @@ -216,7 +218,7 @@ def handle(self, task_object: dict, properties: dict, **kwargs): included_models=local_import_models, filter_on='pmd:TSO') merged_model.merge_included_entity.extend( - [ModelEntity(data_source='PDN', **model).__dict__ for model in additional_models]) + [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]] @@ -233,7 +235,7 @@ def handle(self, task_object: dict, properties: dict, **kwargs): 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', **model).__dict__ for model in + 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) @@ -265,7 +267,7 @@ def handle(self, task_object: dict, properties: dict, **kwargs): if replacement_models: logger.info( f"Replacement model(s) found: {[model['pmd:fileName'] for model in replacement_models]}") - replaced_entities = [ModelEntity(data_source='OPDM', **model).__dict__ for model in + 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)