diff --git a/site/cds_rdm/inspire_harvester/reader.py b/site/cds_rdm/inspire_harvester/reader.py index 22c2d346..4c320e37 100644 --- a/site/cds_rdm/inspire_harvester/reader.py +++ b/site/cds_rdm/inspire_harvester/reader.py @@ -15,6 +15,8 @@ from cds_rdm.inspire_harvester.transform.resource_types import ALL_DOCUMENT_TYPES +INSPIRE_LITERATURE_API = "https://inspirehep.net/api/literature" + class InspireHTTPReader(BaseReader): """INSPIRE HTTP Reader.""" @@ -40,40 +42,81 @@ def __init__( super().__init__(origin, mode, *args, **kwargs) + def _build_url(self, q, **params): + """Build an INSPIRE literature search URL.""" + query_params = {"q": q, **params} + return f"{INSPIRE_LITERATURE_API}?{urlencode(query_params)}" + + def _get_json(self, url, headers): + """Fetch JSON from INSPIRE or raise ReaderError.""" + current_app.logger.info(f"Querying URL: {url}.") + response = requests.get(url, headers=headers) + if response.status_code != 200: + error_message = ( + f"Error occurred while getting JSON data from INSPIRE. " + f"See URL: {url}. Error message: {response.text}. " + f"Status code: {response.status_code}" + ) + current_app.logger.error(error_message) + raise ReaderError(error_message) + current_app.logger.debug("Request response is successful (200).") + return response.json() + def _iter(self, url, *args, **kwargs): """Yields HTTP response.""" # header set to include additional data (external file URLs and more detailed metadata headers = {"Accept": "application/vnd+inspire.record.expanded+json"} - initial_url = url - - while url: # Continue until there is no "next" link - current_app.logger.info(f"Querying URL: {url}.") - response = requests.get(url, headers=headers) - data = response.json() - if response.status_code == 200: - current_app.logger.debug("Request response is successful (200).") + harvest_url = url + seen_ids = set() + first_pass = True + + while True: + page_url = harvest_url + reported_total = None + new_in_pass = 0 + + while page_url: + data = self._get_json(page_url, headers) total = data["hits"]["total"] hits = data["hits"]["hits"] - if total == 0: - current_app.logger.warning( - f"No results found when querying INSPIRE. See URL: {url}." - ) - elif url == initial_url: - current_app.logger.info(f"Records found: {total}.") + if reported_total is None: + reported_total = total + if total == 0: + current_app.logger.warning( + f"No results found when querying INSPIRE. See URL: {page_url}." + ) + else: + current_app.logger.info(f"Records found: {total}.") for inspire_record in hits: + record_id = str(inspire_record["id"]) + if record_id in seen_ids: + continue + seen_ids.add(record_id) + new_in_pass += 1 current_app.logger.debug( - f"Sending INSPIRE record #{inspire_record['id']} to transformer." + f"Sending INSPIRE record #{record_id} to transformer." ) yield inspire_record - else: - error_message = f"Error occurred while getting JSON data from INSPIRE. See URL: {url}. Error message: {response.text}. Status code: {response.status_code}" - current_app.logger.error(error_message) - raise ReaderError(error_message) - # Get the next page URL if available - url = data.get("links", {}).get("next") + page_url = data.get("links", {}).get("next") + + if len(seen_ids) == reported_total: + return + + if not first_pass and new_in_pass == 0: + current_app.logger.warning( + "Harvest retry added no new INSPIRE records; stopping. " + f"| details: harvested={len(seen_ids)}, reported={reported_total}" + ) + return + + first_pass = False + current_app.logger.info( + "Harvested fewer INSPIRE records than reported; harvesting again. " + f"| details: harvested={len(seen_ids)}, reported={reported_total}" + ) def read(self, item=None, *args, **kwargs): """Builds a query depending on the input data.""" @@ -120,9 +163,7 @@ def read(self, item=None, *args, **kwargs): ) query_params = {"q": f"{q} AND du >= {self._since}"} - base_url = "https://inspirehep.net/api/literature" - encoded_query = urlencode(query_params) - url = f"{base_url}?{encoded_query}" + url = self._build_url(query_params["q"]) current_app.logger.info( f"Resulting query: {query_params['q']}. URL for harvesting data from INSPIRE: {url}." diff --git a/site/tests/inspire_harvester/test_harvester_job.py b/site/tests/inspire_harvester/test_harvester_job.py index ad17050e..e9934078 100644 --- a/site/tests/inspire_harvester/test_harvester_job.py +++ b/site/tests/inspire_harvester/test_harvester_job.py @@ -441,3 +441,113 @@ def mock_requests_get_pagination( tranformation(created_record2.to_dict()["hits"]["hits"][0]["id"], expected_result_2) tranformation(created_record3.to_dict()["hits"]["hits"][0]["id"], expected_result_3) + + +def test_inspire_job_recovers_pagination_shift(running_app, scientific_community, caplog): + """Full harvest persists a record skipped by mid-pagination INSPIRE shifts.""" + page_1_file = DATA_DIR / "inspire_response_15_records_page_1.json" + page_2_file = DATA_DIR / "inspire_response_15_records_page_2.json" + with open(page_1_file) as f: + page_1_data = json.load(f) + with open(page_2_file) as f: + page_2_data = json.load(f) + + by_id = { + str(hit["id"]): hit + for hit in page_1_data["hits"]["hits"] + page_2_data["hits"]["hits"] + } + # Known-good thesis fixtures from test_inspire_job. + record_a = by_id["2802969"] + record_b = by_id["1452604"] + record_c = by_id["2840463"] # skipped in first pass + skipped_id = str(record_c["id"]) + + page_1 = { + "hits": {"total": 3, "hits": [record_a, record_b]}, + "links": { + "next": ( + "https://inspirehep.net/api/literature" + "?q=_oai.sets%3AForCDS+AND+du+%3E%3D+2024-11-15+AND+du+%3C%3D+2025-01-09" + "&size=2&page=2" + ) + }, + } + # After a live update, record C moved off page 2, so the first pass never + # sees it and harvests 2 of the reported 3 records. + page_2_shifted = { + "hits": {"total": 3, "hits": []}, + "links": {}, + } + # By the time the reader harvests again, record C is back in view. + page_2_retry = { + "hits": {"total": 3, "hits": [record_c]}, + "links": {}, + } + page_2_calls = {"n": 0} + + ds_config = { + "config": { + "readers": [ + { + "type": "inspire-http-reader", + "args": { + "since": "2024-11-15", + "until": "2025-01-09", + }, + }, + ], + "transformers": [{"type": "inspire-json-transformer"}], + "writers": [ + { + "type": "async", + "args": { + "writer": { + "type": "inspire-writer", + } + }, + } + ], + "batch_size": 100, + "write_many": True, + } + } + + def mock_requests_get_shift( + url, + headers={"Accept": "application/vnd+inspire.record.expanded+json"}, + stream=True, + ): + if "page=2" in url: + page_2_calls["n"] += 1 + content = page_2_shifted if page_2_calls["n"] == 1 else page_2_retry + else: + content = page_1 + return mock_requests_get(url, mock_content=content) + + run_harvester_mock(ds_config, mock_requests_get_shift) + + RDMRecord.index.refresh() + + assert ( + "Harvested fewer INSPIRE records than reported; harvesting again." + in caplog.text + ) + assert "harvested=2, reported=3" in caplog.text + assert skipped_id in caplog.text + + for inspire_id in (str(record_a["id"]), str(record_b["id"]), skipped_id): + created = current_rdm_records_service.search( + system_identity, + params={"q": f"metadata.related_identifiers.identifier:{inspire_id}"}, + ) + assert created.total == 1, f"Expected CDS record for INSPIRE#{inspire_id}" + + # Explicitly prove the shifted record was persisted. + skipped_record = current_rdm_records_service.search( + system_identity, + params={"q": f"metadata.related_identifiers.identifier:{skipped_id}"}, + ) + assert ( + skipped_record.to_dict()["hits"]["hits"][0]["metadata"]["title"] + == record_c["metadata"]["titles"][0]["title"] + ) diff --git a/site/tests/inspire_harvester/test_reader.py b/site/tests/inspire_harvester/test_reader.py index 70676872..0852ffa3 100644 --- a/site/tests/inspire_harvester/test_reader.py +++ b/site/tests/inspire_harvester/test_reader.py @@ -62,3 +62,117 @@ def test_reader_success(running_app): assert "metadata" in data assert "id" in data assert "links" in data + + +def test_reader_recovers_missing_records(running_app, caplog): + """Test reader harvests again when pagination skipped a record.""" + page_1 = { + "hits": { + "total": 6, + "hits": [{"id": "1"}, {"id": "2"}, {"id": "3"}], + }, + "links": { + "next": "https://inspirehep.net/api/literature?q=test&page=2", + }, + } + # Record 4 moved to page 1 after a live update, so page 2 no longer has it. + page_2 = { + "hits": { + "total": 6, + "hits": [{"id": "5"}, {"id": "6"}], + }, + "links": {}, + } + page_2_retry = { + "hits": { + "total": 6, + "hits": [{"id": "4"}, {"id": "5"}, {"id": "6"}], + }, + "links": {}, + } + page_2_calls = {"n": 0} + + def side_effect(url, headers=None): + mock_response = Mock() + mock_response.status_code = 200 + if "page=2" in url: + page_2_calls["n"] += 1 + mock_response.json.return_value = ( + page_2 if page_2_calls["n"] == 1 else page_2_retry + ) + else: + mock_response.json.return_value = page_1 + return mock_response + + with patch("requests.get", side_effect=side_effect): + reader = InspireHTTPReader(since="2024-01-01", until="2024-01-02") + records = list(reader.read()) + + assert [str(r["id"]) for r in records] == ["1", "2", "3", "5", "6", "4"] + assert "Harvested fewer INSPIRE records than reported; harvesting again." in caplog.text + + +def test_reader_does_not_retry_when_counts_match(running_app, caplog): + """Test reader stops when harvested IDs already match hits.total.""" + page_1 = { + "hits": { + "total": 6, + "hits": [{"id": "1"}, {"id": "2"}, {"id": "3"}], + }, + "links": { + "next": "https://inspirehep.net/api/literature?q=test&page=2", + }, + } + page_2 = { + "hits": { + "total": 6, + "hits": [{"id": "4"}, {"id": "5"}, {"id": "6"}], + }, + "links": {}, + } + page_1_calls = {"n": 0} + + def side_effect(url, headers=None): + mock_response = Mock() + mock_response.status_code = 200 + if "page=2" in url: + mock_response.json.return_value = page_2 + else: + page_1_calls["n"] += 1 + mock_response.json.return_value = page_1 + return mock_response + + with patch("requests.get", side_effect=side_effect) as mock_get: + reader = InspireHTTPReader(since="2024-01-01", until="2024-01-02") + records = list(reader.read()) + + assert [str(r["id"]) for r in records] == ["1", "2", "3", "4", "5", "6"] + assert "harvesting again" not in caplog.text + assert page_1_calls["n"] == 1 + assert all("fields=id" not in call.args[0] for call in mock_get.call_args_list) + + +def test_reader_skips_recovery_for_single_page(running_app, caplog): + """Single-page harvests do not run again when counts already match.""" + page_1 = { + "hits": { + "total": 2, + "hits": [{"id": "1"}, {"id": "2"}], + }, + "links": {}, + } + + def side_effect(url, headers=None): + mock_response = Mock() + mock_response.status_code = 200 + mock_response.json.return_value = page_1 + return mock_response + + with patch("requests.get", side_effect=side_effect) as mock_get: + reader = InspireHTTPReader(since="2024-01-01", until="2024-01-02") + records = list(reader.read()) + + assert [str(r["id"]) for r in records] == ["1", "2"] + assert "harvesting again" not in caplog.text + assert all("fields=id" not in call.args[0] for call in mock_get.call_args_list) + assert mock_get.call_count == 1