From 224597cba89710e5e9041bfc185ff0dc0a903445 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Z=C3=BCbeyde=20Civelek?= Date: Tue, 21 Jul 2026 12:08:01 +0200 Subject: [PATCH 1/3] change(inclusion request): add the user during transform # Conflicts: # cds_migrator_kit/rdm/records/load/load.py --- cds_migrator_kit/rdm/records/load/load.py | 73 +------------------ .../rdm/records/transform/transform.py | 7 ++ .../xml_processing/quality/__init__.py | 8 ++ .../xml_processing/quality/reviewers.py | 68 +++++++++++++++++ .../xml_processing/rules/faser_publication.py | 11 +++ .../xml_processing/rules/research.py | 15 +++- 6 files changed, 108 insertions(+), 74 deletions(-) create mode 100644 cds_migrator_kit/rdm/records/transform/xml_processing/quality/__init__.py create mode 100644 cds_migrator_kit/rdm/records/transform/xml_processing/quality/reviewers.py diff --git a/cds_migrator_kit/rdm/records/load/load.py b/cds_migrator_kit/rdm/records/load/load.py index fdb25568..5cac99eb 100644 --- a/cds_migrator_kit/rdm/records/load/load.py +++ b/cds_migrator_kit/rdm/records/load/load.py @@ -48,61 +48,6 @@ ) -def _is_email(value): - """Return True if the reviewer value looks like an email address.""" - return "@" in value - - -def _parse_reviewer_name(name): - """Split a 'Family, Given' or 'Given Family' string into (family, given). - - ``request_reviewers`` (906__p) stores names as "Given Family" (comma - already resolved), but legacy data can also arrive as "Family, Given". - """ - name = name.strip() - if "," in name: - family, _, given = name.partition(",") - return family.strip(), given.strip() - parts = name.split() - if len(parts) > 1: - return parts[-1], " ".join(parts[:-1]) - return name, "" - - -def find_reviewer(reviewer): - """Resolve a reviewer string (email or name) to a User. - - :param reviewer: email address, or a "Family, Given"/"Given Family" name. - :raises RecordFlaggedCuration: if no matching user is found, so the - record is flagged for manual curation instead of failing outright. - """ - reviewer = reviewer.strip() - if _is_email(reviewer): - user = User.query.filter_by(email=reviewer).one_or_none() - else: - family_name, given_name = _parse_reviewer_name(reviewer) - query = User.query.filter( - db.func.lower(User._user_profile["family_name"].as_string()) - == family_name.lower() - ) - if given_name: - query = query.filter( - db.func.lower(User._user_profile["given_name"].as_string()) - == given_name.lower() - ) - user = query.one_or_none() - - if user is None: - raise RecordFlaggedCuration( - message=f"Reviewer '{reviewer}' could not be matched to an account.", - field="request_reviewers", - stage="load", - value=reviewer, - ) - - return user - - def import_legacy_files(filepath): """Download file from legacy.""" if current_app.config["CDS_MIGRATOR_KIT_ENV"] == "local": @@ -121,7 +66,6 @@ def __init__( self, db_uri=None, data_dir=None, - tmp_dir=None, entries=None, dry_run=False, legacy_pids_to_redirect=None, @@ -484,19 +428,7 @@ def _after_publish_add_inclusion_request(self, request_data, record, entry, uow) return status = request_data.get("status", "accepted") - reviewer_names = request_data.get("reviewers", []) - reviewers = [] - for name in reviewer_names: - try: - reviewers.append(find_reviewer(name)) - except RecordFlaggedCuration as exc: - self.migration_logger.add_information( - legacy_recid, - {"message": exc.message, "value": exc.value}, - ) - # keep a placeholder so the reviewer count/order is preserved; - # resolved to the "-1" sentinel user id below - reviewers.append(None) + reviewers = request_data.get("reviewers", []) created_at = datetime.datetime.fromisoformat(record["created"]) @@ -537,8 +469,7 @@ def create_event(request_model, payload, event_type, user): request.model.created = created_at if reviewers: - reviewers_payload = [{"user": str(r.id) if r else "-1"} for r in reviewers] - request.reviewers = reviewers_payload + request.reviewers = reviewers if status: request.status = status diff --git a/cds_migrator_kit/rdm/records/transform/transform.py b/cds_migrator_kit/rdm/records/transform/transform.py index 210c4752..d4752072 100644 --- a/cds_migrator_kit/rdm/records/transform/transform.py +++ b/cds_migrator_kit/rdm/records/transform/transform.py @@ -770,6 +770,13 @@ def transform(self, entry): del json_data["_clc_sync"] request_data = json_data.pop("request_data", None) + if request_data: + reviewer_errors = request_data.pop("_reviewer_errors", []) + for error in reviewer_errors: + self.migration_logger.add_information( + entry["recid"], + error, + ) record_json_output = { "files": self._files(record_dump), diff --git a/cds_migrator_kit/rdm/records/transform/xml_processing/quality/__init__.py b/cds_migrator_kit/rdm/records/transform/xml_processing/quality/__init__.py new file mode 100644 index 00000000..5b2e0c6d --- /dev/null +++ b/cds_migrator_kit/rdm/records/transform/xml_processing/quality/__init__.py @@ -0,0 +1,8 @@ +# -*- coding: utf-8 -*- +# +# Copyright (C) 2022 CERN. +# +# CDS-RDM is free software; you can redistribute it and/or modify it under +# the terms of the MIT License; see LICENSE file for more details. + +"""RDM records transform quality utilities.""" diff --git a/cds_migrator_kit/rdm/records/transform/xml_processing/quality/reviewers.py b/cds_migrator_kit/rdm/records/transform/xml_processing/quality/reviewers.py new file mode 100644 index 00000000..c05447b6 --- /dev/null +++ b/cds_migrator_kit/rdm/records/transform/xml_processing/quality/reviewers.py @@ -0,0 +1,68 @@ +# -*- coding: utf-8 -*- +# +# Copyright (C) 2022 CERN. +# +# CDS-RDM is free software; you can redistribute it and/or modify it under +# the terms of the MIT License; see LICENSE file for more details. + +"""Reviewer resolution utilities.""" + +from invenio_accounts.models import User +from invenio_db import db + +from cds_migrator_kit.errors import RecordFlaggedCuration + + +def _is_email(value): + """Return True if the reviewer value looks like an email address.""" + return "@" in value + + +def _parse_reviewer_name(name): + """Split a 'Family, Given' or 'Given Family' string into (family, given). + + ``request_reviewers`` (906__p) stores names as "Given Family" (comma + already resolved), but legacy data can also arrive as "Family, Given". + """ + name = name.strip() + if "," in name: + family, _, given = name.partition(",") + return family.strip(), given.strip() + parts = name.split() + if len(parts) > 1: + return parts[-1], " ".join(parts[:-1]) + return name, "" + + +def find_reviewer(reviewer): + """Resolve a reviewer string (email or name) to a User. + + :param reviewer: email address, or a "Family, Given"/"Given Family" name. + :raises RecordFlaggedCuration: if no matching user is found, so the + record is flagged for manual curation instead of failing outright. + """ + reviewer = reviewer.strip() + if _is_email(reviewer): + user = User.query.filter_by(email=reviewer).one_or_none() + else: + family_name, given_name = _parse_reviewer_name(reviewer) + query = User.query.filter( + db.func.lower(User._user_profile["family_name"].as_string()) + == family_name.lower() + ) + if given_name: + query = query.filter( + db.func.lower(User._user_profile["given_name"].as_string()) + == given_name.lower() + ) + user = query.one_or_none() + + if user is None: + raise RecordFlaggedCuration( + message=f"Reviewer '{reviewer}' could not be matched to an account.", + field="request_reviewers", + stage="transform", + value=reviewer, + ) + + return user diff --git a/cds_migrator_kit/rdm/records/transform/xml_processing/rules/faser_publication.py b/cds_migrator_kit/rdm/records/transform/xml_processing/rules/faser_publication.py index 59995a4a..93131091 100644 --- a/cds_migrator_kit/rdm/records/transform/xml_processing/rules/faser_publication.py +++ b/cds_migrator_kit/rdm/records/transform/xml_processing/rules/faser_publication.py @@ -16,6 +16,7 @@ from cds_migrator_kit.transform.xml_processing.rules.base import process_contributors from ...models.faser_publication import faser_publication_model as model +from .research import status as research_status @model.over("collection", "^690C_", override=True) @@ -54,6 +55,16 @@ def access_grants(self, key, value): raise IgnoreKey("access_grants") +@model.over("_approval", "(^591__)", override=True) +def status(self, key, value): + """Translates faser publication approval status.""" + research_status(self, key, value) + request_data = self.setdefault("request_data", {}) + reviewers = request_data.setdefault("reviewers", []) + reviewers.append({"group": "faser-all"}) + raise IgnoreKey("_approval") + + @model.over("faser_contributors", "^700__", override=True) @for_each_value @require(["a"]) diff --git a/cds_migrator_kit/rdm/records/transform/xml_processing/rules/research.py b/cds_migrator_kit/rdm/records/transform/xml_processing/rules/research.py index b6df2ef0..5444ccd0 100644 --- a/cds_migrator_kit/rdm/records/transform/xml_processing/rules/research.py +++ b/cds_migrator_kit/rdm/records/transform/xml_processing/rules/research.py @@ -8,7 +8,7 @@ from idutils.normalizers import normalize_isbn, normalize_issn from isbnlib import NotValidISBNError -from cds_migrator_kit.errors import ManualImportRequired, UnexpectedValue +from cds_migrator_kit.errors import ManualImportRequired, RecordFlaggedCuration, UnexpectedValue from cds_migrator_kit.transform.xml_processing.quality.decorators import ( filter_list_values, for_each_value, @@ -23,6 +23,7 @@ from ...models.base_publication_record import rdm_base_publication_model as model from .base import licenses as _base_licenses from .base import normalize +from ..quality.reviewers import find_reviewer # Unwrapped base functions (strip @for_each_value to avoid double-wrapping). # licenses also has @filter_values beneath @for_each_value, so two levels deep. @@ -476,8 +477,16 @@ def request_reviewers(self, key, value): request_data = self.setdefault("request_data", {}) reviewers = request_data.setdefault("reviewers", []) - if reviewer not in reviewers: - reviewers.append(reviewer) + try: + user = find_reviewer(reviewer) + reviewer_entry = {"user": str(user.id)} + except RecordFlaggedCuration as exc: + reviewer_errors = request_data.setdefault("_reviewer_errors", []) + reviewer_errors.append({"message": exc.message, "value": exc.value}) + reviewer_entry = {"user": "-1"} + + if reviewer_entry not in reviewers: + reviewers.append(reviewer_entry) raise IgnoreKey("request_reviewers") From e5c990947d4084065dba1b5fd0d774351512cc7c Mon Sep 17 00:00:00 2001 From: Karolina Przerwa Date: Wed, 22 Jul 2026 11:46:06 +0200 Subject: [PATCH 2/3] fix(ep migration): change the file split condition add(ep migration): review step fix(ep migration): usage of uow --- cds_migrator_kit/rdm/migration_config.py | 44 +++- .../rdm/records/load/approval_request.py | 143 +++++++++---- .../rdm/records/load/ep_approval_entry.py | 67 ++++-- .../rdm/records/load/ep_approval_load.py | 193 ++++++++++++------ cds_migrator_kit/rdm/records/load/load.py | 37 +++- .../xml_processing/rules/research.py | 4 +- cds_migrator_kit/rdm/streams.yaml | 66 +++++- scripts/opensearch-init.sh | 24 ++- scripts/snapshot.sh | 90 ++++++++ 9 files changed, 535 insertions(+), 133 deletions(-) create mode 100755 scripts/snapshot.sh diff --git a/cds_migrator_kit/rdm/migration_config.py b/cds_migrator_kit/rdm/migration_config.py index 459e9b52..75068681 100644 --- a/cds_migrator_kit/rdm/migration_config.py +++ b/cds_migrator_kit/rdm/migration_config.py @@ -491,6 +491,7 @@ def resolve_record_pid(pid): # Access group already added to the record as 506__m when draft created, no need to add: # https://gitlab.cern.ch/cds-team/cds-legacy/-/blob/master/src/wn-cdsweb/lib/python/invenio/websubmit_functions/EPPHAPP_Test_values.py#L249-266 "EP Restricted Draft": [], + "PH-EP Restricted Draft": [], # CERN E-guide restricted docs: https://cds.cern.ch/admin/webaccess/webaccessadmin.py/showroledetails?id_role=69 CERN personnel has view rights } @@ -539,10 +540,51 @@ def resolve_record_pid(pid): CDS_COMMITTEE_APPROVAL_COMMUNITIES = { "dd13404c-bcd6-4b15-aeef-38d678c61ff1": { + # faser "label": "EP approval", # shown in UI buttons/headings "referee_group": "cds-ph-ep-publication", # CERN e-group slug "report_number": { - "prefix": "CERN-EP", # literal prefix, e.g. "CERN-EP" + "prefix": "CERN-TH-EP", # literal prefix, e.g. "CERN-EP" + "include_year": True, # append the current year after prefix + "counter_digits": 3, # zero-padding width, e.g. 3 → "001" + }, + }, + "7c568753-550b-4b48-8d76-461181973100": { + # aleph + "label": "EP approval", # shown in UI buttons/headings + "referee_group": "cds-ph-ep-publication", # CERN e-group slug + "report_number": { + "prefix": "CERN-TH-EP", # literal prefix, e.g. "CERN-EP" + "include_year": True, # append the current year after prefix + "counter_digits": 3, # zero-padding width, e.g. 3 → "001" + }, + }, + "bb0ac2d2-b90b-498e-bf64-1986791d1032": { + # l3 + "label": "EP approval", # shown in UI buttons/headings + "referee_group": "cds-ph-ep-publication", # CERN e-group slug + "report_number": { + "prefix": "CERN-PH-EP", # literal prefix, e.g. "CERN-EP" + "include_year": True, # append the current year after prefix + "counter_digits": 3, # zero-padding width, e.g. 3 → "001" + }, + }, + "473e34c5-4fe1-44fc-a4c2-3305cf6adcba": { + # opal + "label": "EP approval", # shown in UI buttons/headings + "referee_group": "cds-ph-ep-publication", # CERN e-group slug + "report_number": { + "prefix": "CERN-PH-EP", # literal prefix, e.g. "CERN-EP" + "include_year": True, # append the current year after prefix + "counter_digits": 3, # zero-padding width, e.g. 3 → "001" + }, + }, + "b6553d89-ea62-4a7c-9f5b-e76b5bfdb733": { + # delphi + "label": "EP approval", # shown in UI buttons/headings + "referee_group": "cds-ph-ep-publication", # CERN e-group slug + "report_number": { + "prefix": "CERN-PH-EP", # literal prefix, e.g. "CERN-EP" "include_year": True, # append the current year after prefix "counter_digits": 3, # zero-padding width, e.g. 3 → "001" }, diff --git a/cds_migrator_kit/rdm/records/load/approval_request.py b/cds_migrator_kit/rdm/records/load/approval_request.py index db4fdc6c..09a80b72 100644 --- a/cds_migrator_kit/rdm/records/load/approval_request.py +++ b/cds_migrator_kit/rdm/records/load/approval_request.py @@ -18,7 +18,10 @@ from invenio_pidstore.models import PersistentIdentifier, PIDStatus from invenio_rdm_records.records.api import RDMParent from invenio_records_resources.services.uow import RecordCommitOp -from invenio_requests.customizations.event_types import LogEventType +from invenio_requests.customizations.event_types import ( + LogEventType, + ReviewersUpdatedType, +) from invenio_requests.proxies import current_events_service, current_requests_service from invenio_requests.resolvers.registry import ResolverRegistry @@ -26,6 +29,7 @@ EP_APPROVAL_WAITING_STATUS = "waiting" EP_APPROVAL_APPROVED_STATUS = "approved" +EP_APPROVAL_REVIEWING_STATUS = "reviewing" class ApprovalRequest: @@ -45,13 +49,14 @@ def __init__( self.resource_type = resource_type self.dry_run = dry_run self.waiting_entry = None + self.reviewing_entry = None self.approved_entry = None self.report_number = None self.approved_at = None def validate(self): """Validate EP approval data before creating any records.""" - waiting_entry, approved_entry, report_number = self._parse_history() + waiting_entry, approved_entry, reviewing, report_number = self._parse_history() existing = self._existing_request() if existing: @@ -69,10 +74,16 @@ def validate(self): self.waiting_entry = waiting_entry self.approved_entry = approved_entry + self.reviewing_entry = reviewing self.report_number = report_number - def create(self, restricted_record_state): - """Create and approve EP approval request after restricted record exists.""" + def create(self, restricted_record_state, uow=None): + """Create and approve EP approval request after restricted record exists. + + If ``uow`` is provided, the request is registered on it without + committing, so the caller can group this atomically with other + operations (e.g. the public record creation and linking). + """ if self.dry_run: return @@ -91,14 +102,15 @@ def create(self, restricted_record_state): self._create_request( restricted_recid, restricted_parent, + uow=uow, ) self._mint_apprn_pid(restricted_record_state["latest_version_object_uuid"]) def _parse_history(self): """Return waiting/approved history entries and the report number.""" - if len(self.ep_approval) != 2: + if len(self.ep_approval) > 3: raise UnexpectedValue( - message="EP approval history has more/less than 2 entries", + message="EP approval history has more/less than 3 entries", stage="load", priority="critical", ) @@ -111,6 +123,14 @@ def _parse_history(self): ), None, ) + reviewing = next( + ( + item + for item in history + if item.get("status") == EP_APPROVAL_REVIEWING_STATUS + ), + None, + ) approved = next( ( item @@ -163,7 +183,7 @@ def _parse_history(self): priority="critical", ) - return waiting, approved, report_number + return waiting, approved, reviewing, report_number @staticmethod def parse_legacy_datetime(value): @@ -278,6 +298,43 @@ def _create_accept_log_event(self, request, uow): uow.register(RecordCommitOp(event, indexer=current_events_service.indexer)) + + def _create_reviewing_log_event(self, request, uow): + """Create the reviewers-updated timeline event with the legacy reviewer as created_by.""" + if not self.reviewing_entry: + return + + reviewer_ref = self._resolve_user_by_email( + self.reviewing_entry.get("submitted_by"), + "reviewer", + ) + request.reviewers = [reviewer_ref] + + event = current_events_service.record_cls.create( + {}, + request=request.model, + request_id=str(request.id), + type=ReviewersUpdatedType, + ) + event.update( + { + "payload": { + "event": "reviewers_updated", + "content": self.reviewing_entry.get("description", ""), + "reviewers": [reviewer_ref], + } + } + ) + event.created_by = ResolverRegistry.resolve_entity_proxy( + reviewer_ref, raise_=True + ) + + reviewing_at = self.parse_legacy_datetime(self.reviewing_entry.get("date")) + if reviewing_at: + event.model.created = reviewing_at + + uow.register(RecordCommitOp(event, indexer=current_events_service.indexer)) + def _apply_approved_entry(self, request, uow): """Update an existing request to accepted using the legacy approved entry.""" payload = dict(request.get("payload") or {}) @@ -291,38 +348,52 @@ def _apply_approved_entry(self, request, uow): self._create_accept_log_event(request, uow) - def _create_request(self, restricted_recid, restricted_parent): - """Create request from waiting entry, then update it with approved entry.""" + def _create_request(self, restricted_recid, restricted_parent, uow=None): + """Create request from waiting entry, then update it with approved entry. + + If ``uow`` is provided, it is used as-is and left uncommitted for the + caller to commit; otherwise a unit of work is created and committed + here. + """ + if uow is not None: + self._build_request(restricted_recid, restricted_parent, uow) + return + + with UnitOfWork() as inner_uow: + self._build_request(restricted_recid, restricted_parent, inner_uow) + inner_uow.commit() + + def _build_request(self, restricted_recid, restricted_parent, uow): + """Register the request creation and its updates on the given uow.""" expires_at = self.parse_legacy_datetime(self.waiting_entry.get("deadline")) referee_group = self._get_referee_group(restricted_parent) - with UnitOfWork() as uow: - request_item = current_requests_service.create( - system_identity, - data={ - "title": f'EP approval for "{self.title}"', - "payload": {}, - }, - request_type=CommitteeApprovalRequest, - receiver={"group": referee_group}, - creator=self._resolve_user_by_email( - self.waiting_entry.get("submitted_by"), "submitter" - ), - topic={"record": restricted_recid}, - expires_at=expires_at, - uow=uow, - ) - request = request_item._record - request.number = f"lrecid:{self.legacy_recid}:ep-approval" - request.status = "submitted" + request_item = current_requests_service.create( + system_identity, + data={ + "title": f'EP approval for "{self.title}"', + "payload": {}, + }, + request_type=CommitteeApprovalRequest, + receiver={"group": referee_group}, + creator=self._resolve_user_by_email( + self.waiting_entry.get("submitted_by"), "submitter" + ), + topic={"record": restricted_recid}, + expires_at=expires_at, + uow=uow, + ) + request = request_item._record + request.number = f"lrecid:{self.legacy_recid}:ep-approval" + request.status = "submitted" - submitted_at = self.parse_legacy_datetime(self.waiting_entry.get("date")) - if submitted_at: - request.model.created = submitted_at + submitted_at = self.parse_legacy_datetime(self.waiting_entry.get("date")) + if submitted_at: + request.model.created = submitted_at - self._apply_approved_entry(request, uow) + self._create_reviewing_log_event(request, uow) + self._apply_approved_entry(request, uow) - uow.register( - RecordCommitOp(request, indexer=current_requests_service.indexer) - ) - uow.commit() + uow.register( + RecordCommitOp(request, indexer=current_requests_service.indexer) + ) diff --git a/cds_migrator_kit/rdm/records/load/ep_approval_entry.py b/cds_migrator_kit/rdm/records/load/ep_approval_entry.py index 365b3734..de15e082 100644 --- a/cds_migrator_kit/rdm/records/load/ep_approval_entry.py +++ b/cds_migrator_kit/rdm/records/load/ep_approval_entry.py @@ -47,6 +47,20 @@ def _apply_metadata(self, split): def _apply_entry_modifications(self, split): """Apply record/parent level modifications.""" + @staticmethod + def _is_restricted_file(file_data): + """Return whether a file belongs to the restricted split. + + A file is restricted either because it is an EPPHAPP draft file, or + because it carries its own file-level access restriction independent + of the EPPHAPP workflow (e.g. a record that was never restricted as + a whole, but ships a mix of public and individually-restricted + files). + """ + return bool( + file_data.get("type") == EPPHAPP_FILE_TYPE or file_data.get("access") + ) + def _log_removed_identifiers(self, removed, split_type): recid = self.entry.get("record", {}).get("recid") self.migration_logger.add_information( @@ -96,20 +110,9 @@ def _build_versions(self, split): current_version_files = OrderedDict() for key, file_data in version_data.get("files", {}).items(): - if file_data.get("type") == EPPHAPP_FILE_TYPE: + if self._is_restricted_file(file_data): continue - if file_data.get("access"): - raise UnexpectedValue( - message=( - "Public split contains restricted files after excluding " - f"EPPHAPP files: {[key]}" - ), - stage="load", - recid=split["record"]["recid"], - priority="critical", - ) - current_version_files[key] = deepcopy(file_data) if not current_version_files: @@ -192,9 +195,32 @@ def _add_cern_scientific_community(self, entry): class RestrictedEntry(MetadataEntry): """Build the restricted EP approval split entry.""" - def _has_epphapp_files(self, split): + def _apply_entry_modifications(self, split): + self._remove_cern_scientific_community(split) + + def _remove_cern_scientific_community(self, entry): + """Drop the CERN Scientific community from the restricted split. + + The restricted record holds the internal-only EPPHAPP draft and must + not be discoverable via the broader community; only PublicEntry adds + CDS_CERN_SCIENTIFIC_COMMUNITY_ID (see _add_cern_scientific_community). + """ + communities = entry.get("parent", {}).get("json", {}).get("communities", {}) + ids = [ + cid + for cid in communities.get("ids", []) + if cid != CDS_CERN_SCIENTIFIC_COMMUNITY_ID + ] + communities["ids"] = ids + if communities.get("default") == CDS_CERN_SCIENTIFIC_COMMUNITY_ID: + communities["default"] = ids[0] if ids else None + entry.setdefault("parent", {}).setdefault("json", {})[ + "communities" + ] = communities + + def _has_restricted_files(self, split): return any( - file_data.get("type") == EPPHAPP_FILE_TYPE + self._is_restricted_file(file_data) for version_data in split.get("versions", {}).values() for file_data in version_data.get("files", {}).values() ) @@ -203,14 +229,14 @@ def _build_versions(self, split): new_versions = OrderedDict() versioned_files = OrderedDict() previous_signature = None - has_epphapp_files = self._has_epphapp_files(split) + has_restricted_files = self._has_restricted_files(split) - if not has_epphapp_files: + if not has_restricted_files: self.migration_logger.add_information( split["record"]["recid"], { "message": ( - "No EPPHAPP files found; public files used for the " + "No restricted files found; public files used for the " "restricted record." ), "value": "public files", @@ -221,10 +247,11 @@ def _build_versions(self, split): current_version_files = OrderedDict() for key, file_data in version_data.get("files", {}).items(): - is_epphapp = file_data.get("type") == EPPHAPP_FILE_TYPE + is_restricted = self._is_restricted_file(file_data) - # If draft file exists, use that otherwise use the public files. - if not is_epphapp and has_epphapp_files: + # If restricted files exist, use only those; otherwise fall + # back to using all (public) files for the restricted record. + if not is_restricted and has_restricted_files: continue current_version_files[key] = deepcopy(file_data) diff --git a/cds_migrator_kit/rdm/records/load/ep_approval_load.py b/cds_migrator_kit/rdm/records/load/ep_approval_load.py index c7550ad4..84307777 100644 --- a/cds_migrator_kit/rdm/records/load/ep_approval_load.py +++ b/cds_migrator_kit/rdm/records/load/ep_approval_load.py @@ -38,7 +38,6 @@ def __init__( self, db_uri, data_dir, - tmp_dir, entries=None, dry_run=False, legacy_pids_to_redirect=None, @@ -68,11 +67,29 @@ def _load(self, entry): 1. create the restricted record 2. create and approve the EP approval request 3. create the public record and link both with related_identifiers + + Steps 1-4 are grouped into a single unit of work so a failure at any + point rolls back everything, instead of leaving an orphaned + restricted record and/or approval request committed behind. """ if not entry: return try: recid = entry.get("record", {}).get("recid") + + # The same legacy recid can be cross-listed under multiple EP + # collections (e.g. a joint ALEPH/DELPHI/L3/OPAL paper appears in + # all four experiments' dumps). Once one pass has fully migrated + # it, later passes should skip cleanly instead of failing on the + # already-created approval request. + if CDSRecordServiceLoad._have_migrated_recid(recid): + self.migration_logger.add_information( + recid, + state={"message": "Record already migrated", "value": recid}, + ) + self.migration_logger.finalise_record(recid) + return + ep_approval = entry.get("record", {}).get("ep_approval") if not ep_approval: raise UnexpectedValue( @@ -105,7 +122,6 @@ def _load(self, entry): migration_logger=self.migration_logger, ).build() - # 1. Create restricted record restricted_record_service = CDSRecordServiceLoad( dry_run=self.dry_run, collection=self.collection, @@ -115,12 +131,6 @@ def _load(self, entry): legacy_pids_to_redirect=self.legacy_pids_to_redirect, _is_final_record=False, ) - restricted_record_state = restricted_record_service._load(restricted_entry) - - # 2. Create and approve EP approval request - self.approval_request.create(restricted_record_state) - - # 3. Create public record and link both records public_record_service = CDSRecordServiceLoad( dry_run=self.dry_run, collection=self.collection, @@ -130,31 +140,73 @@ def _load(self, entry): legacy_pids_to_redirect=self.legacy_pids_to_redirect, _is_final_record=True, ) - public_record_state = public_record_service._load(public_entry) - if not self.dry_run: + if self.dry_run: + # 1. Create restricted record + restricted_record_state = restricted_record_service._load( + restricted_entry + ) + # 2. Create and approve EP approval request + self.approval_request.create(restricted_record_state) + # 3. Create public record + public_record_service._load(public_entry) + return + + with UnitOfWork(db.session) as uow: + # 1. Create restricted record + restricted_record_state = restricted_record_service._load( + restricted_entry, uow=uow + ) + + # 2. Create and approve EP approval request + self.approval_request.create(restricted_record_state, uow=uow) + + # 3. Create public record + public_record_state = public_record_service._load( + public_entry, uow=uow + ) + if not public_record_state: + raise UnexpectedValue( + message="Public record is required for EP approval.", + stage="load", + recid=recid, + priority="critical", + ) + # Link the records with related_identifiers self._append_related_identifier( public_record_state["latest_version"], restricted_record_state["latest_version"], "isversionof", self.approval_request.resource_type, + uow=uow, ) self._append_related_identifier( restricted_record_state["latest_version"], public_record_state["latest_version"], "isvariantformof", self.approval_request.resource_type, + uow=uow, ) - # 4. Link the records with related_identifiers + # 4. Write EP approval metadata on both parents self._link_parent_ep_approvals( - restricted_record_state, public_record_state, legacy_recid=recid + restricted_record_state, + public_record_state, + legacy_recid=recid, + uow=uow, ) public_record_state["internal_version"] = restricted_record_state[ "latest_version" ] + + uow.commit() + + # The public record is the final one; finalise it only now + # that the whole split has actually committed (see the + # matching `uow is None` guard in CDSRecordServiceLoad._load). + self.migration_logger.finalise_record(recid) except (UnexpectedValue, ManualImportRequired) as e: self.migration_logger.add_log(e, record=entry) except Exception as e: @@ -168,51 +220,71 @@ def _load(self, entry): self.migration_logger.add_log(exc, record=entry) def _link_parent_ep_approvals( - self, restricted_record_state, public_record_state, legacy_recid + self, restricted_record_state, public_record_state, legacy_recid, uow=None ): - """Write parent metadata and link public/restricted records.""" - approved_entry = self.approval_request.approved_entry - report_number = self.approval_request.report_number - approval_iso = self.approval_request.approved_at.isoformat() + """Write parent metadata and link public/restricted records. - if not self.dry_run: - if not restricted_record_state or not public_record_state: - raise UnexpectedValue( - message="Both public and restricted records are required for EP approval.", - stage="load", - recid=legacy_recid, - priority="critical", - ) - restricted_recid = restricted_record_state["latest_version"] - public_recid = public_record_state["latest_version"] - restricted_parent = RDMParent.get_record( - restricted_record_state["parent_object_uuid"] + If ``uow`` is provided, the writes are registered on it without + committing, so the caller can group this atomically with other + operations. + """ + if self.dry_run: + return + + if not restricted_record_state or not public_record_state: + raise UnexpectedValue( + message="Both public and restricted records are required for EP approval.", + stage="load", + recid=legacy_recid, + priority="critical", ) - public_parent = RDMParent.get_record( - public_record_state["parent_object_uuid"] + + if uow is not None: + self._write_parent_ep_approvals( + restricted_record_state, public_record_state, uow ) + return - with UnitOfWork() as uow: - self._write_parent_ep_approval( - restricted_parent, - { - "reportnumber": report_number, - "datetime": approval_iso, - "approved_internal_version": restricted_recid, - "approved_public_version": public_recid, - "source_public_version": restricted_recid, - }, - uow, - ) - self._write_parent_ep_approval( - public_parent, - { - "reportnumber": report_number, - "source_internal_version": restricted_recid, - }, - uow, - ) - uow.commit() + with UnitOfWork() as inner_uow: + self._write_parent_ep_approvals( + restricted_record_state, public_record_state, inner_uow + ) + inner_uow.commit() + + def _write_parent_ep_approvals( + self, restricted_record_state, public_record_state, uow + ): + """Register the EP approval parent writes for both records on the given uow.""" + report_number = self.approval_request.report_number + approval_iso = self.approval_request.approved_at.isoformat() + restricted_recid = restricted_record_state["latest_version"] + public_recid = public_record_state["latest_version"] + restricted_parent = RDMParent.get_record( + restricted_record_state["parent_object_uuid"] + ) + public_parent = RDMParent.get_record( + public_record_state["parent_object_uuid"] + ) + + self._write_parent_ep_approval( + restricted_parent, + { + "reportnumber": report_number, + "datetime": approval_iso, + "approved_internal_version": restricted_recid, + "approved_public_version": public_recid, + "source_public_version": restricted_recid, + }, + uow, + ) + self._write_parent_ep_approval( + public_parent, + { + "reportnumber": report_number, + "source_internal_version": restricted_recid, + }, + uow, + ) def _write_parent_ep_approval( self, @@ -227,10 +299,17 @@ def _write_parent_ep_approval( uow.register(ParentRecordCommitOp(parent)) def _append_related_identifier( - self, record_id, target_id, relation_id, resource_type + self, record_id, target_id, relation_id, resource_type, uow=None ): - """Append the related identifier to the record.""" - draft = current_rdm_records_service.edit(system_identity, id_=record_id) + """Append the related identifier to the record. + + If ``uow`` is provided, it is passed through to the record service + calls so they register on it without committing, letting the caller + group this atomically with other operations. + """ + draft = current_rdm_records_service.edit( + system_identity, id_=record_id, uow=uow + ) data = draft.data related = list(data.get("metadata", {}).get("related_identifiers", [])) @@ -244,9 +323,9 @@ def _append_related_identifier( related.append(entry) data.setdefault("metadata", {})["related_identifiers"] = related current_rdm_records_service.update_draft( - system_identity, id_=draft.id, data=data + system_identity, id_=draft.id, data=data, uow=uow ) - current_rdm_records_service.publish(system_identity, id_=draft.id) + current_rdm_records_service.publish(system_identity, id_=draft.id, uow=uow) return True def _cleanup(self, *args, **kwargs): diff --git a/cds_migrator_kit/rdm/records/load/load.py b/cds_migrator_kit/rdm/records/load/load.py index 5cac99eb..b0d1bc0b 100644 --- a/cds_migrator_kit/rdm/records/load/load.py +++ b/cds_migrator_kit/rdm/records/load/load.py @@ -257,6 +257,14 @@ 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( + kw in specific_file_restrictions for kw in ("firerole:", "allow ") + ): + # bare CERN e-group name, e.g. + # "cds-ph-ep-publications-referee-non-lhc [CERN]" + groups.add(_normalize_group_name(specific_file_restrictions)) else: if not any( kw in specific_file_restrictions @@ -776,8 +784,13 @@ def _after_load_clc_sync(self, record_state): ) db.session.add(sync) - def _load(self, entry): - """Use the services to load the entries.""" + def _load(self, entry, uow=None): + """Use the services to load the entries. + + If ``uow`` is provided, operations are registered on it without + committing, so the caller can group this load atomically with other + operations (e.g. the EP approval record split). + """ if entry: recid = entry.get("record", {}).get("recid", {}) if self._should_skip_recid(recid): @@ -803,16 +816,28 @@ def _load(self, entry): if self.dry_run: self._dry_load(entry) recid_state_after_load = 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._after_load_clc_sync(recid_state_after_load) else: - with UnitOfWork(db.session) as uow: - recid_state_after_load = self._load_versions(entry, uow) + with UnitOfWork(db.session) as 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 ) self._after_load_clc_sync(recid_state_after_load) - uow.commit() - if self._is_final_record: + inner_uow.commit() + if self._is_final_record and uow is None: + # When an external uow is provided, the caller owns the + # commit boundary and is responsible for finalising the + # record only after it actually commits. self.migration_logger.finalise_record(recid) return recid_state_after_load except (UnexpectedValue, ManualImportRequired) as e: diff --git a/cds_migrator_kit/rdm/records/transform/xml_processing/rules/research.py b/cds_migrator_kit/rdm/records/transform/xml_processing/rules/research.py index 5444ccd0..3cb4794f 100644 --- a/cds_migrator_kit/rdm/records/transform/xml_processing/rules/research.py +++ b/cds_migrator_kit/rdm/records/transform/xml_processing/rules/research.py @@ -823,8 +823,8 @@ def ep_approval(self, key, value): ep_report_number = value.get("b", "").strip() stamp_info = value.get("g", "").strip() doc_type = value.get("c", "").strip() - if status not in ["waiting", "approved"]: - raise UnexpectedValue(subfield="a", field=key, value=value) + if status not in ["waiting", "approved", "reviewing"]: + raise UnexpectedValue(subfield="s", field=key, value=value) return { k: v for k, v in { diff --git a/cds_migrator_kit/rdm/streams.yaml b/cds_migrator_kit/rdm/streams.yaml index f10b9d13..41460876 100644 --- a/cds_migrator_kit/rdm/streams.yaml +++ b/cds_migrator_kit/rdm/streams.yaml @@ -14,6 +14,20 @@ records: - "c2c46ab3-5fb4-4d86-83c6-5d9dc8392d6f" load: legacy_pids_to_redirect: cds_migrator_kit/rdm/data/lep_exp/aleph/duplicated_pids.json + aleph_ep: + data_dir: cds_migrator_kit/rdm/data/lep_ep/aleph + plots: true + create_inclusion_request: true + extract: + dirpath: cds_migrator_kit/rdm/data/lep_ep/aleph/dump/ + transform: + files_dump_dir: cds_migrator_kit/rdm/data/lep_ep/aleph/files/ + missing_users: cds_migrator_kit/rdm/data/users + communities_ids: + - "7c568753-550b-4b48-8d76-461181973100" + - "c2c46ab3-5fb4-4d86-83c6-5d9dc8392d6f" + load: + legacy_pids_to_redirect: cds_migrator_kit/rdm/data/lep_ep/aleph/duplicated_pids.json aleph_drafts: data_dir: cds_migrator_kit/rdm/data/lep_exp/aleph_drafts restricted: "True" @@ -42,6 +56,20 @@ records: - "c2c46ab3-5fb4-4d86-83c6-5d9dc8392d6f" load: legacy_pids_to_redirect: cds_migrator_kit/rdm/data/lep_exp/l3/duplicated_pids.json + l3_ep: + data_dir: cds_migrator_kit/rdm/data/lep_ep/l3 + plots: true + create_inclusion_request: true + extract: + dirpath: cds_migrator_kit/rdm/data/lep_ep/l3/dump/ + transform: + files_dump_dir: cds_migrator_kit/rdm/data/lep_ep/l3/files/ + missing_users: cds_migrator_kit/rdm/data/users + communities_ids: + - "bb0ac2d2-b90b-498e-bf64-1986791d1032" + - "c2c46ab3-5fb4-4d86-83c6-5d9dc8392d6f" + load: + legacy_pids_to_redirect: cds_migrator_kit/rdm/data/lep_ep/l3/duplicated_pids.json opal: data_dir: cds_migrator_kit/rdm/data/lep_exp/opal plots: true @@ -56,6 +84,20 @@ records: - "c2c46ab3-5fb4-4d86-83c6-5d9dc8392d6f" load: legacy_pids_to_redirect: cds_migrator_kit/rdm/data/lep_exp/opal/duplicated_pids.json + opal_ep: + data_dir: cds_migrator_kit/rdm/data/lep_ep/opal + plots: true + create_inclusion_request: true + extract: + dirpath: cds_migrator_kit/rdm/data/lep_ep/opal/dump/ + transform: + files_dump_dir: cds_migrator_kit/rdm/data/lep_ep/opal/files/ + missing_users: cds_migrator_kit/rdm/data/users + communities_ids: + - "473e34c5-4fe1-44fc-a4c2-3305cf6adcba" + - "c2c46ab3-5fb4-4d86-83c6-5d9dc8392d6f" + load: + legacy_pids_to_redirect: cds_migrator_kit/rdm/data/lep_ep/opal/duplicated_pids.json delphi: data_dir: cds_migrator_kit/rdm/data/lep_exp/delphi plots: true @@ -70,6 +112,20 @@ records: - "c2c46ab3-5fb4-4d86-83c6-5d9dc8392d6f" load: legacy_pids_to_redirect: cds_migrator_kit/rdm/data/lep_exp/delphi/duplicated_pids.json + delphi_ep: + data_dir: cds_migrator_kit/rdm/data/lep_ep/delphi + plots: true + create_inclusion_request: true + extract: + dirpath: cds_migrator_kit/rdm/data/lep_ep/delphi/dump/ + transform: + files_dump_dir: cds_migrator_kit/rdm/data/lep_ep/delphi/files/ + missing_users: cds_migrator_kit/rdm/data/users + communities_ids: + - "b6553d89-ea62-4a7c-9f5b-e76b5bfdb733" + - "c2c46ab3-5fb4-4d86-83c6-5d9dc8392d6f" + load: + legacy_pids_to_redirect: cds_migrator_kit/rdm/data/lep_ep/delphi/duplicated_pids.json delphi_priv: data_dir: cds_migrator_kit/rdm/data/lep_exp/delphi_priv restricted: "True" @@ -452,8 +508,6 @@ records: faser-ep: plots: true data_dir: cds_migrator_kit/rdm/data/faser - tmp_dir: cds_migrator_kit/rdm/tmp/faser - log_dir: cds_migrator_kit/rdm/log/faser create_inclusion_request: true extract: dirpath: cds_migrator_kit/rdm/data/faser-ep/dump/ @@ -490,7 +544,7 @@ records: files_dump_dir: cds_migrator_kit/rdm/data/dep_ts/files/ missing_users: cds_migrator_kit/rdm/data/users communities_ids: - - "058c4659-2ad5-47c7-98ad-5e42333b09b3" + - "97ffcddd-ee1f-4e21-ac7b-b78612dc81d8" mt_dep: plots: true data_dir: cds_migrator_kit/rdm/data/dep_mt @@ -500,7 +554,7 @@ records: files_dump_dir: cds_migrator_kit/rdm/data/dep_mt/files/ missing_users: cds_migrator_kit/rdm/data/users communities_ids: - - "f508d592-580f-4d3e-8083-d8119cac1913" + - "ad4469fa-5a17-41b9-825f-8010c379ba23" sb_dep: plots: true data_dir: cds_migrator_kit/rdm/data/dep_sb @@ -510,7 +564,7 @@ records: files_dump_dir: cds_migrator_kit/rdm/data/dep_sb/files/ missing_users: cds_migrator_kit/rdm/data/users communities_ids: - - "87d73294-8f1f-4a63-a75e-21a87aa25980" + - "ca71c1ab-2eb3-4a43-a593-b8b140fe6123" est_dep: plots: true data_dir: cds_migrator_kit/rdm/data/dep_est @@ -520,4 +574,4 @@ records: files_dump_dir: cds_migrator_kit/rdm/data/dep_est/files/ missing_users: cds_migrator_kit/rdm/data/users communities_ids: - - "9a9c7474-18d2-41ee-b3ec-0708139dd0ee" + - "d9c138bd-d39c-4797-a06c-c6f8db45f3ce" diff --git a/scripts/opensearch-init.sh b/scripts/opensearch-init.sh index 4b97766e..f0929d0e 100755 --- a/scripts/opensearch-init.sh +++ b/scripts/opensearch-init.sh @@ -32,10 +32,6 @@ usage() { [[ "$ACTION" == "dump" || "$ACTION" == "restore" ]] || usage -if ! command -v multielasticdump &>/dev/null; then - echo "multielasticdump not found. Install it with: npm install -g multielasticdump" >&2 - exit 1 -fi if ! command -v elasticdump &>/dev/null; then echo "elasticdump not found. Install it with: npm install -g elasticdump" >&2 exit 1 @@ -57,7 +53,25 @@ case "$ACTION" in dump) mkdir -p "$DIR" echo "Dumping index data (no mapping/settings/alias) from $HOST to $DIR ..." - multielasticdump --direction=dump --input="$HOST" --output="$DIR" --includeType=data + # Not using multielasticdump here: it discovers indices via GET + # /_aliases, then checks `'error' in response` to detect an ES-style + # error payload. If the cluster happens to have a real index literally + # named "error" (e.g. invenio-logging's error log index), that check + # false-positives on the alias listing itself and multielasticdump exits + # before dumping anything. /_cat/indices isn't keyed by index name, so + # it's immune to that collision - dump each real index individually + # with plain elasticdump instead, which produces the same per-index + # files multielasticdump would have. + indices="$(curl -sf "$HOST/_cat/indices?h=index" | sort)" + if [[ -z "$indices" ]]; then + echo "No indices found at $HOST" >&2 + exit 1 + fi + while IFS= read -r index; do + [[ -n "$index" ]] || continue + echo "=== Dumping $index ===" + elasticdump --input="$HOST/$index" --output="$DIR/$index.json" --type=data + done <<<"$indices" ;; restore) [[ -d "$DIR" ]] || { echo "Dump directory not found: $DIR" >&2; exit 1; } diff --git a/scripts/snapshot.sh b/scripts/snapshot.sh new file mode 100755 index 00000000..6bf9d11e --- /dev/null +++ b/scripts/snapshot.sh @@ -0,0 +1,90 @@ +#!/usr/bin/env bash +# Dump/restore Postgres + OpenSearch together as one named snapshot, so both +# stay tied to the same point in time instead of drifting apart. +# +# Usage: +# scripts/snapshot.sh dump +# scripts/snapshot.sh restore +# +# Layout for a snapshot named "before-ep-approval-fix": +# snapshots/before-ep-approval-fix/db.sql +# snapshots/before-ep-approval-fix/opensearch/*.json +set -euo pipefail + +SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" +SNAPSHOTS_ROOT="${SNAPSHOTS_ROOT:-snapshots}" + +# opensearch-init.sh's restore path runs `invenio index destroy`/`invenio +# index init`, which need the actual app (cds-rdm) virtualenv, not whatever +# happens to be active in the caller's shell (e.g. cds-migrator-kit's own). +CDS_RDM_VENV_ACTIVATE="${CDS_RDM_VENV_ACTIVATE:-$SCRIPT_DIR/../../cds-rdm/.venv/bin/activate}" + +ACTION="${1:-}" +NAME="${2:-}" + +usage() { + echo "Usage: $0 {dump|restore} " >&2 + exit 1 +} + +[[ "$ACTION" == "dump" || "$ACTION" == "restore" ]] || usage +[[ -n "$NAME" ]] || usage + +[[ -f "$CDS_RDM_VENV_ACTIVATE" ]] || { + echo "cds-rdm virtualenv not found at $CDS_RDM_VENV_ACTIVATE" >&2 + echo "Set CDS_RDM_VENV_ACTIVATE to its activate script if cds-rdm lives elsewhere." >&2 + exit 1 +} + +run_opensearch_init() { + # Subshell: activate cds-rdm's venv only for this call, without leaking it + # into the rest of snapshot.sh (postgres-init.sh needs no app venv at all). + ( + # shellcheck disable=SC1090 + source "$CDS_RDM_VENV_ACTIVATE" + "$SCRIPT_DIR/opensearch-init.sh" "$@" + ) +} + +SNAPSHOT_DIR="$SNAPSHOTS_ROOT/$NAME" +DB_FILE="$SNAPSHOT_DIR/db.sql" +OS_DIR="$SNAPSHOT_DIR/opensearch" + +case "$ACTION" in + dump) + mkdir -p "$SNAPSHOT_DIR" + echo "=== Snapshot '$NAME': dumping Postgres ===" + "$SCRIPT_DIR/postgres-init.sh" dump "$DB_FILE" + echo "=== Snapshot '$NAME': dumping OpenSearch ===" + run_opensearch_init dump "$OS_DIR" + echo "Snapshot '$NAME' written to $SNAPSHOT_DIR" + ;; + restore) + [[ -f "$DB_FILE" ]] || { echo "Snapshot db dump not found: $DB_FILE" >&2; exit 1; } + [[ -d "$OS_DIR" ]] || { echo "Snapshot OpenSearch dump not found: $OS_DIR" >&2; exit 1; } + + if [[ "${FORCE:-}" != "1" ]]; then + read -r -p "This will DESTROY the current database and OpenSearch indices and restore snapshot '$NAME'. Continue? [y/N] " reply + [[ "$reply" =~ ^[Yy]$ ]] || { echo "Aborted."; exit 1; } + fi + + echo "=== Snapshot '$NAME': restoring Postgres + OpenSearch in parallel ===" + (FORCE=1 "$SCRIPT_DIR/postgres-init.sh" restore "$DB_FILE" 2>&1 | sed -u 's/^/[postgres] /') & + pg_pid=$! + (run_opensearch_init restore "$OS_DIR" 2>&1 | sed -u 's/^/[opensearch] /') & + os_pid=$! + + pg_status=0 + os_status=0 + wait "$pg_pid" || pg_status=$? + wait "$os_pid" || os_status=$? + + if [[ "$pg_status" -ne 0 || "$os_status" -ne 0 ]]; then + echo "Restore failed (postgres exit=$pg_status, opensearch exit=$os_status)" >&2 + exit 1 + fi + echo "Snapshot '$NAME' restored." + ;; +esac + +echo "Done." From 73a888b1fc204d3b13d6855ca54c0cf27b1077e9 Mon Sep 17 00:00:00 2001 From: Karolina Przerwa Date: Wed, 22 Jul 2026 17:12:39 +0200 Subject: [PATCH 3/3] change(files): add original filepath to migrated metadata --- .../rdm/records/transform/transform.py | 5 +++-- tests/cds-rdm/test_ep_approval_entry.py | 18 +++++++++++++----- tests/cds-rdm/test_load_reviewers.py | 2 +- 3 files changed, 17 insertions(+), 8 deletions(-) diff --git a/cds_migrator_kit/rdm/records/transform/transform.py b/cds_migrator_kit/rdm/records/transform/transform.py index d4752072..13320266 100644 --- a/cds_migrator_kit/rdm/records/transform/transform.py +++ b/cds_migrator_kit/rdm/records/transform/transform.py @@ -569,9 +569,9 @@ def field_experiments(record_json, custom_fields_dict): subj.append({"subject": experiment}) json_output["metadata"]["subjects"] = subj raise UnexpectedValue( - subfield="u", + subfield="e", value=experiment, - field="author", + field="693", message=f"Experiment {experiment} not found", stage="vocabulary match", ) @@ -1017,6 +1017,7 @@ def compute_files(file_dump, versions_dict): "description": file_dump["description"], "name": file_dump["name"], "status": file_dump["status"], + "original_path": file_dump["path"], "comment": file_dump["comment"], }, "mimetype": file_dump["mime"], diff --git a/tests/cds-rdm/test_ep_approval_entry.py b/tests/cds-rdm/test_ep_approval_entry.py index d0d18470..6d0ba064 100644 --- a/tests/cds-rdm/test_ep_approval_entry.py +++ b/tests/cds-rdm/test_ep_approval_entry.py @@ -272,7 +272,11 @@ def test_public_raises_when_no_public_files(self): entry, _make_approval_request(), _make_migration_logger() ).build() - def test_public_raises_on_restricted_files(self): + def test_public_excludes_restricted_non_epphapp_files(self): + """A record that isn't restricted as a whole can still ship individually + restricted, non-EPPHAPP files (e.g. via 506__m); those must be excluded + from the public split (and land in RestrictedEntry instead), not raise. + """ versions = OrderedDict( [ ( @@ -288,6 +292,7 @@ def test_public_raises_on_restricted_files(self): "type": "Main", "creation_date": "2020-01-01", }, + PUBLIC_FILE_KEY: _public_file(), }, "publication_date": "2020-01-01", "access": {"access_obj": {"record": None, "files": None}}, @@ -296,10 +301,13 @@ def test_public_raises_on_restricted_files(self): ] ) entry = _make_entry(versions) - with pytest.raises(UnexpectedValue, match="restricted files"): - PublicEntry( - entry, _make_approval_request(), _make_migration_logger() - ).build() + result = PublicEntry( + entry, _make_approval_request(), _make_migration_logger() + ).build() + + files = result["versions"][1]["files"] + assert "restricted.pdf" not in files + assert PUBLIC_FILE_KEY in files def test_public_multiple_distinct_versions(self): versions = OrderedDict( diff --git a/tests/cds-rdm/test_load_reviewers.py b/tests/cds-rdm/test_load_reviewers.py index 7858bd5f..2195b060 100644 --- a/tests/cds-rdm/test_load_reviewers.py +++ b/tests/cds-rdm/test_load_reviewers.py @@ -11,7 +11,7 @@ from invenio_accounts.testutils import create_test_user from cds_migrator_kit.errors import RecordFlaggedCuration -from cds_migrator_kit.rdm.records.load.load import ( +from cds_migrator_kit.rdm.records.transform.xml_processing.quality.reviewers import ( _is_email, _parse_reviewer_name, find_reviewer,