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
146 changes: 91 additions & 55 deletions cds_migrator_kit/rdm/records/load/load.py
Original file line number Diff line number Diff line change
Expand Up @@ -32,9 +32,9 @@
from invenio_records_resources.services.uow import RecordCommitOp
from invenio_requests.customizations.event_types import (
LogEventType,
ReviewersUpdatedType,
)
from invenio_requests.proxies import current_events_service, current_requests_service
from invenio_requests.records.models import RequestMetadata
from marshmallow import ValidationError
from psycopg2.errors import UniqueViolation
from sqlalchemy.exc import IntegrityError
Expand All @@ -46,6 +46,52 @@
)


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 UnexpectedValue: if no matching user is found.
"""
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()

return user


def import_legacy_files(filepath):
"""Download file from legacy."""
if current_app.config["CDS_MIGRATOR_KIT_ENV"] == "local":
Expand Down Expand Up @@ -406,15 +452,21 @@ def _after_publish_mint_recid(self, record, entry, version):
# but then we get a double redirection
legacy_recid_minter(legacy_recid, record._record.parent.model.id)

def _after_publish_add_inclusion_request(self, request_data, record, entry):
def _after_publish_add_inclusion_request(self, request_data, record, entry, uow):
"""Create community inclusion request after publish."""
legacy_recid = entry["record"]["recid"]
request_number = f"lrecid:{legacy_recid}"

# Defensive/idempotency guard: skip if a request for this record was
# already committed by a previous (partial) load attempt, otherwise
# re-creating it would violate the unique constraint on `number`.
if RequestMetadata.query.filter_by(number=request_number).first() is not None:
return

status = request_data.get("status", "accepted")
reviewer_names = request_data.get("reviewers", [])
reviewers = User.query.filter(User._displayname.in_(reviewer_names)).all()
# TODO: check if reviewers are missing and log a warning
reviewers = [find_reviewer(name) for name in reviewer_names]

legacy_recid = entry["record"]["recid"]
created_at = datetime.datetime.fromisoformat(record["created"])

parent = record._record.parent
Expand All @@ -438,59 +490,43 @@ def create_event(request_model, payload, event_type, user):

uow.register(RecordCommitOp(event, indexer=current_events_service.indexer))

with UnitOfWork() as uow:
request_item = current_requests_service.create(
system_identity,
data={"title": record["metadata"]["title"]},
request_type=CommunityInclusion,
receiver=receiver,
creator=creator,
topic={"record": record.id},
uow=uow,
)

request = request_item._record
request.status = "submitted"
request.number = f"lrecid:{legacy_recid}"
request.model.created = created_at

if reviewers:
reviewers_payload = [{"user": str(r.id)} for r in reviewers]
request.reviewers = reviewers_payload

for reviewer in reviewers:
create_event(
request.model,
{
"event": "reviewers_updated",
"content": _("added a reviewer"),
"reviewers": [{"user": str(reviewer.id)}],
},
event_type=ReviewersUpdatedType,
user="system",
)

if status:
request.status = status
if status == "accepted":
parent_to_request_relation = (
parent.communities._m2m_model_cls.query.filter_by(
record_id=parent.id, community_id=community.id
).one()
)
parent_to_request_relation.request_id = request.id
request_item = current_requests_service.create(
system_identity,
data={"title": record["metadata"]["title"]},
request_type=CommunityInclusion,
receiver=receiver,
creator=creator,
topic={"record": record.id},
uow=uow,
)

create_event(
request.model,
{"event": status},
event_type=LogEventType,
user="system",
request = request_item._record
request.status = "submitted"
request.number = request_number
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

if status:
request.status = status
if status == "accepted":
parent_to_request_relation = (
parent.communities._m2m_model_cls.query.filter_by(
record_id=parent.id, community_id=community.id
).one()
)
parent_to_request_relation.request_id = request.id

uow.register(
RecordCommitOp(request, indexer=current_requests_service.indexer)
create_event(
request.model,
{"event": status},
event_type=LogEventType,
user="system",
)
uow.commit()

uow.register(RecordCommitOp(request, indexer=current_requests_service.indexer))

def _after_publish_update_files_created(self, record, entry, version):
"""Update the created date of the files post publish."""
Expand Down Expand Up @@ -520,7 +556,7 @@ def _after_publish(self, identity, published_record, entry, version, uow):

if self.create_inclusion_request and request_data:
self._after_publish_add_inclusion_request(
request_data, published_record, entry
request_data, published_record, entry, uow
)
# db.session.commit()

Expand Down
20 changes: 13 additions & 7 deletions cds_migrator_kit/rdm/records/transform/transform.py
Original file line number Diff line number Diff line change
Expand Up @@ -654,17 +654,23 @@ def field_journal(record_json):
journal = record_json.get("custom_fields", {}).get("journal:journal", {})
if journal:
if not journal.get("title"):
raise UnexpectedValue(
subfield="a",
value=journal,
field="journal",
message="Title is missing in journal field",
stage="vocabulary match",
raise RecordFlaggedCuration(
message="found partial journal field, to be checked",
stage="transform",
field="773",
)
return journal
return {}

_cf = json_entry.get("custom_fields", {})
try:
journal = field_journal(json_entry)
except RecordFlaggedCuration as e:
self.migration_logger.add_information(
json_entry["recid"],
{"message": e.message, "value": e.value},
)
journal = {}
custom_fields = {
"cern:experiments": [],
"cern:departments": [],
Expand All @@ -678,7 +684,7 @@ def field_journal(record_json):
"cern:committees": _cf.get("cern:committees"),
"cern:oa_funding_model": _cf.get("cern:oa_funding_model"),
"thesis:thesis": _cf.get("thesis:thesis", {}),
"journal:journal": field_journal(json_entry),
"journal:journal": journal,
"imprint:imprint": _cf.get("imprint:imprint", {}),
"meeting:meeting": _cf.get("meeting:meeting", {}),
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -49,7 +49,9 @@ def _sub(v, code):
def isbn(self, key, value):
_custom_fields = self.get("custom_fields", {})
_isbn = StringValue(value.get("a", "")).parse()
_isbn_u = StringValue(value.get("u", "")).parse()

_isbn = _isbn or _isbn_u
if _isbn:
try:
_isbn = normalize_isbn(_isbn)
Expand Down Expand Up @@ -455,13 +457,17 @@ def organisation(self, key, value):
@for_each_value
def request_reviewers(self, key, value):
name = StringValue(value.get("p", "")).parse().strip()
email = StringValue(value.get("m", "")).parse().strip()

if "," in name:
last, first = (part.strip() for part in name.split(",", 1))
if email:
reviewer = email
else:
last, first = name, ""
if "," in name:
last, first = (part.strip() for part in name.split(",", 1))
else:
last, first = name, ""

reviewer = " ".join(part for part in (first, last) if part)
reviewer = " ".join(part for part in (first, last) if part)

if reviewer:
request_data = self.setdefault("request_data", {})
Expand Down
5 changes: 5 additions & 0 deletions cds_migrator_kit/rdm/streams.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ records:
tmp_dir: cds_migrator_kit/rdm/tmp/lep_exp/aleph
log_dir: cds_migrator_kit/rdm/log/lep_exp/aleph
plots: true
create_inclusion_request: true
extract:
dirpath: cds_migrator_kit/rdm/data/lep_exp/aleph/dump/
transform:
Expand Down Expand Up @@ -36,6 +37,7 @@ records:
tmp_dir: cds_migrator_kit/rdm/tmp/lep_exp/l3
log_dir: cds_migrator_kit/rdm/log/lep_exp/l3
plots: true
create_inclusion_request: true
extract:
dirpath: cds_migrator_kit/rdm/data/lep_exp/l3/dump/
transform:
Expand All @@ -51,6 +53,7 @@ records:
tmp_dir: cds_migrator_kit/rdm/tmp/lep_exp/opal
log_dir: cds_migrator_kit/rdm/log/lep_exp/opal
plots: true
create_inclusion_request: true
extract:
dirpath: cds_migrator_kit/rdm/data/lep_exp/opal/dump/
transform:
Expand All @@ -66,6 +69,7 @@ records:
tmp_dir: cds_migrator_kit/rdm/tmp/lep_exp/delphi
log_dir: cds_migrator_kit/rdm/log/lep_exp/delphi
plots: true
create_inclusion_request: true
extract:
dirpath: cds_migrator_kit/rdm/data/lep_exp/delphi/dump/
transform:
Expand All @@ -82,6 +86,7 @@ records:
log_dir: cds_migrator_kit/rdm/log/lep_exp/delphi_priv
restricted: "True"
plots: true
create_inclusion_request: true
extract:
dirpath: cds_migrator_kit/rdm/data/lep_exp/delphi_priv/dump/
transform:
Expand Down
46 changes: 46 additions & 0 deletions cds_migrator_kit/rdm/users/log.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,46 @@
# -*- coding: utf-8 -*-
#
# Copyright (C) 2026 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.

"""CDS-RDM submitter/reviewer accounts migration logger module."""

import logging


class SubmitterLogger:
"""Migrator submitter/reviewer accounts logger."""

@classmethod
def initialize(cls, log_dir):
"""Attach file and stream handlers to the submitter-migrator logger."""
formatter = logging.Formatter(
fmt="%(asctime)s %(levelname)-8s %(message)s", datefmt="%Y-%m-%d %H:%M:%S"
)
logger = logging.getLogger("submitter-migrator")
logger.setLevel(logging.INFO)

# info and above to file
fh = logging.FileHandler(log_dir / "info.log")
fh.setLevel(logging.INFO)
fh.setFormatter(formatter)
logger.addHandler(fh)

# errors to their own file
fh = logging.FileHandler(log_dir / "error.log")
fh.setLevel(logging.ERROR)
fh.setFormatter(formatter)
logger.addHandler(fh)

# info to stream/stdout
sh = logging.StreamHandler()
sh.setFormatter(formatter)
sh.setLevel(logging.INFO)
logger.addHandler(sh)

@classmethod
def get_logger(cls):
"""Get migration logger."""
return logging.getLogger("submitter-migrator")
5 changes: 5 additions & 0 deletions cds_migrator_kit/rdm/users/runner.py
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@

from cds_migrator_kit.rdm.affiliations.log import AffiliationsLogger
from cds_migrator_kit.rdm.users.api import CDSMigrationUserAPI
from cds_migrator_kit.rdm.users.log import SubmitterLogger

from .transform import people_marc21, users_migrator_marc21

Expand Down Expand Up @@ -55,6 +56,9 @@ def __init__(self, stream_definition, dirpath, missing_users_dir, log_dir, dry_r
"""Constructor."""
self.log_dir = Path(log_dir)
self.log_dir.mkdir(parents=True, exist_ok=True)

SubmitterLogger.initialize(self.log_dir)

self.stream = Stream(
stream_definition.name,
extract=stream_definition.extract_cls(dirpath),
Expand All @@ -65,6 +69,7 @@ def __init__(self, stream_definition, dirpath, missing_users_dir, log_dir, dry_r
dry_run=dry_run,
missing_users_dir=missing_users_dir,
user_api_cls=CDSMigrationUserAPI,
logger=SubmitterLogger.get_logger(),
),
)

Expand Down
Loading
Loading