Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 12 additions & 0 deletions invenio.cfg
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
12 changes: 12 additions & 0 deletions site/cds_rdm/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -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."""
53 changes: 53 additions & 0 deletions site/cds_rdm/inspire_harvester/load/draft.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,17 +9,46 @@

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,
raise_unexpected_operation_error,
)


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."""

Expand All @@ -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)
Expand Down
13 changes: 12 additions & 1 deletion site/cds_rdm/inspire_harvester/load/validator.py
Original file line number Diff line number Diff line change
Expand Up @@ -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", ""):
Expand Down
36 changes: 36 additions & 0 deletions site/cds_rdm/inspire_harvester/utils.py
Original file line number Diff line number Diff line change
Expand Up @@ -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}"

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

instead of using the API from production and put potentially extra load on our production system, shall we use the producion database directly? We have a read_only user so that is an extra measure to isolate the operation, wdyt @kpsherva ?

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

discussed IRL to go with the API approach and reevaluate.

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 []:
Expand Down
26 changes: 23 additions & 3 deletions site/cds_rdm/inspire_harvester/writer.py
Original file line number Diff line number Diff line change
Expand Up @@ -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


Expand Down Expand Up @@ -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:
Expand All @@ -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
1 change: 1 addition & 0 deletions site/tests/conftest.py
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
62 changes: 61 additions & 1 deletion site/tests/inspire_harvester/test_update_create_CDS_records.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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)

Expand All @@ -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,
Expand Down
Loading