diff --git a/CHANGELOG.md b/CHANGELOG.md index 1169f23d..2cc24827 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -75,6 +75,10 @@ and this project adheres to [Semantic Versioning](http://semver.org/spec/v2.0.0. - Downloaded data cached for one dataset, edition or service is no longer served for another. Previously cached downloads are fetched again once. - Error responses and interrupted downloads are no longer cached. - Layers that join several tables (some ES-CB and ES-IB editions) are now filtered and paged correctly, and their column names no longer carry table prefixes. + - Service metadata and every page are now retried when the service fails (a 5xx error, a timeout or a dropped connection), instead of one failure ending the run. Requests the service rejects (4xx) are not retried. + - Requests time out after 3 minutes, so a hung connection is retried instead of blocking the run. + - A cached page that cannot be read is fetched again instead of ending the run. + - The key field of a joined layer is read from the layer's metadata, which ES-IB answers where it refuses a query. - Fixed `use_variant_as_determination` so determination dates are retained. - `metrics:area` measured from the geometry is now correct in CRSs that are in metres but not equal-area, such as Web Mercator: they are reprojected to an equal-area CRS first. Source areas in hectares are converted to m² also where missing values are filled in, and empty values are filled in, not only 0. Invalid geometries are repaired before they are measured, as they are for the output, so that e.g. a self-intersecting polygon no longer gets an area of 0. - Converters no longer publish duplicate `id`s (#282): row-numbered ids count over all source files instead of restarting per file, and ids the source repeats get a `~` suffix. diff --git a/fiboa_cli/conversion/converter_rest.py b/fiboa_cli/conversion/converter_rest.py index 05ffd299..3c9d661c 100644 --- a/fiboa_cli/conversion/converter_rest.py +++ b/fiboa_cli/conversion/converter_rest.py @@ -1,13 +1,57 @@ +import json import os import re import time import zlib from urllib.parse import urlencode +import aiohttp import geopandas as gpd import requests from vecorel_cli.vecorel.util import get_fs, stream_file +REST_BACKOFF = (1, 5, 15, 60) # seconds before each retry +REST_TIMEOUT = 180 # seconds + + +class RESTError(RuntimeError): + def __init__(self, message, status=None): + super().__init__(message) + self.status = status + + +def _raise_for_esri_error(payload): + # Esri answers a failed request with 200 and an error body + error = payload.get("error") if isinstance(payload, dict) else None + if error is None: + return + status = error.get("code") + message = "; ".join([error.get("message") or "Unknown error", *(error.get("details") or [])]) + raise RESTError(f"{message} (code {status})", status) + + +def _is_transient(error): + if isinstance( + error, + ( + ConnectionError, + TimeoutError, + requests.ConnectionError, + requests.exceptions.ChunkedEncodingError, + requests.Timeout, + aiohttp.ClientConnectionError, + aiohttp.ClientPayloadError, + ), + ): + return True + status = getattr(error, "status", None) + if status is None: + status = getattr(getattr(error, "response", None), "status_code", None) + try: + return int(status) >= 500 + except (TypeError, ValueError): + return False + class EsriRESTConverterMixin: cache_folder = None @@ -72,34 +116,23 @@ def get_data(self, paths, **kwargs): return base_url = paths[0] # loop over paths to support more than 1 source - source_fs = get_fs(base_url) + source_fs = get_fs( + base_url, client_kwargs={"timeout": aiohttp.ClientTimeout(total=REST_TIMEOUT)} + ) cache_fs, cache_folder = self.get_cache(self.cache_folder) - service_metadata = requests.get(base_url, {"f": "pjson"}).json() + service_metadata = self._rest_json(base_url, {"f": "pjson"}) layer = self.rest_layer_filter(service_metadata["layers"]) page_size = service_metadata["maxRecordCount"] layer_url = f"{base_url}/{layer['id']}/query" - # Joined layers qualify every field with the table name; discover the - # real key field before paging on it ("OBJECTID" alone fails there). - probe = requests.get( - layer_url, - { - "f": "json", - "where": "1=1", - "outFields": "*", - "resultRecordCount": 1, - "returnGeometry": "false", - }, - ).json() - attribute = self.rest_attribute - if probe.get("features"): - names = list(probe["features"][0]["attributes"].keys()) - attribute = next( - (n for n in names if n == self.rest_attribute), - next( - (n for n in names if n.endswith("." + self.rest_attribute)), self.rest_attribute - ), - ) + # Joined layers qualify every field with the table name, so read the key + # field from the layer's metadata: es_ib refuses a "where=1=1" probe there. + layer_metadata = self._rest_json(f"{base_url}/{layer['id']}", {"f": "pjson"}) + names = [field["name"] for field in layer_metadata.get("fields") or []] + attribute = next( + (n for n in names if n == self.rest_attribute), + next((n for n in names if n.endswith("." + self.rest_attribute)), self.rest_attribute), + ) base_where = self.rest_params.get("where") # Page by half-open id windows rather than orderByFields + "id > last": @@ -125,45 +158,47 @@ def get_data(self, paths, **kwargs): windows = self._cached_windows(cache_fs, cache_folder, prefix) page = 0 lo = min_id - 1 + + def page_file(lo, hi): + return os.path.join(cache_folder, f"{prefix}{lo}-{hi}.{self.rest_format}") + while lo < max_id: cached = lo in windows hi = windows[lo] if cached else lo + page_size + cache_file = page_file(lo, hi) if cached: - url = os.path.join(cache_folder, f"{prefix}{lo}-{hi}.{self.rest_format}") - else: + try: + data = gpd.read_file(cache_file) + except Exception as e: + self.warning(f"Cached page {cache_file} is unreadable, fetching it again: {e}") + cache_fs.rm(cache_file) + cached = False + hi = lo + page_size + cache_file = page_file(lo, hi) + if not cached: clause = f"{attribute}>{lo} AND {attribute}<={hi}" get_dict["where"] = f"{clause} AND ({base_where})" if base_where else clause url = f"{layer_url}?{urlencode(get_dict)}" - if cache_fs is not None: - cache_file = os.path.join(cache_folder, f"{prefix}{lo}-{hi}.{self.rest_format}") - try: - with cache_fs.open(cache_file, mode="wb") as file: - stream_file(source_fs, url, file) - except Exception: - # A download that broke off must not survive as a cached page - if cache_fs.exists(cache_file): - cache_fs.rm(cache_file) - raise - url = cache_file - - try: - data = gpd.read_file(url) - except Exception as e: - # An error response from the server must not survive as a cached page - if cache_fs is not None and cache_fs.exists(url): - cache_fs.rm(url) - raise RuntimeError(f"Could not read ids ({lo} ... {hi}] of {layer_url}: {e}") from e + try: + data = self._rest_retry( + f"ids ({lo} ... {hi}]", + lambda: self._rest_page(source_fs, cache_fs, url, cache_file), + ) + except Exception as e: + raise RuntimeError( + f"Could not read ids ({lo} ... {hi}] of {layer_url}: {e}" + ) from e if len(data) == 0 and not cached: # An id gap wider than a page: ask once where the ids resume, and # let the empty page cover the whole gap on later runs, instead of # paging through a span that may hold millions of absent ids. resume = self._rest_id_bound(layer_url, attribute, base_where, "ASC", floor=lo) - if resume - 1 > hi and cache_fs is not None: + if resume - 1 > hi: gap = os.path.join( cache_folder, f"{prefix}{lo}-{resume - 1}.{self.rest_format}" ) - cache_fs.mv(url, gap) + cache_fs.mv(cache_file, gap) hi = max(hi, resume - 1) lo = hi @@ -177,7 +212,7 @@ def get_data(self, paths, **kwargs): def _cached_windows(cache_fs, cache_folder, prefix): """The cached (lo, hi] windows, keyed by lo. The name carries both bounds, so a page is only ever read as exactly the window it was fetched for.""" - if cache_fs is None or not cache_fs.exists(cache_folder): + if not cache_fs.exists(cache_folder): return {} pattern = re.compile(re.escape(prefix) + r"(-?\d+)-(-?\d+)\.") windows = {} @@ -187,7 +222,53 @@ def _cached_windows(cache_fs, cache_folder, prefix): windows[int(match.group(1))] = int(match.group(2)) return windows - def _rest_id_bound(self, layer_url, attribute, base_where, direction, floor=-1, attempts=5): + def _rest_retry(self, what, action, backoff=REST_BACKOFF): + """Retry `action` on 5xx errors, timeouts and dropped connections.""" + for attempt in range(len(backoff) + 1): + try: + return action() + except Exception as e: + if attempt == len(backoff) or not _is_transient(e): + raise + delay = backoff[attempt] + self.warning(f"{what}: {e}, retrying in {delay} s ({attempt + 1}/{len(backoff)})") + time.sleep(delay) + + def _rest_json(self, url, params): + def ask(): + response = requests.get(url, params, timeout=REST_TIMEOUT) + response.raise_for_status() + payload = response.json() + _raise_for_esri_error(payload) + return payload + + return self._rest_retry(url, ask) + + @staticmethod + def _rest_page(source_fs, cache_fs, url, cache_file): + """Download and read one page; neither a download that broke off nor an + error response survives as a cached page.""" + try: + with cache_fs.open(cache_file, mode="wb") as file: + stream_file(source_fs, url, file) + try: + return gpd.read_file(cache_file) + except Exception: + with cache_fs.open(cache_file, mode="rb") as file: + head = file.read(64 * 1024) + if head.lstrip().startswith(b'{"error"'): + try: + payload = json.loads(head) + except ValueError: + payload = None + _raise_for_esri_error(payload) + raise + except Exception: + if cache_fs.exists(cache_file): + cache_fs.rm(cache_file) + raise + + def _rest_id_bound(self, layer_url, attribute, base_where, direction, floor=-1): clause = f"{attribute}>{floor}" params = { "f": "json", @@ -197,14 +278,7 @@ def _rest_id_bound(self, layer_url, attribute, base_where, direction, floor=-1, "orderByFields": f"{attribute} {direction}", "resultRecordCount": 1, } - # This is the one sorted query left, and it is the one a tired server - # gives up on: the Balearic proxy answers two in three with a 502. The - # pages themselves are range queries and do not need this. - for attempt in range(attempts): - try: - response = requests.get(layer_url, params).json() - return int(next(iter(response["features"][0]["attributes"].values()))) - except Exception: - if attempt == attempts - 1: - raise - time.sleep(2**attempt) + response = self._rest_json(layer_url, params) + if not response.get("features"): + raise RuntimeError(f"No features of {layer_url} match {params['where']}") + return int(next(iter(response["features"][0]["attributes"].values()))) diff --git a/fiboa_cli/datasets/es_cm.py b/fiboa_cli/datasets/es_cm.py index d2a16e55..7efef9f6 100644 --- a/fiboa_cli/datasets/es_cm.py +++ b/fiboa_cli/datasets/es_cm.py @@ -1,7 +1,5 @@ import re -import requests - from fiboa_cli.conversion.converter_rest import EsriRESTConverterMixin from fiboa_cli.datasets.es_base import ESBaseConverter @@ -44,7 +42,7 @@ class ESCMConverter(EsriRESTConverterMixin, ESBaseConverter): def get_urls(self): # Always use the year-named service: the unnamed "Recintos_sigpac" service is # whatever year is current (2025 in August 2026) and keys on OBJECTID instead. - services = requests.get(self.rest_base_url, {"f": "pjson"}).json()["services"] + services = self._rest_json(self.rest_base_url, {"f": "pjson"})["services"] layer = next( s["name"] for s in services diff --git a/tests/test_converter_rest.py b/tests/test_converter_rest.py index d428203f..234924cd 100644 --- a/tests/test_converter_rest.py +++ b/tests/test_converter_rest.py @@ -12,9 +12,10 @@ import geopandas as gpd import pytest +import requests from shapely.geometry import Point -from fiboa_cli.conversion.converter_rest import EsriRESTConverterMixin +from fiboa_cli.conversion.converter_rest import REST_BACKOFF, EsriRESTConverterMixin, RESTError from fiboa_cli.conversion.fiboa_converter import FiboaBaseConverter BASE_URL = "https://example.test/arcgis/rest/services/SIXPAC_2024/MapServer" @@ -43,6 +44,19 @@ def _collection(oids, key="OBJECTID"): return {"type": "FeatureCollection", "features": [_feature(o, key) for o in oids]} +class Response: + status_code = 200 + + def __init__(self, payload): + self._payload = payload + + def raise_for_status(self): + pass + + def json(self): + return self._payload + + class FakeService: """Answers the three metadata calls and records every page request.""" @@ -51,26 +65,20 @@ def __init__(self, key="OBJECTID", ids=FEATURE_IDS, page_size=PAGE_SIZE): self.ids = ids self.page_size = page_size self.pages = [] # the `where` of every page the mixin downloaded + self.timeouts = [] def get(self, url, params=None, **kwargs): params = params or {} - - class Response: - def __init__(self, payload): - self._payload = payload - - def json(self): - return self._payload - + self.timeouts.append(kwargs.get("timeout")) if params.get("f") == "pjson": + if url.endswith(f"/{LAYER_ID}"): # the layer's fields name the real key + return Response({"fields": [{"name": self.key}, {"name": "USO_SIGPAC"}]}) return Response( { "layers": [{"id": LAYER_ID, "name": "Recintos"}], "maxRecordCount": self.page_size, } ) - if params.get("outFields") == "*": # the probe for the real key field - return Response({"features": [{"attributes": {self.key: self.ids[0]}}]}) floor = int(re.search(r">(-?\d+)", params["where"]).group(1)) above = [o for o in self.ids if o > floor] bound = max(above) if params["orderByFields"].endswith("DESC") else min(above) @@ -97,6 +105,11 @@ def _rebind_logger(): LoggerMixin.logger = None +@pytest.fixture(autouse=True) +def _no_backoff(monkeypatch): + monkeypatch.setattr("fiboa_cli.conversion.converter_rest.time.sleep", lambda _s: None) + + @pytest.fixture def service(monkeypatch): fake = FakeService() @@ -110,6 +123,12 @@ def _read(converter, cache, service=None): return list(converter.get_data([BASE_URL])) +def test_runs_without_a_cache_folder(service): + pages = list(RESTConverter().get_data([BASE_URL])) + + assert [len(data) for data, *_ in pages] == [2, 2, 1] + + def test_pages_by_id_window_and_caches_per_service(service, tmp_path): converter = RESTConverter() pages = _read(converter, tmp_path) @@ -140,7 +159,8 @@ def test_converter_where_is_kept(service, tmp_path): def test_qualified_key_field_is_discovered(monkeypatch, tmp_path): - """A joined layer qualifies every field with its table name.""" + """A joined layer qualifies every field with its table name; the layer's + own metadata names them, where a one-row probe would be a query.""" fake = FakeService(key="RECINTOS.OBJECTID") monkeypatch.setattr("fiboa_cli.conversion.converter_rest.requests.get", fake.get) monkeypatch.setattr("fiboa_cli.conversion.converter_rest.stream_file", fake.stream) @@ -192,11 +212,114 @@ def broken(_source_fs, url, file): monkeypatch.setattr("fiboa_cli.conversion.converter_rest.stream_file", broken) - with pytest.raises(ConnectionError): + with pytest.raises(RuntimeError, match="connection reset"): _read(RESTConverter(), tmp_path) assert os.listdir(tmp_path) == [] +def test_a_rejected_page_is_not_asked_for_again(service, tmp_path, monkeypatch): + attempts = [] + + def rejected(_source_fs, url, file): + attempts.append(url) + file.write(b'{"error":{"code":400,"message":"Invalid query","details":[]}}') + + monkeypatch.setattr("fiboa_cli.conversion.converter_rest.stream_file", rejected) + + with pytest.raises(RuntimeError, match=r"Invalid query \(code 400\)"): + _read(RESTConverter(), tmp_path) + assert len(attempts) == 1 + assert os.listdir(tmp_path) == [] + + +def test_an_error_page_is_asked_for_again(service, tmp_path, monkeypatch): + attempts = [] + real = service.stream + + def tired(source_fs, url, file): + attempts.append(url) + if len(attempts) == 1: + file.write(b'{"error":{"code":500,"message":"Error performing query"}}') + return + return real(source_fs, url, file) + + monkeypatch.setattr("fiboa_cli.conversion.converter_rest.stream_file", tired) + + pages = _read(RESTConverter(), tmp_path) + + assert [len(data) for data, *_ in pages] == [2, 2, 1] + assert len(attempts) == 4 # three pages, the first of them twice + + +def test_an_unreadable_cached_page_is_fetched_again(service, tmp_path): + broken = tmp_path / f"test_rest_{SERVICE}_{LAYER_ID}_r0-2.geojson" + broken.write_text('{"type":"FeatureColl') + + pages = _read(RESTConverter(), tmp_path) + + assert [len(data) for data, *_ in pages] == [2, 2, 1] + assert len(service.pages) == 3 + assert json.loads(broken.read_text())["features"] + + +def test_a_page_is_asked_for_again(service, tmp_path, monkeypatch): + """Hundreds of pages per edition: one refusal is normal, not fatal.""" + attempts = [] + real = service.stream + + def flaky(source_fs, url, file): + attempts.append(url) + if len(attempts) == 1: + file.write(b'{"type":"FeatureColl') + raise ConnectionError("connection reset") + return real(source_fs, url, file) + + monkeypatch.setattr("fiboa_cli.conversion.converter_rest.stream_file", flaky) + + pages = _read(RESTConverter(), tmp_path) + + assert [len(data) for data, *_ in pages] == [2, 2, 1] + assert len(attempts) == 4 # three pages, the first of them twice + + +def test_service_metadata_is_asked_for_again(service, tmp_path, monkeypatch): + """Esri answers a failed request with 200 and an error body.""" + answers = [{"error": {"code": 502, "message": "Bad Gateway"}}] + real = service.get + + def tired(url, params=None, **kwargs): + if answers: + return Response(answers.pop(0)) + return real(url, params, **kwargs) + + monkeypatch.setattr("fiboa_cli.conversion.converter_rest.requests.get", tired) + + pages = _read(RESTConverter(), tmp_path) + + assert [len(data) for data, *_ in pages] == [2, 2, 1] + assert answers == [] + + +def test_rejected_service_metadata_is_not_asked_for_again(service, tmp_path, monkeypatch): + calls = [] + + def rejected(url, params=None, **kwargs): + calls.append(url) + return Response({"error": {"code": 499, "message": "Token Required"}}) + + monkeypatch.setattr("fiboa_cli.conversion.converter_rest.requests.get", rejected) + + with pytest.raises(RESTError, match="Token Required"): + _read(RESTConverter(), tmp_path) + assert len(calls) == 1 + + +def test_every_request_has_a_timeout(service, tmp_path): + _read(RESTConverter(), tmp_path) + + assert service.timeouts and all(service.timeouts) + + def test_pages_cached_under_the_old_sorted_scheme_are_not_reused(service, tmp_path): """The old file name carries neither service nor filter, so nothing proves where its rows came from; such pages are fetched again under the new key.""" @@ -302,11 +425,18 @@ class Years(RESTConverter): def test_id_bound_is_retried(monkeypatch): """The one sorted query left is the one a tired server gives up on.""" - class Answer: - def json(self): - return {"features": [{"attributes": {"OBJECTID": 7}}]} + class BadGateway(Response): + status_code = 502 + + def raise_for_status(self): + raise requests.HTTPError("502 Bad Gateway", response=self) - answers = [RuntimeError("502"), Answer()] + answers = [ + requests.ConnectionError("connection reset"), + requests.exceptions.ChunkedEncodingError("connection broken"), + BadGateway(None), + Response({"features": [{"attributes": {"OBJECTID": 7}}]}), + ] def get(url, params=None, **kwargs): answer = answers.pop(0) @@ -315,12 +445,25 @@ def get(url, params=None, **kwargs): return answer monkeypatch.setattr("fiboa_cli.conversion.converter_rest.requests.get", get) - monkeypatch.setattr("fiboa_cli.conversion.converter_rest.time.sleep", lambda _s: None) assert RESTConverter()._rest_id_bound("url", "OBJECTID", None, "ASC") == 7 assert answers == [] +def test_retries_are_given_up_after_the_backoff(monkeypatch): + calls = [] + + def get(url, params=None, **kwargs): + calls.append(url) + raise requests.Timeout("read timed out") + + monkeypatch.setattr("fiboa_cli.conversion.converter_rest.requests.get", get) + + with pytest.raises(requests.Timeout): + RESTConverter()._rest_id_bound("url", "OBJECTID", None, "ASC") + assert len(calls) == len(REST_BACKOFF) + 1 + + def test_download_files_passes_the_rest_url_through(tmp_path): converter = RESTConverter() assert converter.download_files({"REST": BASE_URL}, str(tmp_path)) == [BASE_URL] diff --git a/tests/test_converters.py b/tests/test_converters.py index d84c9f84..fa5dc117 100644 --- a/tests/test_converters.py +++ b/tests/test_converters.py @@ -58,16 +58,19 @@ class Response: def __init__(self, payload): self._payload = payload + def raise_for_status(self): + pass + def json(self): return self._payload - def fake_get(url, params=None): + def fake_get(url, params=None, **kwargs): params = params or {} if params.get("f") == "pjson": + if url.endswith("/0"): + return Response({"fields": [{"name": "OBJECTID"}]}) # maxRecordCount above the page length below, so paging stops after one page return Response({"layers": [{"id": 0}], "maxRecordCount": 1000}) - if params.get("outFields") == "*": # probe for the real key field - return Response({"features": [{"attributes": {"OBJECTID": 1}}]}) # the two id bounds the window paging starts from bound = 1000 if params["orderByFields"].endswith("DESC") else 1 return Response({"features": [{"attributes": {"OBJECTID": bound}}]})