From 472b6717021ee4663537719d12911bde87a547fc Mon Sep 17 00:00:00 2001 From: Saksham Date: Wed, 2 Sep 2026 17:12:49 +0200 Subject: [PATCH 1/5] add(scripts): Add script to rebuild rdm_records_state.json --- scripts/rebuild_records_state.py | 299 +++++++++++++++++++++++++++++++ 1 file changed, 299 insertions(+) create mode 100644 scripts/rebuild_records_state.py diff --git a/scripts/rebuild_records_state.py b/scripts/rebuild_records_state.py new file mode 100644 index 00000000..57ec473c --- /dev/null +++ b/scripts/rebuild_records_state.py @@ -0,0 +1,299 @@ +"""Regenerate ``rdm_records_state.json`` straight from the CDS-RDM database. + +Root cause - previously, the code had a bug: + + for version in versions.keys(): + draft = self._pre_publish(...) + published_record = current_rdm_records_service.publish(...) + self._after_publish(...) + records.append(published_record._record) # left outside the loop + +It ran exactly once per record, for only the latest version. Now fixed in code already. + +``invenio migration stats run`` uses only the ``versions[]`` in that state +file to attribute downloads, so any multi-version record migrated in that +window lost every download event belonging to any version except the latest. + +This script rebuilds the state file from scratch, entirely from CDS-RDM, +one stream/collection at a time - the same collection name used everywhere +else (``invenio migration run --collection ``, and the folder name +under ``cds_migrator_kit/rdm/data/`` and ``/``): + +1. Resolve the collection's community id(s) from ``streams.yaml`` / + ``streams_done.yaml`` / ``streams_shelved.yaml`` (same as + cds_migrator_kit/runner/runner.py). +2. Resolve those community ids to every parent record in them using ``RDMParentCommunity``. +3. Resolve each parent to its legacy recid via the ``lrecid`` PID minted at migration time. +4. Fetch every published version for that parent. +5. Rebuild the state entry field-for-field like ``_load_record_state`` does, + reading current file metadata straight off each version. + +Output is written next to the original, as ``/rdm_records_state.fixed.json``. + +On restart the script reads its own log file for DONE lines and skips any +legacy recid already completed. + +Usage: + invenio shell scripts/rebuild_records_state.py + + main(collection="it", dry_run=True) + main(collection="it", dry_run=False) +""" + +import json +import traceback +from pathlib import Path + +import yaml +from flask import current_app +from invenio_pidstore.models import PersistentIdentifier +from invenio_rdm_records.records.api import RDMRecord +from invenio_rdm_records.records.models import RDMParentCommunity, RDMRecordMetadata + +STREAM_CONFIG_FILES = ( + "cds_migrator_kit/rdm/streams_done.yaml", + "cds_migrator_kit/rdm/streams.yaml", + "cds_migrator_kit/rdm/streams_shelved.yaml", +) + +log_fp = None + + +def log(msg): + print(msg) + if log_fp is not None: + log_fp.write(msg + "\n") + log_fp.flush() + + +def load_completed_recids(log_path): + """Return the set of legacy recids already marked DONE in a previous run.""" + completed = set() + path = Path(log_path) + if not path.exists(): + return completed + with open(log_path, "r") as f: + for line in f: + if line.startswith("DONE: legacy_recid="): + completed.add(line.strip().split("=")[1]) + return completed + + +def mark_done(legacy_recid): + log(f"DONE: legacy_recid={legacy_recid}") + + +def load_communities_ids(collection, config_files=STREAM_CONFIG_FILES): + """Look up a collection's ``transform.communities_ids`` in streams.yaml. + + Reads the same shape ``Runner._read_config()`` does + (cds_migrator_kit/runner/runner.py:28-31): a top-level ``records`` key, + one entry per collection name. + """ + for path in config_files: + if not Path(path).exists(): + continue + with open(path) as f: + config = yaml.safe_load(f) or {} + collection_config = config.get("records", {}).get(collection) + if collection_config: + return collection_config["transform"]["communities_ids"] + raise ValueError( + f"collection {collection!r} not found in any of {config_files} - " + "pass its community id(s) directly via community_ids= instead" + ) + + +def default_output_path(collection): + """Same folder ``RecordStateLogger`` writes to (cds_migrator_kit/reports/log.py), + a different filename - this never overwrites ``rdm_records_state.json``.""" + base_path = current_app.config["CDS_MIGRATOR_KIT_LOGS_PATH"] + return str(Path(base_path) / collection / "rdm_records_state.fixed.json") + + +def find_legacy_recids(community_ids): + """Return {legacy_recid: parent_object_uuid} for every migrated record in + any of the given communities. + + Two plain queries, nothing read from disk: community ids -> parent + uuids (``RDMParentCommunity``), then parent uuids -> legacy recid via + the ``lrecid`` PID (mirrors the lookup + ``CDSMigrationEntryLoad._have_migrated_recid`` does one recid at a time, + cds_migrator_kit/rdm/records/load/load.py). + """ + parent_uuids = [ + row.record_id + for row in RDMParentCommunity.query.filter( + RDMParentCommunity.community_id.in_(community_ids) + ) + ] + pids = PersistentIdentifier.query.filter( + PersistentIdentifier.pid_type == "lrecid", + PersistentIdentifier.object_uuid.in_(parent_uuids), + ) + return {pid.pid_value: str(pid.object_uuid) for pid in pids} + + +def get_rdm_versions(parent_object_uuid): + """Return {version_index: RDMRecord} for every published version of a parent. + + Soft-deleted records are excluded by ``RDMRecord.get_record()`` default. + """ + version_models = ( + RDMRecordMetadata.query.filter_by(parent_id=parent_object_uuid) + .order_by(RDMRecordMetadata.created) + .all() + ) + records = {} + for m in version_models: + try: + rdm_record = RDMRecord.get_record(str(m.id)) + records[rdm_record.versions.index] = rdm_record + except Exception as exc: + log(f"could not load record {m.id}: {exc}") + return records + + +def convert_file_format(file_entries, bucket_id): + """Mirror ``RecordLoad._load_record_state.convert_file_format``.""" + return [ + { + "legacy_file_id": entry["metadata"]["legacy_file_id"], + "bucket_id": bucket_id, + "file_key": entry["key"], + "file_id": entry["file_id"], + "size": str(entry["size"]), + } + for entry in file_entries.values() + ] + + +def extract_record_version(record): + """Mirror ``RecordLoad._load_record_state.extract_record_version``.""" + bucket_id = str(record.files.bucket_id) + files = record.__class__.files.dump( + record, record.files, include_entries=True + ).get("entries", {}) + return { + "new_recid": record.pid.pid_value, + "version": record.versions.index, + "files": convert_file_format(files, bucket_id), + } + + +def build_state_entry(legacy_recid, rdm_versions): + """Rebuild one ``rdm_records_state.json`` entry from live RDM versions. + + Mirrors ``RecordLoad._load_record_state`` + (cds_migrator_kit/rdm/records/load/entities/record.py) field-for-field, + reading current per-version file metadata straight off each version + instead of relying on the state the buggy loop bookkeeping produced. + """ + recid_state = {"legacy_recid": str(legacy_recid), "versions": []} + parent_recid = None + + for version_index in sorted(rdm_versions): + record = rdm_versions[version_index] + + if parent_recid is None: + parent_recid = record.parent.pid.pid_value + recid_state["parent_recid"] = parent_recid + recid_state["parent_object_uuid"] = str(record.parent.id) + + recid_state["versions"].append(extract_record_version(record)) + + if "latest_version" not in recid_state: + latest = record.get_latest_by_parent(record.parent) + recid_state["latest_version"] = latest["id"] + recid_state["latest_version_object_uuid"] = str(latest.id) + + return recid_state + + +def write_state_file(filepath, entries): + """Write entries in the same JSON-list-of-compact-objects format + ``RecordStateLogger.finalise()`` uses (cds_migrator_kit/reports/log.py), + which is what ``invenio migration stats run --filepath`` expects.""" + Path(filepath).parent.mkdir(parents=True, exist_ok=True) + with open(filepath, "w", encoding="utf-8") as f: + f.write("[\n") + for i, entry in enumerate(entries): + json_str = json.dumps(entry, ensure_ascii=False, separators=(",", ":")) + comma = "," if i < len(entries) - 1 else "" + f.write(f"{json_str}{comma}\n") + f.write("]") + + +def main(collection, dry_run=True, output_path=None): + """Rebuild ``rdm_records_state.json`` for one stream/collection. + + :param collection: the stream name, as used by + ``invenio migration run --collection`` - looked up in streams.yaml + to find the community id(s) and to build the default output path. + :param output_path: defaults to + ``//rdm_records_state.fixed.json``, + next to (never over) the original state file. + """ + global log_fp + + community_ids = load_communities_ids(collection) + output_path = output_path or default_output_path(collection) + log_file = f"{output_path}.log" + + completed_recids = load_completed_recids(log_file) + log_fp = open(log_file, "a") + + if completed_recids: + log(f"resuming — {len(completed_recids)} recid(s) already completed, skipping them") + + try: + legacy_recids = find_legacy_recids(community_ids) + log( + f"starting [collection={collection}, community_ids={community_ids}, " + f"dry_run={dry_run}, legacy_recids={len(legacy_recids)}]" + ) + + stats = {"checked": 0, "skipped_done": 0, "no_versions": 0, "fixed": 0, "errors": 0} + entries = [] + + for i, (legacy_recid, parent_object_uuid) in enumerate(legacy_recids.items(), start=1): + stats["checked"] += 1 + + if legacy_recid in completed_recids: + stats["skipped_done"] += 1 + continue + + log(f"[{i}/{len(legacy_recids)}] legacy_recid={legacy_recid}") + + try: + rdm_versions = get_rdm_versions(parent_object_uuid) + if not rdm_versions: + log(f"legacy_recid={legacy_recid} - no published versions found, skipping") + stats["no_versions"] += 1 + continue + + entries.append(build_state_entry(legacy_recid, rdm_versions)) + stats["fixed"] += 1 + + if not dry_run: + mark_done(legacy_recid) + + except Exception as exc: + log(f"unexpected error for legacy_recid={legacy_recid}: {exc}") + log(traceback.format_exc()) + stats["errors"] += 1 + + if not dry_run: + write_state_file(output_path, entries) + log(f"wrote {output_path}") + else: + log(f"dry run — would write {output_path}") + + log( + f"\nsummary:\nchecked={stats['checked']}\nalready_done={stats['skipped_done']}\n" + f"no_versions={stats['no_versions']}\nfixed={stats['fixed']}\n" + f"errors={stats['errors']}" + ) + finally: + log_fp.close() + log_fp = None From 64709dcadf526091d2f2cc41f0341b0565e51d68 Mon Sep 17 00:00:00 2001 From: Saksham Date: Mon, 21 Sep 2026 10:58:35 +0200 Subject: [PATCH 2/5] refactor(rebuild_record_state): Pass community_ids directly --- scripts/rebuild_records_state.py | 61 ++++++-------------------------- 1 file changed, 10 insertions(+), 51 deletions(-) diff --git a/scripts/rebuild_records_state.py b/scripts/rebuild_records_state.py index 57ec473c..94bdabed 100644 --- a/scripts/rebuild_records_state.py +++ b/scripts/rebuild_records_state.py @@ -50,11 +50,6 @@ from invenio_rdm_records.records.api import RDMRecord from invenio_rdm_records.records.models import RDMParentCommunity, RDMRecordMetadata -STREAM_CONFIG_FILES = ( - "cds_migrator_kit/rdm/streams_done.yaml", - "cds_migrator_kit/rdm/streams.yaml", - "cds_migrator_kit/rdm/streams_shelved.yaml", -) log_fp = None @@ -79,38 +74,6 @@ def load_completed_recids(log_path): return completed -def mark_done(legacy_recid): - log(f"DONE: legacy_recid={legacy_recid}") - - -def load_communities_ids(collection, config_files=STREAM_CONFIG_FILES): - """Look up a collection's ``transform.communities_ids`` in streams.yaml. - - Reads the same shape ``Runner._read_config()`` does - (cds_migrator_kit/runner/runner.py:28-31): a top-level ``records`` key, - one entry per collection name. - """ - for path in config_files: - if not Path(path).exists(): - continue - with open(path) as f: - config = yaml.safe_load(f) or {} - collection_config = config.get("records", {}).get(collection) - if collection_config: - return collection_config["transform"]["communities_ids"] - raise ValueError( - f"collection {collection!r} not found in any of {config_files} - " - "pass its community id(s) directly via community_ids= instead" - ) - - -def default_output_path(collection): - """Same folder ``RecordStateLogger`` writes to (cds_migrator_kit/reports/log.py), - a different filename - this never overwrites ``rdm_records_state.json``.""" - base_path = current_app.config["CDS_MIGRATOR_KIT_LOGS_PATH"] - return str(Path(base_path) / collection / "rdm_records_state.fixed.json") - - def find_legacy_recids(community_ids): """Return {legacy_recid: parent_object_uuid} for every migrated record in any of the given communities. @@ -224,22 +187,18 @@ def write_state_file(filepath, entries): f.write("]") -def main(collection, dry_run=True, output_path=None): +def main(community_ids, output_path, log_file, dry_run=True): """Rebuild ``rdm_records_state.json`` for one stream/collection. - :param collection: the stream name, as used by - ``invenio migration run --collection`` - looked up in streams.yaml - to find the community id(s) and to build the default output path. - :param output_path: defaults to - ``//rdm_records_state.fixed.json``, - next to (never over) the original state file. + :param community_ids: the community UUIDs passed inside the streams.yaml, + as used by ``invenio migration run --collection``. + :param output_path: should be `//rdm_records_state.fixed.json` + next to the original state file (no overwriting). + :param log_file: Log file. + :param dry_run: Pass dfry_run False to write generate state json data to file. """ global log_fp - community_ids = load_communities_ids(collection) - output_path = output_path or default_output_path(collection) - log_file = f"{output_path}.log" - completed_recids = load_completed_recids(log_file) log_fp = open(log_file, "a") @@ -249,8 +208,8 @@ def main(collection, dry_run=True, output_path=None): try: legacy_recids = find_legacy_recids(community_ids) log( - f"starting [collection={collection}, community_ids={community_ids}, " - f"dry_run={dry_run}, legacy_recids={len(legacy_recids)}]" + f"starting, community_ids={str(community_ids)}, " + f"dry_run={dry_run}, found legacy_recids={len(legacy_recids)}]" ) stats = {"checked": 0, "skipped_done": 0, "no_versions": 0, "fixed": 0, "errors": 0} @@ -276,7 +235,7 @@ def main(collection, dry_run=True, output_path=None): stats["fixed"] += 1 if not dry_run: - mark_done(legacy_recid) + log(f"DONE: legacy_recid={legacy_recid}") except Exception as exc: log(f"unexpected error for legacy_recid={legacy_recid}: {exc}") From 2e44fe9ac63cfbe5e0b76acaf22b44cae41ca0d8 Mon Sep 17 00:00:00 2001 From: Saksham Date: Fri, 2 Oct 2026 16:37:40 +0200 Subject: [PATCH 3/5] refactor(rebuild_record_state): Write record state in batches to improve perf --- scripts/rebuild_records_state.py | 101 ++++++++++++++++++++++--------- 1 file changed, 72 insertions(+), 29 deletions(-) diff --git a/scripts/rebuild_records_state.py b/scripts/rebuild_records_state.py index 94bdabed..adb01202 100644 --- a/scripts/rebuild_records_state.py +++ b/scripts/rebuild_records_state.py @@ -14,38 +14,44 @@ file to attribute downloads, so any multi-version record migrated in that window lost every download event belonging to any version except the latest. -This script rebuilds the state file from scratch, entirely from CDS-RDM, -one stream/collection at a time - the same collection name used everywhere -else (``invenio migration run --collection ``, and the folder name -under ``cds_migrator_kit/rdm/data/`` and ``/``): - -1. Resolve the collection's community id(s) from ``streams.yaml`` / - ``streams_done.yaml`` / ``streams_shelved.yaml`` (same as - cds_migrator_kit/runner/runner.py). -2. Resolve those community ids to every parent record in them using ``RDMParentCommunity``. -3. Resolve each parent to its legacy recid via the ``lrecid`` PID minted at migration time. -4. Fetch every published version for that parent. -5. Rebuild the state entry field-for-field like ``_load_record_state`` does, - reading current file metadata straight off each version. - -Output is written next to the original, as ``/rdm_records_state.fixed.json``. +This script rebuilds the state file from scratch, entirely from CDS-RDM, for +a given collection's community id(s) (from that collection's +``transform.communities_ids`` in ``streams.yaml`` / ``streams_done.yaml`` / +``streams_shelved.yaml``): + +1. Resolve the given community id(s) to every parent record in them using + ``RDMParentCommunity``. +2. Resolve each parent to its legacy recid via the ``lrecid`` PID minted at + migration time. +3. Rebuild the state entry via the shared + ``cds_migrator_kit.rdm.records.load.entities.record.build_record_state_from_db`` + - the same function ``CDSMigrationEntryLoad._should_skip_recid`` uses to + backfill a single record's state when it's found missing for an + already-migrated record (cds_migrator_kit/rdm/records/load/load.py). + +Output is written next to the original, as ``/rdm_records_state.fixed.json``, +in batches of ``BATCH_SIZE`` entries at a time rather than all at once - a +batch is merged into whatever's already on disk and the complete file is +rewritten, so a crash loses at most one in-progress batch, and the file +never has to be held in memory in full. On restart the script reads its own log file for DONE lines and skips any legacy recid already completed. Usage: - invenio shell scripts/rebuild_records_state.py + invenio shell - main(collection="it", dry_run=True) - main(collection="it", dry_run=False) + main(community_ids=[""], + output_path="/migration/tmp//rdm_records_state.fixed.json", + log_file="/migration/tmp//records_state.log", + dry_run=True) """ import json import traceback from pathlib import Path -import yaml -from flask import current_app +from invenio_db import db from invenio_pidstore.models import PersistentIdentifier from invenio_rdm_records.records.api import RDMRecord from invenio_rdm_records.records.models import RDMParentCommunity, RDMRecordMetadata @@ -104,7 +110,6 @@ def get_rdm_versions(parent_object_uuid): """ version_models = ( RDMRecordMetadata.query.filter_by(parent_id=parent_object_uuid) - .order_by(RDMRecordMetadata.created) .all() ) records = {} @@ -173,6 +178,9 @@ def build_state_entry(legacy_recid, rdm_versions): return recid_state +BATCH_SIZE = 500 + + def write_state_file(filepath, entries): """Write entries in the same JSON-list-of-compact-objects format ``RecordStateLogger.finalise()`` uses (cds_migrator_kit/reports/log.py), @@ -187,15 +195,29 @@ def write_state_file(filepath, entries): f.write("]") +def flush_batch(output_path, batch): + """Merge a batch of new entries into whatever is already on disk and + rewrite the complete file. Called every ``BATCH_SIZE`` entries instead + of once per entry (too slow, file never readable) or once for the whole + run (loses everything on a crash) - entries already flushed by a + previous batch are read back and re-written, not held in memory for the + rest of the run.""" + existing = [] + if Path(output_path).exists(): + with open(output_path, encoding="utf-8") as f: + existing = json.load(f) + write_state_file(output_path, existing + batch) + + def main(community_ids, output_path, log_file, dry_run=True): """Rebuild ``rdm_records_state.json`` for one stream/collection. - :param community_ids: the community UUIDs passed inside the streams.yaml, - as used by ``invenio migration run --collection``. + :param community_ids: the community UUIDs for a collection, as set in + its ``transform.communities_ids`` in streams.yaml. :param output_path: should be `//rdm_records_state.fixed.json` next to the original state file (no overwriting). - :param log_file: Log file. - :param dry_run: Pass dfry_run False to write generate state json data to file. + :param log_file: resumable progress log - DONE lines mark completed recids. + :param dry_run: pass dry_run=False to actually write entries to disk. """ global log_fp @@ -205,6 +227,24 @@ def main(community_ids, output_path, log_file, dry_run=True): if completed_recids: log(f"resuming — {len(completed_recids)} recid(s) already completed, skipping them") + batch = [] + # legacy_recids pending in `batch` - only marked DONE once their batch + # is actually flushed, so a crash before that flush leaves them absent + # from completed_recids on resume, and state generation re-runs for + # them instead of being silently skipped as already done. + batch_recids = [] + + def flush_pending(): + if batch: + flush_batch(output_path, batch) + for r in batch_recids: + log(f"DONE: legacy_recid={r}") + batch.clear() + batch_recids.clear() + # Do not hold the Record objects in DB session once the batch is flushed + # We are only reading from DB so this is safe + db.session.expunge_all() + try: legacy_recids = find_legacy_recids(community_ids) log( @@ -213,7 +253,6 @@ def main(community_ids, output_path, log_file, dry_run=True): ) stats = {"checked": 0, "skipped_done": 0, "no_versions": 0, "fixed": 0, "errors": 0} - entries = [] for i, (legacy_recid, parent_object_uuid) in enumerate(legacy_recids.items(), start=1): stats["checked"] += 1 @@ -231,11 +270,13 @@ def main(community_ids, output_path, log_file, dry_run=True): stats["no_versions"] += 1 continue - entries.append(build_state_entry(legacy_recid, rdm_versions)) + batch.append(build_state_entry(legacy_recid, rdm_versions)) stats["fixed"] += 1 if not dry_run: - log(f"DONE: legacy_recid={legacy_recid}") + batch_recids.append(legacy_recid) + if len(batch) >= BATCH_SIZE: + flush_pending() except Exception as exc: log(f"unexpected error for legacy_recid={legacy_recid}: {exc}") @@ -243,7 +284,7 @@ def main(community_ids, output_path, log_file, dry_run=True): stats["errors"] += 1 if not dry_run: - write_state_file(output_path, entries) + flush_pending() log(f"wrote {output_path}") else: log(f"dry run — would write {output_path}") @@ -254,5 +295,7 @@ def main(community_ids, output_path, log_file, dry_run=True): f"errors={stats['errors']}" ) finally: + if not dry_run: + flush_pending() log_fp.close() log_fp = None From fd8095fce34db0ed8dc806cb85e486b207991655 Mon Sep 17 00:00:00 2001 From: Saksham Date: Mon, 5 Oct 2026 15:48:35 +0200 Subject: [PATCH 4/5] refactor(rebuild_record_state): Batch queries and load in memory before creating state json file --- scripts/rebuild_records_state.py | 427 ++++++++++++++++++------------- 1 file changed, 244 insertions(+), 183 deletions(-) diff --git a/scripts/rebuild_records_state.py b/scripts/rebuild_records_state.py index adb01202..6d9b8839 100644 --- a/scripts/rebuild_records_state.py +++ b/scripts/rebuild_records_state.py @@ -21,22 +21,20 @@ 1. Resolve the given community id(s) to every parent record in them using ``RDMParentCommunity``. -2. Resolve each parent to its legacy recid via the ``lrecid`` PID minted at - migration time. -3. Rebuild the state entry via the shared - ``cds_migrator_kit.rdm.records.load.entities.record.build_record_state_from_db`` - - the same function ``CDSMigrationEntryLoad._should_skip_recid`` uses to - backfill a single record's state when it's found missing for an - already-migrated record (cds_migrator_kit/rdm/records/load/load.py). - -Output is written next to the original, as ``/rdm_records_state.fixed.json``, -in batches of ``BATCH_SIZE`` entries at a time rather than all at once - a -batch is merged into whatever's already on disk and the complete file is -rewritten, so a crash loses at most one in-progress batch, and the file -never has to be held in memory in full. - -On restart the script reads its own log file for DONE lines and skips any -legacy recid already completed. +2. Resolve each parent to its legacy recid and its parent recid in one PID + query (``pid_type`` ``lrecid`` and ``recid``). +3. Load every published version of those parents and the file rows. Each + query is limited to ``CHUNK_SIZE`` ids so one result set is not the whole + community. ``rdm_records_metadata.parent_id`` is not indexed, so this is a + handful of table scans instead of one scan per record. +4. Resolve each record uuid to its recid through the ``recid`` PID. The + version query does not read the record JSON. +5. Build the state list in memory. The latest version is the highest + ``index`` among those rows. File entries are not stored on the record + JSON (``FilesField(store=False)``); they come from ``rdm_records_files``. + +Output is written once, next to the original, as +``/rdm_records_state.fixed.json``. Usage: invenio shell @@ -52,10 +50,13 @@ from pathlib import Path from invenio_db import db -from invenio_pidstore.models import PersistentIdentifier -from invenio_rdm_records.records.api import RDMRecord -from invenio_rdm_records.records.models import RDMParentCommunity, RDMRecordMetadata - +from invenio_files_rest.models import FileInstance, ObjectVersion +from invenio_pidstore.models import PersistentIdentifier, PIDStatus +from invenio_rdm_records.records.models import ( + RDMFileRecordMetadata, + RDMParentCommunity, + RDMRecordMetadata, +) log_fp = None @@ -67,118 +68,210 @@ def log(msg): log_fp.flush() -def load_completed_recids(log_path): - """Return the set of legacy recids already marked DONE in a previous run.""" - completed = set() - path = Path(log_path) - if not path.exists(): - return completed - with open(log_path, "r") as f: - for line in f: - if line.startswith("DONE: legacy_recid="): - completed.add(line.strip().split("=")[1]) - return completed - - -def find_legacy_recids(community_ids): - """Return {legacy_recid: parent_object_uuid} for every migrated record in - any of the given communities. - - Two plain queries, nothing read from disk: community ids -> parent - uuids (``RDMParentCommunity``), then parent uuids -> legacy recid via - the ``lrecid`` PID (mirrors the lookup - ``CDSMigrationEntryLoad._have_migrated_recid`` does one recid at a time, - cds_migrator_kit/rdm/records/load/load.py). - """ - parent_uuids = [ - row.record_id - for row in RDMParentCommunity.query.filter( - RDMParentCommunity.community_id.in_(community_ids) - ) - ] - pids = PersistentIdentifier.query.filter( - PersistentIdentifier.pid_type == "lrecid", - PersistentIdentifier.object_uuid.in_(parent_uuids), - ) - return {pid.pid_value: str(pid.object_uuid) for pid in pids} +CHUNK_SIZE = 10000 + +def chunked(items, size=CHUNK_SIZE): + """Yield ``items`` in slices of ``size``.""" + items = list(items) + for start in range(0, len(items), size): + yield items[start : start + size] -def get_rdm_versions(parent_object_uuid): - """Return {version_index: RDMRecord} for every published version of a parent. - Soft-deleted records are excluded by ``RDMRecord.get_record()`` default. +def find_parent_pids(community_ids): + """Return legacy recids and parent recids for migrated records. + + Community ids map to parent uuids via ``RDMParentCommunity``. One PID + query per chunk then splits ``lrecid`` and ``recid`` for those parents. + + :returns: ``({legacy_recid: parent_uuid}, {parent_uuid: parent_recid})`` """ - version_models = ( - RDMRecordMetadata.query.filter_by(parent_id=parent_object_uuid) - .all() - ) - records = {} - for m in version_models: - try: - rdm_record = RDMRecord.get_record(str(m.id)) - records[rdm_record.versions.index] = rdm_record - except Exception as exc: - log(f"could not load record {m.id}: {exc}") - return records - - -def convert_file_format(file_entries, bucket_id): - """Mirror ``RecordLoad._load_record_state.convert_file_format``.""" - return [ + parent_uuids = list( { - "legacy_file_id": entry["metadata"]["legacy_file_id"], - "bucket_id": bucket_id, - "file_key": entry["key"], - "file_id": entry["file_id"], - "size": str(entry["size"]), + row.record_id + for row in RDMParentCommunity.query.filter( + RDMParentCommunity.community_id.in_(community_ids) + ) } - for entry in file_entries.values() - ] + ) + legacy_recids = {} + parent_recids = {} + for batch in chunked(parent_uuids): + pids = PersistentIdentifier.query.filter( + PersistentIdentifier.pid_type.in_(("lrecid", "recid")), + PersistentIdentifier.object_type == "rec", + PersistentIdentifier.status == PIDStatus.REGISTERED, + PersistentIdentifier.object_uuid.in_(batch), + ) + for pid in pids: + object_uuid = str(pid.object_uuid) + if pid.pid_type == "lrecid": + legacy_recids[pid.pid_value] = object_uuid + else: + parent_recids[object_uuid] = pid.pid_value + db.session.expunge_all() + return legacy_recids, parent_recids -def extract_record_version(record): - """Mirror ``RecordLoad._load_record_state.extract_record_version``.""" - bucket_id = str(record.files.bucket_id) - files = record.__class__.files.dump( - record, record.files, include_entries=True - ).get("entries", {}) - return { - "new_recid": record.pid.pid_value, - "version": record.versions.index, - "files": convert_file_format(files, bucket_id), - } +def load_pids(object_uuids, pid_type): + """Return ``{object_uuid: pid_value}`` for registered PIDs. + ``object_type`` is included so Postgres can use ``idx_object`` + ``(object_type, object_uuid)`` instead of scanning ``pidstore_pid``. + """ + found = {} + for batch in chunked(object_uuids): + pids = PersistentIdentifier.query.filter( + PersistentIdentifier.pid_type == pid_type, + PersistentIdentifier.object_type == "rec", + PersistentIdentifier.status == PIDStatus.REGISTERED, + PersistentIdentifier.object_uuid.in_(batch), + ) + for pid in pids: + found[str(pid.object_uuid)] = pid.pid_value + db.session.expunge_all() + return found -def build_state_entry(legacy_recid, rdm_versions): - """Rebuild one ``rdm_records_state.json`` entry from live RDM versions. - Mirrors ``RecordLoad._load_record_state`` - (cds_migrator_kit/rdm/records/load/entities/record.py) field-for-field, - reading current per-version file metadata straight off each version - instead of relying on the state the buggy loop bookkeeping produced. - """ - recid_state = {"legacy_recid": str(legacy_recid), "versions": []} - parent_recid = None +def load_versions(legacy_recids): + """Return ``({legacy_recid: {version_index: version dict}}, failed_recids)``. - for version_index in sorted(rdm_versions): - record = rdm_versions[version_index] + Version rows are loaded without the JSON column. The recid is the + ``recid`` PID for that row's primary key. Soft-deleted rows (``json`` + is NULL) are excluded, matching ``RDMRecord.get_record()``. + """ + parent_of = {parent: recid for recid, parent in legacy_recids.items()} + versions = {} + parents = list(parent_of) + for n, batch in enumerate(chunked(parents), start=1): + rows = ( + db.session.query( + RDMRecordMetadata.id, + RDMRecordMetadata.parent_id, + RDMRecordMetadata.index, + RDMRecordMetadata.bucket_id, + ) + .filter(RDMRecordMetadata.parent_id.in_(batch)) + .filter(RDMRecordMetadata.json.isnot(None)) + .all() + ) + for rec_id, parent_id, index, bucket_id in rows: + legacy_recid = parent_of[str(parent_id)] + versions.setdefault(legacy_recid, {})[index] = { + "record_uuid": str(rec_id), + "bucket_id": str(bucket_id) if bucket_id else None, + "parent_object_uuid": str(parent_id), + "files": [], + } + log(f"loaded version chunk {n}, legacy_recids with versions={len(versions)}") + db.session.expunge_all() + + record_uuids = [ + version["record_uuid"] + for by_index in versions.values() + for version in by_index.values() + ] + record_recids = load_pids(record_uuids, "recid") + failed = set() + for legacy_recid, by_index in versions.items(): + for version in by_index.values(): + new_recid = record_recids.get(version["record_uuid"]) + if not new_recid: + log( + f"record {version['record_uuid']} has no recid PID, " + f"legacy_recid={legacy_recid}" + ) + failed.add(legacy_recid) + continue + version["new_recid"] = new_recid + return versions, failed - if parent_recid is None: - parent_recid = record.parent.pid.pid_value - recid_state["parent_recid"] = parent_recid - recid_state["parent_object_uuid"] = str(record.parent.id) - recid_state["versions"].append(extract_record_version(record)) +def load_files(record_uuids): + """Return {record_uuid: [file dict, ...]} from ``rdm_records_files``. - if "latest_version" not in recid_state: - latest = record.get_latest_by_parent(record.parent) - recid_state["latest_version"] = latest["id"] - recid_state["latest_version_object_uuid"] = str(latest.id) + ``record_id`` is indexed. ``file_id`` and ``size`` live on the object + version and its file instance, not on the record JSON. + """ + files = {} + uuids = list(record_uuids) + for n, batch in enumerate(chunked(uuids), start=1): + rows = ( + db.session.query( + RDMFileRecordMetadata.record_id, + RDMFileRecordMetadata.key, + RDMFileRecordMetadata.json, + ObjectVersion.file_id, + FileInstance.size, + ) + .join( + ObjectVersion, + ObjectVersion.version_id == RDMFileRecordMetadata.object_version_id, + ) + .join(FileInstance, FileInstance.id == ObjectVersion.file_id) + .filter(RDMFileRecordMetadata.record_id.in_(batch)) + .filter(RDMFileRecordMetadata.json.isnot(None)) + .all() + ) + for record_id, key, file_json, file_id, size in rows: + metadata = (file_json or {}).get("metadata") or {} + files.setdefault(str(record_id), []).append( + { + "legacy_file_id": metadata.get("legacy_file_id"), + "file_key": key, + "file_id": str(file_id), + "size": str(size), + } + ) + del rows + log(f"loaded file chunk {n}, records with files={len(files)}") + db.session.expunge_all() + return files + + +def _file_for_state(bucket_id, file_entry): + """One file entry in the shape ``_load_record_state`` writes.""" + legacy_file_id = file_entry["legacy_file_id"] + if legacy_file_id is None: + raise KeyError("legacy_file_id") + return { + "legacy_file_id": legacy_file_id, + "bucket_id": bucket_id, + "file_key": file_entry["file_key"], + "file_id": file_entry["file_id"], + "size": file_entry["size"], + } - return recid_state +def build_state_entry(legacy_recid, rdm_versions, parent_recid): + """Rebuild one ``rdm_records_state.json`` entry. -BATCH_SIZE = 500 + Same fields as ``RecordLoad._load_record_state`` + (cds_migrator_kit/rdm/records/load/entities/record.py). The latest + version is the highest ``index`` still published. + """ + recid_state = { + "legacy_recid": str(legacy_recid), + "parent_recid": parent_recid, + "versions": [], + } + for version_index in sorted(rdm_versions): + version = rdm_versions[version_index] + if "parent_object_uuid" not in recid_state: + recid_state["parent_object_uuid"] = version["parent_object_uuid"] + recid_state["versions"].append( + { + "new_recid": version["new_recid"], + "version": version_index, + "files": [ + _file_for_state(version["bucket_id"], f) for f in version["files"] + ], + } + ) + latest_version = rdm_versions[max(rdm_versions)] + recid_state["latest_version"] = latest_version["new_recid"] + recid_state["latest_version_object_uuid"] = latest_version["record_uuid"] + return recid_state def write_state_file(filepath, entries): @@ -195,20 +288,6 @@ def write_state_file(filepath, entries): f.write("]") -def flush_batch(output_path, batch): - """Merge a batch of new entries into whatever is already on disk and - rewrite the complete file. Called every ``BATCH_SIZE`` entries instead - of once per entry (too slow, file never readable) or once for the whole - run (loses everything on a crash) - entries already flushed by a - previous batch are read back and re-written, not held in memory for the - rest of the run.""" - existing = [] - if Path(output_path).exists(): - with open(output_path, encoding="utf-8") as f: - existing = json.load(f) - write_state_file(output_path, existing + batch) - - def main(community_ids, output_path, log_file, dry_run=True): """Rebuild ``rdm_records_state.json`` for one stream/collection. @@ -216,86 +295,68 @@ def main(community_ids, output_path, log_file, dry_run=True): its ``transform.communities_ids`` in streams.yaml. :param output_path: should be `//rdm_records_state.fixed.json` next to the original state file (no overwriting). - :param log_file: resumable progress log - DONE lines mark completed recids. + :param log_file: progress log for this run. :param dry_run: pass dry_run=False to actually write entries to disk. """ global log_fp - completed_recids = load_completed_recids(log_file) log_fp = open(log_file, "a") - - if completed_recids: - log(f"resuming — {len(completed_recids)} recid(s) already completed, skipping them") - - batch = [] - # legacy_recids pending in `batch` - only marked DONE once their batch - # is actually flushed, so a crash before that flush leaves them absent - # from completed_recids on resume, and state generation re-runs for - # them instead of being silently skipped as already done. - batch_recids = [] - - def flush_pending(): - if batch: - flush_batch(output_path, batch) - for r in batch_recids: - log(f"DONE: legacy_recid={r}") - batch.clear() - batch_recids.clear() - # Do not hold the Record objects in DB session once the batch is flushed - # We are only reading from DB so this is safe - db.session.expunge_all() - try: - legacy_recids = find_legacy_recids(community_ids) + legacy_recids, parent_recids = find_parent_pids(community_ids) log( f"starting, community_ids={str(community_ids)}, " - f"dry_run={dry_run}, found legacy_recids={len(legacy_recids)}]" + f"dry_run={dry_run}, found legacy_recids={len(legacy_recids)}" ) - - stats = {"checked": 0, "skipped_done": 0, "no_versions": 0, "fixed": 0, "errors": 0} - - for i, (legacy_recid, parent_object_uuid) in enumerate(legacy_recids.items(), start=1): + versions, failed = load_versions(legacy_recids) + record_uuids = [ + version["record_uuid"] + for by_index in versions.values() + for version in by_index.values() + ] + files = load_files(record_uuids) + for by_index in versions.values(): + for version in by_index.values(): + version["files"] = files.get(version["record_uuid"], []) + + stats = {"checked": 0, "no_versions": 0, "fixed": 0, "errors": 0} + entries = [] + for legacy_recid, parent_uuid in legacy_recids.items(): stats["checked"] += 1 - - if legacy_recid in completed_recids: - stats["skipped_done"] += 1 + if legacy_recid in failed: + stats["errors"] += 1 + continue + rdm_versions = versions.get(legacy_recid) + if not rdm_versions: + log( + f"legacy_recid={legacy_recid} - no published versions found, skipping" + ) + stats["no_versions"] += 1 continue - - log(f"[{i}/{len(legacy_recids)}] legacy_recid={legacy_recid}") - try: - rdm_versions = get_rdm_versions(parent_object_uuid) - if not rdm_versions: - log(f"legacy_recid={legacy_recid} - no published versions found, skipping") - stats["no_versions"] += 1 - continue - - batch.append(build_state_entry(legacy_recid, rdm_versions)) + entries.append( + build_state_entry( + legacy_recid, + rdm_versions, + parent_recids[parent_uuid], + ) + ) stats["fixed"] += 1 - - if not dry_run: - batch_recids.append(legacy_recid) - if len(batch) >= BATCH_SIZE: - flush_pending() - except Exception as exc: log(f"unexpected error for legacy_recid={legacy_recid}: {exc}") log(traceback.format_exc()) stats["errors"] += 1 if not dry_run: - flush_pending() - log(f"wrote {output_path}") + write_state_file(output_path, entries) + log(f"wrote {output_path} ({len(entries)} entries)") else: - log(f"dry run — would write {output_path}") + log(f"dry run — would write {output_path} ({len(entries)} entries)") log( - f"\nsummary:\nchecked={stats['checked']}\nalready_done={stats['skipped_done']}\n" + f"\nsummary:\nchecked={stats['checked']}\n" f"no_versions={stats['no_versions']}\nfixed={stats['fixed']}\n" f"errors={stats['errors']}" ) finally: - if not dry_run: - flush_pending() log_fp.close() log_fp = None From fd974ebb2e50987ca3ee717a0c29c417ba909ae8 Mon Sep 17 00:00:00 2001 From: Saksham Date: Mon, 5 Oct 2026 17:04:32 +0200 Subject: [PATCH 5/5] refactor(rebuild_record_state): Handle non-migrated new version edge case --- scripts/rebuild_records_state.py | 92 ++++++++++++++++++++++---------- 1 file changed, 64 insertions(+), 28 deletions(-) diff --git a/scripts/rebuild_records_state.py b/scripts/rebuild_records_state.py index 6d9b8839..17d0bc52 100644 --- a/scripts/rebuild_records_state.py +++ b/scripts/rebuild_records_state.py @@ -28,10 +28,14 @@ community. ``rdm_records_metadata.parent_id`` is not indexed, so this is a handful of table scans instead of one scan per record. 4. Resolve each record uuid to its recid through the ``recid`` PID. The - version query does not read the record JSON. -5. Build the state list in memory. The latest version is the highest - ``index`` among those rows. File entries are not stored on the record - JSON (``FilesField(store=False)``); they come from ``rdm_records_files``. + version query does not read the record JSON. Versions are stored once + per parent, so a redirect ``lrecid`` minted onto that parent gets the + same versions as the migrated record. +5. Build the state list in memory. ``latest_version`` is the highest + ``index`` that still has a migrated file. A record migrated with no + files keeps the highest published index so pageviews have a target. + File entries are not stored on the record JSON + (``FilesField(store=False)``); they come from ``rdm_records_files``. Output is written once, next to the original, as ``/rdm_records_state.fixed.json``. @@ -133,16 +137,18 @@ def load_pids(object_uuids, pid_type): return found -def load_versions(legacy_recids): - """Return ``({legacy_recid: {version_index: version dict}}, failed_recids)``. +def load_versions(parent_uuids): + """Return ``({parent_uuid: {version_index: version dict}}, failed_parents)``. Version rows are loaded without the JSON column. The recid is the ``recid`` PID for that row's primary key. Soft-deleted rows (``json`` is NULL) are excluded, matching ``RDMRecord.get_record()``. + + Keyed by parent so every ``lrecid`` on that parent (the migrated record + and any redirect minted onto it) can share one version list. """ - parent_of = {parent: recid for recid, parent in legacy_recids.items()} versions = {} - parents = list(parent_of) + parents = list(dict.fromkeys(parent_uuids)) for n, batch in enumerate(chunked(parents), start=1): rows = ( db.session.query( @@ -156,14 +162,14 @@ def load_versions(legacy_recids): .all() ) for rec_id, parent_id, index, bucket_id in rows: - legacy_recid = parent_of[str(parent_id)] - versions.setdefault(legacy_recid, {})[index] = { + parent_uuid = str(parent_id) + versions.setdefault(parent_uuid, {})[index] = { "record_uuid": str(rec_id), "bucket_id": str(bucket_id) if bucket_id else None, - "parent_object_uuid": str(parent_id), + "parent_object_uuid": parent_uuid, "files": [], } - log(f"loaded version chunk {n}, legacy_recids with versions={len(versions)}") + log(f"loaded version chunk {n}, parents with versions={len(versions)}") db.session.expunge_all() record_uuids = [ @@ -173,15 +179,15 @@ def load_versions(legacy_recids): ] record_recids = load_pids(record_uuids, "recid") failed = set() - for legacy_recid, by_index in versions.items(): + for parent_uuid, by_index in versions.items(): for version in by_index.values(): new_recid = record_recids.get(version["record_uuid"]) if not new_recid: log( f"record {version['record_uuid']} has no recid PID, " - f"legacy_recid={legacy_recid}" + f"parent_id={parent_uuid}" ) - failed.add(legacy_recid) + failed.add(parent_uuid) continue version["new_recid"] = new_recid return versions, failed @@ -230,10 +236,14 @@ def load_files(record_uuids): def _file_for_state(bucket_id, file_entry): - """One file entry in the shape ``_load_record_state`` writes.""" - legacy_file_id = file_entry["legacy_file_id"] + """One file entry in the shape ``_load_record_state`` writes. + + Files added after migration have no ``legacy_file_id``. Those are not + part of the state file; the caller logs and drops the ``None``. + """ + legacy_file_id = file_entry.get("legacy_file_id") if legacy_file_id is None: - raise KeyError("legacy_file_id") + return None return { "legacy_file_id": legacy_file_id, "bucket_id": bucket_id, @@ -243,32 +253,58 @@ def _file_for_state(bucket_id, file_entry): } +def _version_files(legacy_recid, version_index, version): + """Migrated files for one version. Non-migrated files are logged and skipped.""" + files = [] + for file_entry in version["files"]: + formatted = _file_for_state(version["bucket_id"], file_entry) + if formatted is None: + log( + "skipping (non-migrated record) file with no legacy_file_id: " + f"legacy_recid={legacy_recid} " + f"version={version_index} " + f"new_recid={version['new_recid']} " + f"record_uuid={version['record_uuid']} " + f"file_key={file_entry.get('file_key')}" + ) + continue + files.append(formatted) + return files + + def build_state_entry(legacy_recid, rdm_versions, parent_recid): """Rebuild one ``rdm_records_state.json`` entry. Same fields as ``RecordLoad._load_record_state`` - (cds_migrator_kit/rdm/records/load/entities/record.py). The latest - version is the highest ``index`` still published. + (cds_migrator_kit/rdm/records/load/entities/record.py). ``latest_version`` + is the highest ``index`` that still has a migrated file. A record + migrated with no files keeps the highest published index so pageviews + have a target. """ recid_state = { "legacy_recid": str(legacy_recid), "parent_recid": parent_recid, "versions": [], } + latest_version = None for version_index in sorted(rdm_versions): version = rdm_versions[version_index] if "parent_object_uuid" not in recid_state: recid_state["parent_object_uuid"] = version["parent_object_uuid"] + files = _version_files(legacy_recid, version_index, version) + if not files: + # No migrated files: created after migration, or the files are gone. + continue recid_state["versions"].append( { "new_recid": version["new_recid"], "version": version_index, - "files": [ - _file_for_state(version["bucket_id"], f) for f in version["files"] - ], + "files": files, } ) - latest_version = rdm_versions[max(rdm_versions)] + latest_version = version + if latest_version is None: + latest_version = rdm_versions[max(rdm_versions)] recid_state["latest_version"] = latest_version["new_recid"] recid_state["latest_version_object_uuid"] = latest_version["record_uuid"] return recid_state @@ -285,7 +321,7 @@ def write_state_file(filepath, entries): json_str = json.dumps(entry, ensure_ascii=False, separators=(",", ":")) comma = "," if i < len(entries) - 1 else "" f.write(f"{json_str}{comma}\n") - f.write("]") + f.write("]\n") def main(community_ids, output_path, log_file, dry_run=True): @@ -307,7 +343,7 @@ def main(community_ids, output_path, log_file, dry_run=True): f"starting, community_ids={str(community_ids)}, " f"dry_run={dry_run}, found legacy_recids={len(legacy_recids)}" ) - versions, failed = load_versions(legacy_recids) + versions, failed = load_versions(legacy_recids.values()) record_uuids = [ version["record_uuid"] for by_index in versions.values() @@ -322,10 +358,10 @@ def main(community_ids, output_path, log_file, dry_run=True): entries = [] for legacy_recid, parent_uuid in legacy_recids.items(): stats["checked"] += 1 - if legacy_recid in failed: + if parent_uuid in failed: stats["errors"] += 1 continue - rdm_versions = versions.get(legacy_recid) + rdm_versions = versions.get(parent_uuid) if not rdm_versions: log( f"legacy_recid={legacy_recid} - no published versions found, skipping"