From fafdddb12df82d1c0ac6076c4e6f041ddaaa2fce Mon Sep 17 00:00:00 2001 From: Joseph Rhoads Date: Mon, 28 Sep 2026 14:25:13 +0200 Subject: [PATCH 1/2] Add retry logic and improve backup and rollback for Elasticsearch indexing (#559) Co-authored-by: Cursor Agent --- rorapi/common/es_bulk.py | 58 ++++++ rorapi/management/commands/indexror.py | 3 +- rorapi/management/commands/indexrordump.py | 57 ++++-- rorapi/management/commands/legacyindexgrid.py | 4 +- rorapi/management/commands/setup.py | 27 ++- rorapi/tests/tests_unit/tests_indexrordump.py | 167 ++++++++++++++++++ 6 files changed, 293 insertions(+), 23 deletions(-) create mode 100644 rorapi/common/es_bulk.py create mode 100644 rorapi/tests/tests_unit/tests_indexrordump.py diff --git a/rorapi/common/es_bulk.py b/rorapi/common/es_bulk.py new file mode 100644 index 00000000..31bcc9c0 --- /dev/null +++ b/rorapi/common/es_bulk.py @@ -0,0 +1,58 @@ +"""Retry wrapper for Elasticsearch/OpenSearch bulk indexing.""" + +import logging +import random +import time + +from elasticsearch import ConnectionError as ESConnectionError +from elasticsearch import ConnectionTimeout, TransportError + +logger = logging.getLogger(__name__) + +RETRYABLE_STATUS_CODES = frozenset({429, 502, 503, 504}) +DEFAULT_MAX_ATTEMPTS = 6 +DEFAULT_BASE_DELAY = 1.0 +DEFAULT_MAX_DELAY = 30.0 + + +def is_retryable(exc): + """Return True if the exception is a transient ES/OpenSearch failure.""" + if isinstance(exc, (ESConnectionError, ConnectionTimeout)): + return True + if isinstance(exc, TransportError): + status = getattr(exc, 'status_code', None) + try: + return int(status) in RETRYABLE_STATUS_CODES + except (TypeError, ValueError): + return False + return False + + +def bulk_with_retry( + es_client, + body, + max_attempts=DEFAULT_MAX_ATTEMPTS, + base_delay=DEFAULT_BASE_DELAY, + max_delay=DEFAULT_MAX_DELAY): + """Call ``es_client.bulk(body)`` with exponential backoff on transient errors. + + Retries HTTP 429/502/503/504 and connection/timeout errors. After + ``max_attempts`` failures, re-raises the last ``TransportError``. + """ + last_exc = None + for attempt in range(1, max_attempts + 1): + try: + return es_client.bulk(body) + except TransportError as exc: + last_exc = exc + if not is_retryable(exc) or attempt >= max_attempts: + raise + delay = min(max_delay, base_delay * (2 ** (attempt - 1))) + delay = delay * (0.5 + random.random()) + status = getattr(exc, 'status_code', 'unknown') + logger.warning( + 'Transient Elasticsearch bulk error (status=%s), ' + 'attempt %s/%s, retrying in %.2fs: %s', + status, attempt, max_attempts, delay, exc) + time.sleep(delay) + raise last_exc diff --git a/rorapi/management/commands/indexror.py b/rorapi/management/commands/indexror.py index 228f17bc..a09c8a13 100644 --- a/rorapi/management/commands/indexror.py +++ b/rorapi/management/commands/indexror.py @@ -9,6 +9,7 @@ import pathlib import shutil from rorapi.settings import ES7, ES_VARS, DATA +from rorapi.common.es_bulk import bulk_with_retry from django.core.management.base import BaseCommand from elasticsearch import TransportError @@ -163,7 +164,7 @@ def index(dataset, version): # experimental affiliations_match nested doc org['affiliation_match'] = get_affiliation_match_doc(org) body.append(org) - ES7.bulk(body) + bulk_with_retry(ES7, body) except TransportError: err[index.__name__] = f"Indexing error, reverted index back to previous state" ES7.reindex(body={ diff --git a/rorapi/management/commands/indexrordump.py b/rorapi/management/commands/indexrordump.py index e4ce2444..2c680676 100644 --- a/rorapi/management/commands/indexrordump.py +++ b/rorapi/management/commands/indexrordump.py @@ -1,11 +1,9 @@ import json import os import re -import requests import zipfile -import base64 -from io import BytesIO -from rorapi.settings import ES7, ES_VARS, ROR_DUMP, DATA +from rorapi.settings import ES7, ES_VARS, DATA +from rorapi.common.es_bulk import bulk_with_retry from django.core.management.base import BaseCommand from elasticsearch import TransportError @@ -30,7 +28,7 @@ def get_single_search_names_v2(org): yield name["value"] def get_affiliation_match_doc(org): - doc = { + doc = { 'id': org['id'], 'country': org["locations"][0]["geonames_details"]["country_code"], 'status': org['status'], @@ -40,8 +38,20 @@ def get_affiliation_match_doc(org): } return doc -def index_dump(self, filename, index, dataset): - backup_index = '{}-tmp'.format(index) + +def _maybe_backup_index(index, backup_index): + """Create a backup of ``index`` unless one already exists or the index is empty. + + When ``setup`` has already written ``backup_index``, leave it alone so a failed + bulk load can restore the previous live data rather than an empty new index. + """ + if ES7.indices.exists(backup_index): + return + if not ES7.indices.exists(index): + return + count = ES7.count(index=index).get('count', 0) + if count == 0: + return ES7.reindex(body={ 'source': { 'index': index @@ -51,6 +61,11 @@ def index_dump(self, filename, index, dataset): } }) + +def index_dump(self, filename, index, dataset): + backup_index = '{}-tmp'.format(index) + _maybe_backup_index(index, backup_index) + try: for i in range(0, len(dataset), ES_VARS['BULK_SIZE']): body = [] @@ -70,18 +85,22 @@ def index_dump(self, filename, index, dataset): # experimental affiliations_match nested doc org['affiliation_match'] = get_affiliation_match_doc(org) body.append(org) - ES7.bulk(body) - except TransportError: - self.stdout.write(TransportError) + bulk_with_retry(ES7, body) + except TransportError as e: + self.stdout.write(str(e)) self.stdout.write('Reverting to backup index') - ES7.reindex(body={ - 'source': { - 'index': backup_index - }, - 'dest': { - 'index': index - } - }) + if ES7.indices.exists(backup_index): + ES7.reindex(body={ + 'source': { + 'index': backup_index + }, + 'dest': { + 'index': index + } + }) + ES7.indices.delete(backup_index) + raise + if ES7.indices.exists(backup_index): ES7.indices.delete(backup_index) self.stdout.write('ROR dataset ' + filename + ' indexed') @@ -117,7 +136,7 @@ def handle(self, *args, **options): elif 'schema_v2' in json_file: # Legacy format with schema_v2 in filename is_v2_format = True - + if is_v2_format and (options.get('schema') == 2 or options.get('schema') is None): self.stdout.write('Loading JSON') with open(json_path, 'r') as it: diff --git a/rorapi/management/commands/legacyindexgrid.py b/rorapi/management/commands/legacyindexgrid.py index 00345fde..0bb49306 100644 --- a/rorapi/management/commands/legacyindexgrid.py +++ b/rorapi/management/commands/legacyindexgrid.py @@ -71,8 +71,8 @@ def handle(self, *args, **options): } for n in get_nested_ids(org)] body.append(org) ES.bulk(body) - except TransportError: - self.stdout.write(TransportError) + except TransportError as e: + self.stdout.write(str(e)) ES.reindex(body={ 'source': { 'index': backup_index diff --git a/rorapi/management/commands/setup.py b/rorapi/management/commands/setup.py index d4e84dab..fc9f8bfb 100644 --- a/rorapi/management/commands/setup.py +++ b/rorapi/management/commands/setup.py @@ -5,12 +5,36 @@ from rorapi.management.commands.createindex import Command as CreateIndexCommand from rorapi.management.commands.indexrordump import Command as IndexRorDumpCommand from rorapi.management.commands.getrordump import Command as GetRorDumpCommand -from rorapi.settings import ROR_DUMP +from rorapi.settings import ES7, ES_VARS, ROR_DUMP REQUEST_TIMEOUT_SECONDS = 30 logger = logging.getLogger(__name__) +def backup_live_index(stdout): + """Copy organizations-v2 to organizations-v2-tmp before delete/create. + + Preserves previous live data so index_dump can restore it if bulk indexing fails. + """ + index = ES_VARS['INDEX_V2'] + backup_index = '{}-tmp'.format(index) + if not ES7.indices.exists(index): + stdout.write('No existing {} index to back up'.format(index)) + return + if ES7.indices.exists(backup_index): + ES7.indices.delete(backup_index) + stdout.write('Deleted stale backup index {}'.format(backup_index)) + ES7.reindex(body={ + 'source': { + 'index': index + }, + 'dest': { + 'index': backup_index + } + }) + stdout.write('Backed up {} to {}'.format(index, backup_index)) + + def build_github_headers(): token = ROR_DUMP.get('GITHUB_TOKEN') if not token: @@ -101,6 +125,7 @@ def handle(self, *args, **options): if sha: try: GetRorDumpCommand().handle(*args, **options) + backup_live_index(self.stdout) DeleteIndexCommand().handle(*args, **options) CreateIndexCommand().handle(*args, **options) IndexRorDumpCommand().handle(*args, **options) diff --git a/rorapi/tests/tests_unit/tests_indexrordump.py b/rorapi/tests/tests_unit/tests_indexrordump.py new file mode 100644 index 00000000..cb34d3c1 --- /dev/null +++ b/rorapi/tests/tests_unit/tests_indexrordump.py @@ -0,0 +1,167 @@ +from io import StringIO +from unittest import mock + +from django.core.management.base import OutputWrapper +from django.test import SimpleTestCase +from elasticsearch import TransportError + +from rorapi.common.es_bulk import bulk_with_retry, is_retryable +from rorapi.management.commands import indexrordump +from rorapi.settings import ES_VARS + + +class BulkWithRetryTestCase(SimpleTestCase): + + def test_is_retryable_429(self): + self.assertTrue(is_retryable(TransportError(429, 'Too Many Requests'))) + + def test_is_retryable_503(self): + self.assertTrue(is_retryable(TransportError(503, 'Service Unavailable'))) + + def test_is_retryable_400_not_retryable(self): + self.assertFalse(is_retryable(TransportError(400, 'Bad Request'))) + + @mock.patch('rorapi.common.es_bulk.time.sleep') + def test_retries_then_succeeds(self, sleep_mock): + es_client = mock.Mock() + es_client.bulk.side_effect = [ + TransportError(429, 'Too Many Requests'), + {'_items': []}, + ] + result = bulk_with_retry(es_client, [{'index': {}}], max_attempts=3, base_delay=0.01) + self.assertEqual(result, {'_items': []}) + self.assertEqual(es_client.bulk.call_count, 2) + sleep_mock.assert_called_once() + + @mock.patch('rorapi.common.es_bulk.time.sleep') + def test_exhausted_retries_reraises(self, sleep_mock): + es_client = mock.Mock() + es_client.bulk.side_effect = TransportError(429, 'Too Many Requests') + with self.assertRaises(TransportError) as ctx: + bulk_with_retry(es_client, [{'index': {}}], max_attempts=3, base_delay=0.01) + self.assertEqual(ctx.exception.status_code, 429) + self.assertEqual(es_client.bulk.call_count, 3) + self.assertEqual(sleep_mock.call_count, 2) + + def test_non_retryable_raises_immediately(self): + es_client = mock.Mock() + es_client.bulk.side_effect = TransportError(400, 'Bad Request') + with self.assertRaises(TransportError): + bulk_with_retry(es_client, [{'index': {}}], max_attempts=5, base_delay=0.01) + self.assertEqual(es_client.bulk.call_count, 1) + + +class IndexDumpTestCase(SimpleTestCase): + + def setUp(self): + self.command = mock.Mock() + self.command.stdout = OutputWrapper(StringIO()) + self.dataset = [ + { + 'id': 'https://ror.org/01an7q238', + 'status': 'active', + 'names': [ + {'value': 'University of Example', 'types': ['ror_display', 'label']}, + ], + 'external_ids': [], + 'locations': [ + {'geonames_details': {'country_code': 'US'}}, + ], + 'relationships': [], + } + ] + self.index = ES_VARS['INDEX_V2'] + self.backup_index = '{}-tmp'.format(self.index) + + @mock.patch('rorapi.management.commands.indexrordump.bulk_with_retry') + @mock.patch('rorapi.management.commands.indexrordump.ES7') + def test_retry_success_skips_rollback(self, es7_mock, bulk_mock): + # setup already created -tmp; do not overwrite it + es7_mock.indices.exists.side_effect = lambda name: name == self.backup_index + bulk_mock.return_value = {'_items': []} + + indexrordump.index_dump(self.command, 'test.json', self.index, self.dataset) + + bulk_mock.assert_called_once() + # delete backup after success; no restore reindex + es7_mock.indices.delete.assert_called_once_with(self.backup_index) + reindex_calls = es7_mock.reindex.call_args_list + self.assertEqual(reindex_calls, []) + output = self.command.stdout.getvalue() + self.assertIn('indexed', output) + self.assertNotIn('Reverting', output) + + @mock.patch('rorapi.management.commands.indexrordump.bulk_with_retry') + @mock.patch('rorapi.management.commands.indexrordump.ES7') + def test_exhausted_429_rolls_back_and_reraises(self, es7_mock, bulk_mock): + es7_mock.indices.exists.return_value = True + bulk_mock.side_effect = TransportError(429, 'Too Many Requests') + + with self.assertRaises(TransportError) as ctx: + indexrordump.index_dump(self.command, 'test.json', self.index, self.dataset) + + self.assertEqual(ctx.exception.status_code, 429) + output = self.command.stdout.getvalue() + self.assertIn('Too Many Requests', output) + self.assertIn('Reverting to backup index', output) + self.assertNotIn('indexed', output) + + found_restore = False + for call in es7_mock.reindex.call_args_list: + body = call.kwargs.get('body') + if body is None and call.args: + body = call.args[0] + if body and body.get('source', {}).get('index') == self.backup_index: + found_restore = True + break + self.assertTrue(found_restore, 'expected restore reindex from backup') + es7_mock.indices.delete.assert_called_with(self.backup_index) + + @mock.patch('rorapi.management.commands.indexrordump.bulk_with_retry') + @mock.patch('rorapi.management.commands.indexrordump.ES7') + def test_logging_transport_error_instance_does_not_raise_attribute_error( + self, es7_mock, bulk_mock): + """Regression: writing TransportError class caused AttributeError on endswith.""" + es7_mock.indices.exists.return_value = True + bulk_mock.side_effect = TransportError(429, 'Too Many Requests') + + with self.assertRaises(TransportError): + indexrordump.index_dump(self.command, 'test.json', self.index, self.dataset) + + # If the old bug returned (stdout.write(TransportError class)), Django would + # raise AttributeError before we could re-raise TransportError. + output = self.command.stdout.getvalue() + self.assertIn('TransportError', output) + self.assertIn('429', output) + + +class SetupBackupTestCase(SimpleTestCase): + + @mock.patch('rorapi.management.commands.setup.ES7') + def test_backup_live_index_reindexes_when_live_exists(self, es7_mock): + from rorapi.management.commands.setup import backup_live_index + + index = ES_VARS['INDEX_V2'] + backup = '{}-tmp'.format(index) + es7_mock.indices.exists.side_effect = lambda name: name == index + stdout = OutputWrapper(StringIO()) + + backup_live_index(stdout) + + es7_mock.reindex.assert_called_once() + body = es7_mock.reindex.call_args.kwargs.get('body') or es7_mock.reindex.call_args.args[0] + self.assertEqual(body['source']['index'], index) + self.assertEqual(body['dest']['index'], backup) + self.assertIn('Backed up', stdout.getvalue()) + + @mock.patch('rorapi.management.commands.setup.ES7') + def test_backup_live_index_skips_when_missing(self, es7_mock): + from rorapi.management.commands.setup import backup_live_index + + es7_mock.indices.exists.return_value = False + stdout = OutputWrapper(StringIO()) + + backup_live_index(stdout) + + es7_mock.reindex.assert_not_called() + self.assertIn('No existing', stdout.getvalue()) From 0b061be0841ef385d2359c2547372f0e7961301d Mon Sep 17 00:00:00 2001 From: Joseph Rhoads Date: Mon, 28 Sep 2026 15:56:11 +0200 Subject: [PATCH 2/2] Factor shared index helpers out of indexror and indexrordump (#573) Co-authored-by: Cursor Agent --- rorapi/common/index_helpers.py | 155 ++++++++++++++ rorapi/management/commands/indexror.py | 82 +------- rorapi/management/commands/indexrordump.py | 101 +-------- .../tests/tests_unit/tests_index_helpers.py | 197 ++++++++++++++++++ rorapi/tests/tests_unit/tests_indexrordump.py | 12 +- 5 files changed, 374 insertions(+), 173 deletions(-) create mode 100644 rorapi/common/index_helpers.py create mode 100644 rorapi/tests/tests_unit/tests_index_helpers.py diff --git a/rorapi/common/index_helpers.py b/rorapi/common/index_helpers.py new file mode 100644 index 00000000..2ee8c65d --- /dev/null +++ b/rorapi/common/index_helpers.py @@ -0,0 +1,155 @@ +"""Shared helpers for ROR Elasticsearch indexing commands. + +Used by ``indexror`` and ``indexrordump``. Bulk calls go through +``_bulk``, which uses ``bulk_with_retry`` so transient Elasticsearch +errors back off without duplicating the chunk loop. +""" + +import re + +from elasticsearch import TransportError + +from rorapi.common.es_bulk import bulk_with_retry +from rorapi.settings import ES7, ES_VARS + + +def get_nested_names_v2(org): + for name in org['names']: + yield name['value'] + + +def get_nested_ids_v2(org): + yield org['id'] + yield re.sub('https://', '', org['id']) + yield re.sub('https://ror.org/', '', org['id']) + for ext_id in org['external_ids']: + for eid in ext_id['all']: + yield eid + + +def get_single_search_names_v2(org): + for name in org["names"]: + if "acronym" not in name["types"]: + yield name["value"] + + +def get_affiliation_match_doc(org): + return { + 'id': org['id'], + 'country': org["locations"][0]["geonames_details"]["country_code"], + 'status': org['status'], + 'primary': [n["value"] for n in org["names"] if "ror_display" in n["types"]][0], + 'names': [{"name": n} for n in get_single_search_names_v2(org)], + 'relationships': [{"type": r['type'], "id": r['id']} for r in org['relationships']] + } + + +def enrich_org_for_index(org): + """Attach ``names_ids`` and ``affiliation_match`` fields used by the index.""" + org['names_ids'] = [{ + 'name': n + } for n in get_nested_names_v2(org)] + org['names_ids'] += [{ + 'id': n + } for n in get_nested_ids_v2(org)] + org['affiliation_match'] = get_affiliation_match_doc(org) + return org + + +def build_bulk_body(index, orgs): + """Build an ES bulk body for a chunk of organizations.""" + body = [] + for org in orgs: + body.append({ + 'index': { + '_index': index, + '_id': org['id'] + } + }) + enrich_org_for_index(org) + body.append(org) + return body + + +def _bulk(body): + """Perform one bulk request, retrying transient Elasticsearch errors.""" + return bulk_with_retry(ES7, body) + + +def maybe_backup_index(index, backup_index): + """Prepare rollback for ``index`` and return ``'restore'``, ``'clear'``, or ``'none'``. + + An existing ``backup_index`` is left in place so a backup written by + ``setup`` is not replaced. A live index with documents is copied to + ``backup_index`` (``'restore'``). An existing but empty live index returns + ``'clear'``: reindexing an empty backup does not delete documents a + partial bulk load already wrote, so failure handling must wipe the index. + ``'none'`` means there was no live index; a failed bulk may have created + one, and failure handling deletes it. + """ + if ES7.indices.exists(backup_index): + return 'restore' + if not ES7.indices.exists(index): + return 'none' + count = ES7.count(index=index).get('count', 0) + if count == 0: + return 'clear' + ES7.reindex(body={ + 'source': { + 'index': index + }, + 'dest': { + 'index': backup_index + } + }) + return 'restore' + + +def bulk_index_with_backup(index, dataset, on_transport_error=None, bulk=None): + """Backup ``index`` to ``{index}-tmp``, bulk-index ``dataset``, restore on failure. + + Returns ``True`` if bulk indexing completed without ``TransportError``, + ``False`` if a transport error triggered rollback. Does not re-raise; + callers that must surface the failure should do so from + ``on_transport_error``'s stored exception after this returns. A non-empty + backup is reindexed back onto ``index``. An empty live index is cleared + with delete-by-query, because reindexing an empty backup would leave + partial documents in place. If the index did not exist beforehand and a + failed bulk created it, that index is deleted. Deletes the backup index + when it exists. + """ + if bulk is None: + bulk = _bulk + backup_index = '{}-tmp'.format(index) + rollback = maybe_backup_index(index, backup_index) + + try: + for i in range(0, len(dataset), ES_VARS['BULK_SIZE']): + chunk = dataset[i:i + ES_VARS['BULK_SIZE']] + body = build_bulk_body(index, chunk) + bulk(body) + except TransportError as e: + if on_transport_error is not None: + on_transport_error(e) + if rollback == 'restore' and ES7.indices.exists(backup_index): + ES7.reindex(body={ + 'source': { + 'index': backup_index + }, + 'dest': { + 'index': index + } + }) + ES7.indices.delete(backup_index) + elif rollback == 'clear': + ES7.delete_by_query( + index=index, + body={'query': {'match_all': {}}}, + params={'conflicts': 'proceed', 'refresh': True}, + ) + elif rollback == 'none' and ES7.indices.exists(index): + ES7.indices.delete(index) + return False + if ES7.indices.exists(backup_index): + ES7.indices.delete(backup_index) + return True diff --git a/rorapi/management/commands/indexror.py b/rorapi/management/commands/indexror.py index a09c8a13..795fa204 100644 --- a/rorapi/management/commands/indexror.py +++ b/rorapi/management/commands/indexror.py @@ -1,5 +1,4 @@ import json -import re from functools import wraps from threading import local import zipfile @@ -8,39 +7,11 @@ from os.path import exists import pathlib import shutil -from rorapi.settings import ES7, ES_VARS, DATA -from rorapi.common.es_bulk import bulk_with_retry +from rorapi.settings import ES_VARS, DATA +from rorapi.common.index_helpers import bulk_index_with_backup from django.core.management.base import BaseCommand -from elasticsearch import TransportError - -def get_nested_names_v2(org): - for name in org['names']: - yield name['value'] - -def get_nested_ids_v2(org): - yield org['id'] - yield re.sub('https://', '', org['id']) - yield re.sub('https://ror.org/', '', org['id']) - for ext_id in org['external_ids']: - for eid in ext_id['all']: - yield eid - -def get_single_search_names_v2(org): - for name in org["names"]: - if "acronym" not in name["types"]: - yield name["value"] - -def get_affiliation_match_doc(org): - doc = { - 'id': org['id'], - 'country': org["locations"][0]["geonames_details"]["country_code"], - 'status': org['status'], - 'primary': [n["value"] for n in org["names"] if "ror_display" in n["types"]][0], - 'names': [{"name": n} for n in get_single_search_names_v2(org)], - 'relationships': [{"type": r['type'], "id": r['id']} for r in org['relationships']] - } - return doc + def prepare_files(path, local_file): data = [] @@ -134,49 +105,12 @@ def index(dataset, version): if version != 'v2': err[index.__name__] = f"Only v2 schema version is supported. Received: {version}" return err - index = ES_VARS['INDEX_V2'] - backup_index = '{}-tmp'.format(index) - ES7.reindex(body={ - 'source': { - 'index': index - }, - 'dest': { - 'index': backup_index - } - }) + index_name = ES_VARS['INDEX_V2'] - try: - for i in range(0, len(dataset), ES_VARS['BULK_SIZE']): - body = [] - for org in dataset[i:i + ES_VARS['BULK_SIZE']]: - body.append({ - 'index': { - '_index': index, - '_id': org['id'] - } - }) - org['names_ids'] = [{ - 'name': n - } for n in get_nested_names_v2(org)] - org['names_ids'] += [{ - 'id': n - } for n in get_nested_ids_v2(org)] - # experimental affiliations_match nested doc - org['affiliation_match'] = get_affiliation_match_doc(org) - body.append(org) - bulk_with_retry(ES7, body) - except TransportError: + def on_transport_error(_exc): err[index.__name__] = f"Indexing error, reverted index back to previous state" - ES7.reindex(body={ - 'source': { - 'index': backup_index - }, - 'dest': { - 'index': index - } - }) - if ES7.indices.exists(backup_index): - ES7.indices.delete(backup_index) + + bulk_index_with_backup(index_name, dataset, on_transport_error=on_transport_error) return err class Command(BaseCommand): @@ -189,5 +123,3 @@ def handle(self,*args, **options): dir = options['dir'] version = 'v2' process_files(dir, version) - - diff --git a/rorapi/management/commands/indexrordump.py b/rorapi/management/commands/indexrordump.py index 2c680676..c89c5998 100644 --- a/rorapi/management/commands/indexrordump.py +++ b/rorapi/management/commands/indexrordump.py @@ -2,107 +2,24 @@ import os import re import zipfile -from rorapi.settings import ES7, ES_VARS, DATA -from rorapi.common.es_bulk import bulk_with_retry +from rorapi.settings import ES_VARS, DATA +from rorapi.common.index_helpers import bulk_index_with_backup from django.core.management.base import BaseCommand -from elasticsearch import TransportError HEADERS = {'Accept': 'application/vnd.github.v3+json'} -def get_nested_names_v2(org): - for name in org['names']: - yield name['value'] - -def get_nested_ids_v2(org): - yield org['id'] - yield re.sub('https://', '', org['id']) - yield re.sub('https://ror.org/', '', org['id']) - for ext_id in org['external_ids']: - for eid in ext_id['all']: - yield eid - -def get_single_search_names_v2(org): - for name in org["names"]: - if "acronym" not in name["types"]: - yield name["value"] - -def get_affiliation_match_doc(org): - doc = { - 'id': org['id'], - 'country': org["locations"][0]["geonames_details"]["country_code"], - 'status': org['status'], - 'primary': [n["value"] for n in org["names"] if "ror_display" in n["types"]][0], - 'names': [{"name": n} for n in get_single_search_names_v2(org)], - 'relationships': [{"type": r['type'], "id": r['id']} for r in org['relationships']] - } - return doc - - -def _maybe_backup_index(index, backup_index): - """Create a backup of ``index`` unless one already exists or the index is empty. - - When ``setup`` has already written ``backup_index``, leave it alone so a failed - bulk load can restore the previous live data rather than an empty new index. - """ - if ES7.indices.exists(backup_index): - return - if not ES7.indices.exists(index): - return - count = ES7.count(index=index).get('count', 0) - if count == 0: - return - ES7.reindex(body={ - 'source': { - 'index': index - }, - 'dest': { - 'index': backup_index - } - }) - - def index_dump(self, filename, index, dataset): - backup_index = '{}-tmp'.format(index) - _maybe_backup_index(index, backup_index) + caught = {} - try: - for i in range(0, len(dataset), ES_VARS['BULK_SIZE']): - body = [] - for org in dataset[i:i + ES_VARS['BULK_SIZE']]: - body.append({ - 'index': { - '_index': index, - '_id': org['id'] - } - }) - org['names_ids'] = [{ - 'name': n - } for n in get_nested_names_v2(org)] - org['names_ids'] += [{ - 'id': n - } for n in get_nested_ids_v2(org)] - # experimental affiliations_match nested doc - org['affiliation_match'] = get_affiliation_match_doc(org) - body.append(org) - bulk_with_retry(ES7, body) - except TransportError as e: - self.stdout.write(str(e)) + def on_transport_error(exc): + self.stdout.write(str(exc)) self.stdout.write('Reverting to backup index') - if ES7.indices.exists(backup_index): - ES7.reindex(body={ - 'source': { - 'index': backup_index - }, - 'dest': { - 'index': index - } - }) - ES7.indices.delete(backup_index) - raise + caught['exc'] = exc - if ES7.indices.exists(backup_index): - ES7.indices.delete(backup_index) + ok = bulk_index_with_backup(index, dataset, on_transport_error=on_transport_error) + if not ok: + raise caught['exc'] self.stdout.write('ROR dataset ' + filename + ' indexed') diff --git a/rorapi/tests/tests_unit/tests_index_helpers.py b/rorapi/tests/tests_unit/tests_index_helpers.py new file mode 100644 index 00000000..e0b728d9 --- /dev/null +++ b/rorapi/tests/tests_unit/tests_index_helpers.py @@ -0,0 +1,197 @@ +from unittest import mock + +from django.test import SimpleTestCase +from elasticsearch import TransportError + +from rorapi.common import index_helpers +from rorapi.settings import ES_VARS + + +SAMPLE_ORG = { + 'id': 'https://ror.org/01an7q238', + 'status': 'active', + 'names': [ + {'value': 'University of Example', 'types': ['ror_display', 'label']}, + {'value': 'UoE', 'types': ['acronym']}, + {'value': 'Example U', 'types': ['alias']}, + ], + 'external_ids': [ + {'type': 'grid', 'all': ['grid.1234.5']}, + ], + 'locations': [ + {'geonames_details': {'country_code': 'US'}}, + ], + 'relationships': [ + {'type': 'parent', 'id': 'https://ror.org/05dxps055'}, + ], +} + + +class IndexDocBuildersTestCase(SimpleTestCase): + + def test_get_nested_names_v2(self): + self.assertEqual( + list(index_helpers.get_nested_names_v2(SAMPLE_ORG)), + ['University of Example', 'UoE', 'Example U'], + ) + + def test_get_nested_ids_v2(self): + self.assertEqual( + list(index_helpers.get_nested_ids_v2(SAMPLE_ORG)), + [ + 'https://ror.org/01an7q238', + 'ror.org/01an7q238', + '01an7q238', + 'grid.1234.5', + ], + ) + + def test_get_single_search_names_v2_skips_acronyms(self): + self.assertEqual( + list(index_helpers.get_single_search_names_v2(SAMPLE_ORG)), + ['University of Example', 'Example U'], + ) + + def test_get_affiliation_match_doc(self): + doc = index_helpers.get_affiliation_match_doc(SAMPLE_ORG) + self.assertEqual(doc['id'], SAMPLE_ORG['id']) + self.assertEqual(doc['country'], 'US') + self.assertEqual(doc['status'], 'active') + self.assertEqual(doc['primary'], 'University of Example') + self.assertEqual( + doc['names'], + [{'name': 'University of Example'}, {'name': 'Example U'}], + ) + self.assertEqual( + doc['relationships'], + [{'type': 'parent', 'id': 'https://ror.org/05dxps055'}], + ) + + def test_enrich_org_for_index(self): + org = dict(SAMPLE_ORG) + index_helpers.enrich_org_for_index(org) + self.assertIn({'name': 'University of Example'}, org['names_ids']) + self.assertIn({'id': '01an7q238'}, org['names_ids']) + self.assertEqual(org['affiliation_match']['primary'], 'University of Example') + + +class BulkIndexWithBackupTestCase(SimpleTestCase): + + def setUp(self): + self.index = ES_VARS['INDEX_V2'] + self.backup_index = '{}-tmp'.format(self.index) + self.dataset = [dict(SAMPLE_ORG)] + + def _mock_live_index(self, es7_mock, backup_already=False): + state = {'backup': backup_already} + + def exists(name): + if name == self.backup_index: + return state['backup'] + return name == self.index + + def reindex(*args, **kwargs): + body = kwargs.get('body') + if body is None and args: + body = args[0] + if body and body.get('dest', {}).get('index') == self.backup_index: + state['backup'] = True + + es7_mock.indices.exists.side_effect = exists + es7_mock.reindex.side_effect = reindex + es7_mock.count.return_value = {'count': 2} + + @mock.patch('rorapi.common.index_helpers.ES7') + def test_success_backs_up_bulks_and_deletes_backup(self, es7_mock): + self._mock_live_index(es7_mock) + bulk_mock = mock.Mock(return_value={'_items': []}) + + ok = index_helpers.bulk_index_with_backup( + self.index, self.dataset, bulk=bulk_mock) + + self.assertTrue(ok) + bulk_mock.assert_called_once() + body = bulk_mock.call_args.args[0] + self.assertEqual(body[0]['index']['_index'], self.index) + self.assertEqual(body[0]['index']['_id'], SAMPLE_ORG['id']) + self.assertEqual(body[1]['affiliation_match']['primary'], 'University of Example') + es7_mock.indices.delete.assert_called_once_with(self.backup_index) + # backup reindex only (no restore) + self.assertEqual(es7_mock.reindex.call_count, 1) + + @mock.patch('rorapi.common.index_helpers.ES7') + def test_existing_backup_is_not_overwritten(self, es7_mock): + self._mock_live_index(es7_mock, backup_already=True) + bulk_mock = mock.Mock(return_value={'_items': []}) + + ok = index_helpers.bulk_index_with_backup( + self.index, self.dataset, bulk=bulk_mock) + + self.assertTrue(ok) + es7_mock.reindex.assert_not_called() + es7_mock.indices.delete.assert_called_once_with(self.backup_index) + + @mock.patch('rorapi.common.index_helpers.ES7') + def test_transport_error_rolls_back_and_calls_hook(self, es7_mock): + self._mock_live_index(es7_mock) + bulk_mock = mock.Mock(side_effect=TransportError(429, 'Too Many Requests')) + on_error = mock.Mock() + + ok = index_helpers.bulk_index_with_backup( + self.index, self.dataset, on_transport_error=on_error, bulk=bulk_mock) + + self.assertFalse(ok) + on_error.assert_called_once() + self.assertEqual(es7_mock.reindex.call_count, 2) + restore_body = es7_mock.reindex.call_args_list[1].kwargs.get('body') + if restore_body is None: + restore_body = es7_mock.reindex.call_args_list[1].args[0] + self.assertEqual(restore_body['source']['index'], self.backup_index) + self.assertEqual(restore_body['dest']['index'], self.index) + es7_mock.indices.delete.assert_called_once_with(self.backup_index) + + @mock.patch('rorapi.common.index_helpers.ES7') + def test_empty_index_failure_clears_partial_documents(self, es7_mock): + def exists(name): + return name == self.index + + es7_mock.indices.exists.side_effect = exists + es7_mock.count.return_value = {'count': 0} + bulk_mock = mock.Mock(side_effect=TransportError(429, 'Too Many Requests')) + on_error = mock.Mock() + + ok = index_helpers.bulk_index_with_backup( + self.index, self.dataset, on_transport_error=on_error, bulk=bulk_mock) + + self.assertFalse(ok) + on_error.assert_called_once() + es7_mock.reindex.assert_not_called() + es7_mock.indices.delete.assert_not_called() + es7_mock.delete_by_query.assert_called_once_with( + index=self.index, + body={'query': {'match_all': {}}}, + params={'conflicts': 'proceed', 'refresh': True}, + ) + + @mock.patch('rorapi.common.index_helpers.ES7') + def test_missing_index_failure_deletes_index_created_by_bulk(self, es7_mock): + live_checks = {'n': 0} + + def exists(name): + if name == self.backup_index: + return False + live_checks['n'] += 1 + return live_checks['n'] > 1 + + es7_mock.indices.exists.side_effect = exists + bulk_mock = mock.Mock(side_effect=TransportError(503, 'Service Unavailable')) + on_error = mock.Mock() + + ok = index_helpers.bulk_index_with_backup( + self.index, self.dataset, on_transport_error=on_error, bulk=bulk_mock) + + self.assertFalse(ok) + on_error.assert_called_once() + es7_mock.reindex.assert_not_called() + es7_mock.delete_by_query.assert_not_called() + es7_mock.indices.delete.assert_called_once_with(self.index) diff --git a/rorapi/tests/tests_unit/tests_indexrordump.py b/rorapi/tests/tests_unit/tests_indexrordump.py index cb34d3c1..f516be67 100644 --- a/rorapi/tests/tests_unit/tests_indexrordump.py +++ b/rorapi/tests/tests_unit/tests_indexrordump.py @@ -73,8 +73,8 @@ def setUp(self): self.index = ES_VARS['INDEX_V2'] self.backup_index = '{}-tmp'.format(self.index) - @mock.patch('rorapi.management.commands.indexrordump.bulk_with_retry') - @mock.patch('rorapi.management.commands.indexrordump.ES7') + @mock.patch('rorapi.common.index_helpers.bulk_with_retry') + @mock.patch('rorapi.common.index_helpers.ES7') def test_retry_success_skips_rollback(self, es7_mock, bulk_mock): # setup already created -tmp; do not overwrite it es7_mock.indices.exists.side_effect = lambda name: name == self.backup_index @@ -91,8 +91,8 @@ def test_retry_success_skips_rollback(self, es7_mock, bulk_mock): self.assertIn('indexed', output) self.assertNotIn('Reverting', output) - @mock.patch('rorapi.management.commands.indexrordump.bulk_with_retry') - @mock.patch('rorapi.management.commands.indexrordump.ES7') + @mock.patch('rorapi.common.index_helpers.bulk_with_retry') + @mock.patch('rorapi.common.index_helpers.ES7') def test_exhausted_429_rolls_back_and_reraises(self, es7_mock, bulk_mock): es7_mock.indices.exists.return_value = True bulk_mock.side_effect = TransportError(429, 'Too Many Requests') @@ -117,8 +117,8 @@ def test_exhausted_429_rolls_back_and_reraises(self, es7_mock, bulk_mock): self.assertTrue(found_restore, 'expected restore reindex from backup') es7_mock.indices.delete.assert_called_with(self.backup_index) - @mock.patch('rorapi.management.commands.indexrordump.bulk_with_retry') - @mock.patch('rorapi.management.commands.indexrordump.ES7') + @mock.patch('rorapi.common.index_helpers.bulk_with_retry') + @mock.patch('rorapi.common.index_helpers.ES7') def test_logging_transport_error_instance_does_not_raise_attribute_error( self, es7_mock, bulk_mock): """Regression: writing TransportError class caused AttributeError on endswith."""