From c1f3f582e9bcb819dc79ad9b3219372da5302cdd Mon Sep 17 00:00:00 2001 From: Saksham Date: Thu, 13 Aug 2026 16:37:38 +0200 Subject: [PATCH] fix(rdm.load): Run CLC Sync manually after UnitOfWork commit --- cds_migrator_kit/rdm/migration_config.py | 3 +- cds_migrator_kit/rdm/records/load/load.py | 46 +++++++++++++++-------- 2 files changed, 33 insertions(+), 16 deletions(-) diff --git a/cds_migrator_kit/rdm/migration_config.py b/cds_migrator_kit/rdm/migration_config.py index 51b226f7..68ec9ef6 100644 --- a/cds_migrator_kit/rdm/migration_config.py +++ b/cds_migrator_kit/rdm/migration_config.py @@ -504,7 +504,8 @@ def resolve_record_pid(pid): # SubjectsValidationComponent, *DefaultRecordsComponents, CDSResourcePublication, - ClcSyncComponent, + # Component disabled to not trigger due to auto=True and the sync task is triggered manually in migration code + # ClcSyncComponent, # component disabled, this part is handled separately in migration code # due to two conflicting DB updates causing StaleDataError # MintAlternateIdentifierComponent, diff --git a/cds_migrator_kit/rdm/records/load/load.py b/cds_migrator_kit/rdm/records/load/load.py index 35ef5319..35ec4d8b 100644 --- a/cds_migrator_kit/rdm/records/load/load.py +++ b/cds_migrator_kit/rdm/records/load/load.py @@ -6,6 +6,7 @@ # the terms of the MIT License; see LICENSE file for more details. """CDS-RDM migration load module.""" + import datetime import json import os @@ -14,6 +15,7 @@ import arrow from cds_rdm.clc_sync.models import CDSToCLCSyncModel +from cds_rdm.clc_sync.proxies import current_clc_sync_service from cds_rdm.legacy.models import CDSMigrationLegacyRecord from cds_rdm.legacy.resolver import get_pid_by_legacy_recid from cds_rdm.minters import legacy_recid_minter @@ -194,6 +196,22 @@ def _load_record_access(self, draft, access_dict): record.access = access_dict["access_obj"] record.commit() + def _after_commit_run_clc_sync(self, record_state): + """Run the CLC sync after UOW commit.""" + if not self._is_final_record: + return + if self.clc_sync: + clc_sync_entry = current_clc_sync_service.read( + system_identity, record_state["parent_recid"] + ).to_dict() + clc_sync_entry["record"] = current_rdm_records_service.read( + system_identity, record_state["latest_version"] + ).to_dict() + clc_sync_entry["auto_sync"] = True + current_clc_sync_service.update( + system_identity, clc_sync_entry["id"], clc_sync_entry + ) + def _after_publish_update_dois(self, identity, record, entry, uow): """Update migrated DOIs post publish.""" if not self._is_final_record: @@ -257,9 +275,7 @@ def _normalize_group_name(subject): elif specific_file_restrictions == "restricted": # https://cds.cern.ch/admin/webaccess/webaccessadmin.py/showroledetails?id_role=69 groups.add("cern-personnel") - elif specific_file_restrictions.strip().endswith( - "[CERN]" - ) and not any( + elif specific_file_restrictions.strip().endswith("[CERN]") and not any( kw in specific_file_restrictions for kw in ("firerole:", "allow ") ): # bare CERN e-group name, e.g. @@ -525,12 +541,14 @@ def _after_publish(self, identity, published_record, entry, version, uow): request_data = entry["record"].get("_request_data", {}) if request_data and not self.create_inclusion_request: - raise ManualImportRequired(message="Detected request data, enable the requests", - field="validation", - stage="load", - recid=entry["record"]["recid"], - priority="warning", - subfield=None,) + raise ManualImportRequired( + message="Detected request data, enable the requests", + field="validation", + stage="load", + recid=entry["record"]["recid"], + priority="warning", + subfield=None, + ) if self.create_inclusion_request and request_data: self._after_publish_add_inclusion_request( request_data, published_record, entry, uow @@ -819,15 +837,11 @@ def _load(self, entry, uow=None): elif uow is not None: recid_state_after_load = self._load_versions(entry, uow) if recid_state_after_load: - self._save_original_dumped_record( - entry, recid_state_after_load - ) + self._save_original_dumped_record(entry, recid_state_after_load) self._after_load_clc_sync(recid_state_after_load) else: with UnitOfWork(db.session) as inner_uow: - recid_state_after_load = self._load_versions( - entry, inner_uow - ) + recid_state_after_load = self._load_versions(entry, inner_uow) if recid_state_after_load: self._save_original_dumped_record( entry, recid_state_after_load @@ -839,6 +853,8 @@ def _load(self, entry, uow=None): # commit boundary and is responsible for finalising the # record only after it actually commits. self.migration_logger.finalise_record(recid) + # Run the CLC sync after UOW commit + self._after_commit_run_clc_sync(recid_state_after_load) return recid_state_after_load except (UnexpectedValue, ManualImportRequired) as e: self.migration_logger.add_log(e, record=entry)