diff --git a/invenio.cfg b/invenio.cfg index feac7470..5af79c22 100644 --- a/invenio.cfg +++ b/invenio.cfg @@ -428,6 +428,18 @@ Available options are: 'local', 'development', 'sandbox', 'production'. If set to 'production', the side banner is not shown. """ +CDS_HARVESTER_SANDBOX_ALLOW_CREATE_PROD_MISSING_RECORDS = False +"""Sandbox only: create a record that already exists on prod but is missing here. + +Needs ``CDS_ENVIRONMENT_NAME == "sandbox"``. When both are set we: +1. allow create even if the entry has a CDS DOI +2. ask prod for the parent and version ids from the INSPIRE CDSRDM value +3. create the sandbox record with those same ids +""" + +CDS_HARVESTER_PROD_API_URL = "https://repository.cern" +"""Prod base URL used to resolve a CDSRDM id into parent.id and version id.""" + # Invenio-Files-REST # ================== XROOTD_ENABLED = False diff --git a/site/cds_rdm/config.py b/site/cds_rdm/config.py index ac57e7bc..95bcfc82 100644 --- a/site/cds_rdm/config.py +++ b/site/cds_rdm/config.py @@ -67,3 +67,15 @@ CDS_HARVESTER_USER_EMAIL = None """Email of the INSPIRE harvester service user.""" + +CDS_HARVESTER_SANDBOX_ALLOW_CREATE_PROD_MISSING_RECORDS = False +"""Sandbox only: create a record that already exists on prod but is missing here. + +Needs ``CDS_ENVIRONMENT_NAME == "sandbox"``. When both are set we: +1. allow create even if the entry has a CDS DOI +2. ask prod for the parent and version ids from the INSPIRE CDSRDM value +3. create the sandbox record with those same ids +""" + +CDS_HARVESTER_PROD_API_URL = "https://repository.cern" +"""Prod base URL used to resolve a CDSRDM id into parent.id and version id.""" diff --git a/site/cds_rdm/inspire_harvester/load/draft.py b/site/cds_rdm/inspire_harvester/load/draft.py index a9b0845f..f9149e9a 100644 --- a/site/cds_rdm/inspire_harvester/load/draft.py +++ b/site/cds_rdm/inspire_harvester/load/draft.py @@ -9,10 +9,13 @@ from flask import current_app from invenio_db import db +from invenio_drafts_resources.services.records.uow import ParentRecordCommitOp from invenio_rdm_records.proxies import current_rdm_records_service from invenio_rdm_records.services.errors import ValidationErrorWithMessageAsList +from invenio_records_resources.services.uow import RecordCommitOp, UnitOfWork from invenio_vocabularies.datastreams.errors import WriterError from marshmallow import ValidationError +from sqlalchemy.exc import IntegrityError from cds_rdm.inspire_harvester.logger import ( format_validation_error, @@ -20,6 +23,32 @@ ) +def _remint_recid(obj, new_pid_value, uow): + """Swap the auto-minted recid for the one from prod. + + Create always mints fresh parent/version ids. On sandbox we overwrite + those values so the sandbox record keeps the same ids as prod. + + After resolve, ``obj.pid`` may be rebuilt from the record JSON and not + sit in the DB session. Merge it first so we UPDATE the real row instead + of INSERTing a second one (that looks like "pid already exists"). + """ + # pid from JSON is transient until merged into the session + type(obj).pid.session_merge(obj) + pid = obj.pid + try: + with uow.session.begin_nested(): + pid.pid_value = new_pid_value + uow.session.add(pid) + except IntegrityError as exc: + raise WriterError( + "Cannot reuse CDSRDM recid - already exists. " + f"| details: pid={new_pid_value}" + ) from exc + # Tell the PID field to write the new id into the record JSON as well. + obj.pid = pid + + class DraftLifecycleManager: """Manages draft creation, editing, versioning, and publishing.""" @@ -31,6 +60,30 @@ def create(self, entry): """Create a new draft from entry data.""" return current_rdm_records_service.create(self.identity, data=entry) + def create_reusing_prod_pids(self, entry, parent_pid, record_pid): + """Create a draft and keep prod's parent and version ids. + + Even one published version has two ids: a parent and a + version. We create normally, then replace both auto-minted + ids with the ones from prod, in one unit of work so it is atomic. + """ + with UnitOfWork() as uow: + draft = current_rdm_records_service.create( + self.identity, data=entry, uow=uow + ) + draft_obj = current_rdm_records_service.draft_cls.pid.resolve( + draft.id, registered_only=False + ) + # Parent first, then the version/record id. + _remint_recid(draft_obj.parent, parent_pid, uow) + uow.register(ParentRecordCommitOp(draft_obj.parent)) + _remint_recid(draft_obj, record_pid, uow) + uow.register(RecordCommitOp(draft_obj)) + uow.commit() + + # Return the draft under the reminted version id. + return current_rdm_records_service.read_draft(self.identity, record_pid) + def edit(self, record_pid): """Open an edit draft for an existing published record.""" return current_rdm_records_service.edit(self.identity, record_pid) diff --git a/site/cds_rdm/inspire_harvester/load/validator.py b/site/cds_rdm/inspire_harvester/load/validator.py index 476333c4..a1a4668f 100644 --- a/site/cds_rdm/inspire_harvester/load/validator.py +++ b/site/cds_rdm/inspire_harvester/load/validator.py @@ -51,7 +51,18 @@ class CdsDoiCreateRule(ValidationRule): """Block create when the entry carries a CDS-minted DOI.""" def check(self, stream_entry, *, record=None, record_pid=None, matcher=None): - """Return an error if the entry DOI uses the CDS DataCite prefix.""" + """Return an error if the entry DOI uses the CDS DataCite prefix. + + On sandbox we allow create with a CDS DOI when the record is missing + locally (it already exists on prod). Everywhere else, block create so + we update the existing record instead of minting a duplicate. + """ + # Same gate as the writer remint path: flag on + sandbox only. + if ( + current_app.config["CDS_HARVESTER_SANDBOX_ALLOW_CREATE_PROD_MISSING_RECORDS"] + and current_app.config.get("CDS_ENVIRONMENT_NAME") == "sandbox" + ): + return None doi = stream_entry.entry.get("pids", {}).get("doi", {}) prefix = current_app.config["DATACITE_PREFIX"] if prefix not in doi.get("identifier", ""): diff --git a/site/cds_rdm/inspire_harvester/utils.py b/site/cds_rdm/inspire_harvester/utils.py index 47e42e05..8c0e5f78 100644 --- a/site/cds_rdm/inspire_harvester/utils.py +++ b/site/cds_rdm/inspire_harvester/utils.py @@ -9,12 +9,48 @@ from collections import Counter +import requests +from flask import current_app from invenio_access.permissions import system_identity from invenio_records_resources.proxies import current_service_registry +from invenio_vocabularies.datastreams.errors import WriterError from opensearchpy import RequestError from sqlalchemy.exc import NoResultFound +def fetch_prod_parent_and_version(cdsrdm_id): + """Ask production what parent and version ids belong to this CDSRDM value. + + INSPIRE's CDSRDM can be either the parent id or a version id. They look + the same (xxxxx-xxxxx), so we cannot tell from the string alone. Production's + record API always returns both: parent.id and id (the version). + """ + base = current_app.config["CDS_HARVESTER_PROD_API_URL"].rstrip("/") + url = f"{base}/api/records/{cdsrdm_id}" + try: + response = requests.get( + url, headers={"Accept": "application/json"}, timeout=60 + ) + response.raise_for_status() + except requests.RequestException as exc: + raise WriterError( + "Could not fetch the production CDS record needed to remint " + "parent and version ids on sandbox. " + f"| details: cdsrdm_id={cdsrdm_id}, error={exc}" + ) from exc + data = response.json() + parent_id = data.get("parent", {}).get("id") + version_id = data.get("id") + if not parent_id or not version_id: + raise WriterError( + "Production CDS returned a record without parent.id or id, " + "so sandbox cannot remint the same parent and version ids. " + f"| details: cdsrdm_id={cdsrdm_id}, " + f"parent_id={parent_id}, version_id={version_id}" + ) + return parent_id, version_id + + def retrieve_identifiers(identifiers, scheme): """Yield identifier values for the given scheme.""" for ident in identifiers or []: diff --git a/site/cds_rdm/inspire_harvester/writer.py b/site/cds_rdm/inspire_harvester/writer.py index 5ae39a91..b2eb6042 100644 --- a/site/cds_rdm/inspire_harvester/writer.py +++ b/site/cds_rdm/inspire_harvester/writer.py @@ -36,7 +36,11 @@ UpdateEngine, UpdateEngineConflict, ) -from cds_rdm.inspire_harvester.utils import compare_metadata +from cds_rdm.inspire_harvester.utils import ( + compare_metadata, + fetch_prod_parent_and_version, +) +from cds_rdm.schemes import cds_rdm_regexp from cds_rdm.utils import compact_text @@ -312,12 +316,28 @@ def _create_record( stream_entry.errors.append(f"[inspire_id={inspire_id}] {msg}") return False entry = {k: v for k, v in stream_entry.entry.items() if k != "_inspire_ctx"} + ctx = stream_entry.entry["_inspire_ctx"] file_entries = entry["files"].get("entries") or {} logger.debug(f"Files to create: {len(file_entries)}") logger.debug("Creating new record draft") - draft = self.drafts.create(entry) + # Sandbox only: record exists on prod but not here. Look up both + # parent and version ids on prod and create with those same ids. + if ( + current_app.config["CDS_HARVESTER_SANDBOX_ALLOW_CREATE_PROD_MISSING_RECORDS"] + and current_app.config.get("CDS_ENVIRONMENT_NAME") == "sandbox" + ): + cds_id = ctx.get("cds_id") + if cds_id and cds_rdm_regexp.fullmatch(str(cds_id)): + parent_pid, record_pid = fetch_prod_parent_and_version(cds_id) + draft = self.drafts.create_reusing_prod_pids( + entry, parent_pid, record_pid + ) + else: + draft = self.drafts.create(entry) + else: + draft = self.drafts.create(entry) logger.info(f"New draft is created ({draft.id}).") try: @@ -334,6 +354,6 @@ def _create_record( logger.error(f"Draft {draft.id} is deleted due to errors.") raise - # add_community succeeded — publish without file sync (files already uploaded above) + # Community is on the draft; publish without syncing files again. self.drafts.publish(draft.id, logger) return True diff --git a/site/tests/conftest.py b/site/tests/conftest.py index 191c89c7..529b3cb5 100644 --- a/site/tests/conftest.py +++ b/site/tests/conftest.py @@ -364,6 +364,7 @@ def app_config(app_config, mock_datacite_client, mock_crossref_client): app_config["RDM_RECORD_CLS"] = CDSRDMRecord app_config["RDM_DRAFT_CLS"] = CDSRDMDraft app_config["CDS_HARVESTER_USER_EMAIL"] = "cds-harvester@cern.ch" + app_config["CDS_HARVESTER_SANDBOX_ALLOW_CREATE_PROD_MISSING_RECORDS"] = False return app_config diff --git a/site/tests/inspire_harvester/test_update_create_CDS_records.py b/site/tests/inspire_harvester/test_update_create_CDS_records.py index ba096546..eeb39266 100644 --- a/site/tests/inspire_harvester/test_update_create_CDS_records.py +++ b/site/tests/inspire_harvester/test_update_create_CDS_records.py @@ -2,6 +2,7 @@ from functools import partial from io import BytesIO from pathlib import Path +from unittest.mock import patch from invenio_access.permissions import system_identity from invenio_rdm_records.proxies import current_rdm_records_service @@ -39,7 +40,10 @@ def test_CDS_DOI_create_record_fails( RDMRecord.index.refresh() doi_filters = [ - dsl.Q("term", **{"pids.doi": "10.17181/CERN.LELX.5VJY"}), + dsl.Q( + "term", + **{"pids.doi.identifier.keyword": "10.17181/CERN.LELX.5VJY"}, + ), ] filter = dsl.Q("bool", filter=doi_filters) @@ -49,6 +53,62 @@ def test_CDS_DOI_create_record_fails( assert created_records.total == 0 +def test_CDS_DOI_create_record_allowed( + running_app, location, scientific_community, datastream_config +): + """Sandbox can create a missing prod record and keep the same ids. + + INSPIRE may send a version id as CDSRDM. We mock the prod lookup so the + create path remints both parent and version to match prod. + """ + running_app.app.config["CDS_HARVESTER_SANDBOX_ALLOW_CREATE_PROD_MISSING_RECORDS"] = True + running_app.app.config["CDS_ENVIRONMENT_NAME"] = "sandbox" + cdsrdm_id = "versi-onnnn" + parent_id = "ab12c-de34f" + with open( + DATA_DIR / "record_with_cds_DOI.json", + "r", + ) as f: + new_record = json.load(f) + hit_meta = new_record["hits"]["hits"][0]["metadata"] + hit_meta["external_system_identifiers"] = [ + ident + for ident in hit_meta.get("external_system_identifiers", []) + if ident.get("schema") != "CDS" + ] + hit_meta["external_system_identifiers"].append( + {"schema": "CDSRDM", "value": cdsrdm_id} + ) + + mock_record = partial(mock_requests_get, mock_content=new_record) + with patch( + "cds_rdm.inspire_harvester.writer.fetch_prod_parent_and_version", + return_value=(parent_id, cdsrdm_id), + ): + run_harvester_mock(datastream_config, mock_record) + RDMRecord.index.refresh() + + doi_filters = [ + dsl.Q( + "term", + **{"pids.doi.identifier.keyword": "10.17181/CERN.LELX.5VJY"}, + ), + ] + filter = dsl.Q("bool", filter=doi_filters) + + created_records = current_rdm_records_service.search( + system_identity, extra_filter=filter + ) + assert created_records.total == 1 + created = created_records.to_dict()["hits"]["hits"][0] + assert created["pids"]["doi"]["identifier"] == "10.17181/CERN.LELX.5VJY" + assert created["pids"]["doi"]["provider"] == "datacite" + assert created["parent"]["id"] == parent_id + assert created["id"] == cdsrdm_id + + running_app.app.config["CDS_HARVESTER_SANDBOX_ALLOW_CREATE_PROD_MISSING_RECORDS"] = False + + def test_update_record_with_CDS_DOI_one_doc_type( running_app, location,