Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 `~<n>` suffix.
Expand Down
192 changes: 133 additions & 59 deletions fiboa_cli/conversion/converter_rest.py
Original file line number Diff line number Diff line change
@@ -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,
Comment thread
m-mohr marked this conversation as resolved.
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

Comment thread
m-mohr marked this conversation as resolved.

class EsriRESTConverterMixin:
cache_folder = None
Expand Down Expand Up @@ -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":
Expand All @@ -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
Expand All @@ -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):
Comment thread
m-mohr marked this conversation as resolved.
return {}
pattern = re.compile(re.escape(prefix) + r"(-?\d+)-(-?\d+)\.")
windows = {}
Expand All @@ -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)
Comment thread
m-mohr marked this conversation as resolved.

@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",
Expand All @@ -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())))
4 changes: 1 addition & 3 deletions fiboa_cli/datasets/es_cm.py
Original file line number Diff line number Diff line change
@@ -1,7 +1,5 @@
import re

import requests

from fiboa_cli.conversion.converter_rest import EsriRESTConverterMixin
from fiboa_cli.datasets.es_base import ESBaseConverter

Expand Down Expand Up @@ -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
Expand Down
Loading
Loading