Source code for openworld_radio_twin.providers.overture

import asyncio
import math
import time
from concurrent.futures import ThreadPoolExecutor
from dataclasses import dataclass

from shapely import from_wkb
from shapely.geometry import MultiPolygon, Polygon

from openworld_radio_twin.config import Settings
from openworld_radio_twin.geodesy import LocalFrame
from openworld_radio_twin.models import BuildingFeature
from openworld_radio_twin.providers.base import BuildingProviderError, BuildingQueryResult
from openworld_radio_twin.providers.overture_reader import (
    BUILDING_COLUMNS,
    BUILDING_PART_COLUMNS,
    latest_release,
    read_batches,
)


[docs] @dataclass class CacheEntry: expires_at: float value: BuildingQueryResult
[docs] class OvertureBuildingProvider: """Streams buildings for one local WGS84 bounding box from Overture GeoParquet.""" def __init__(self, settings: Settings) -> None: self.settings = settings self._cache: dict[str, CacheEntry] = {} self._query_lock = asyncio.Lock() self._release: str | None = None
[docs] async def buildings( self, latitude: float, longitude: float, radius_m: int, limit: int, ) -> BuildingQueryResult: key = f"{latitude:.5f}:{longitude:.5f}:{radius_m}:{limit}" cached = self._cache.get(key) if cached and cached.expires_at > time.monotonic(): return cached.value async with self._query_lock: cached = self._cache.get(key) if cached and cached.expires_at > time.monotonic(): return cached.value result = await asyncio.to_thread( self._query, latitude, longitude, radius_m, limit, ) self._cache[key] = CacheEntry( expires_at=time.monotonic() + self.settings.cache_ttl_seconds, value=result, ) return result
def _query( self, latitude: float, longitude: float, radius_m: int, limit: int, ) -> BuildingQueryResult: try: if self._release is None: self._release = latest_release() bbox = self._bbox(longitude, latitude, radius_m) # Buildings and parts live in separate objects; fetch both at once and # decode afterwards, since part heights may fall back to their parents. with ThreadPoolExecutor(max_workers=2, thread_name_prefix="owrt-overture") as pool: building_future = pool.submit( self._read_batches, "building", bbox, BUILDING_COLUMNS ) part_future = pool.submit( self._read_batches, "building_part", bbox, BUILDING_PART_COLUMNS ) building_reader = building_future.result() part_reader = part_future.result() building_features = self._read_features( building_reader, limit * 2, self._release, feature_type="building", ) parent_roofs = { feature.source_dataset_id: ( feature.roof_height_agl_m, feature.height_uncertainty_m, ) for feature in building_features if feature.source_dataset_id is not None } part_features = self._read_features( part_reader, limit * 2, self._release, feature_type="building_part", parent_roofs=parent_roofs, ) features = self._compose_features(building_features, part_features, limit) except BuildingProviderError: raise except Exception as exc: raise BuildingProviderError(f"Overture query failed: {exc}") from exc return BuildingQueryResult( buildings=features, provider="Overture Maps", release=self._release, ) def _read_batches(self, overture_type: str, bbox, columns: tuple[str, ...]): return read_batches( overture_type, bbox, self._release, columns, connect_timeout=self.settings.overture_connect_timeout_seconds, request_timeout=self.settings.overture_request_timeout_seconds, ) @staticmethod def _bbox( longitude: float, latitude: float, radius_m: float, ) -> tuple[float, float, float, float]: frame = LocalFrame.at(longitude, latitude) corners = frame.corners(-radius_m, -radius_m, radius_m, radius_m) longitudes = [point[0] for point in corners] latitudes = [point[1] for point in corners] return min(longitudes), min(latitudes), max(longitudes), max(latitudes) @classmethod def _read_features( cls, reader: object, limit: int, release: str, feature_type: str = "building", parent_roofs: dict[str, tuple[float, float | None]] | None = None, ) -> list[BuildingFeature]: features: list[BuildingFeature] = [] for batch in reader: for row in batch.to_pylist(): if row.get("is_underground"): continue geometry = from_wkb(row["geometry"]) if not isinstance(geometry, (Polygon, MultiPolygon)): continue polygons = [geometry] if isinstance(geometry, Polygon) else list(geometry.geoms) min_height_m = cls._minimum_height(row) parent_id = str(row["building_id"]) if row.get("building_id") else None fallback = (parent_roofs or {}).get(parent_id) if parent_id else None height_m, height_source, uncertainty_m = cls._height( row, fallback_roof=fallback, minimum_height_m=min_height_m, ) sources = row.get("sources") or [] source_record_id = sources[0].get("record_id") if sources else None names = row.get("names") or {} for polygon_index, polygon in enumerate(polygons): exterior = cls._ring(polygon.exterior.coords) holes = [cls._ring(interior.coords) for interior in polygon.interiors] dataset_id = str(row["id"]) feature_id = f"overture:{feature_type}:{dataset_id}:{polygon_index}" features.append( BuildingFeature( feature_id=feature_id, source="Overture Maps", feature_type=feature_type, source_dataset_id=dataset_id, parent_building_id=parent_id, source_record_id=source_record_id, source_release=release, name=names.get("primary"), height_m=height_m, min_height_m=min_height_m, height_source=height_source, height_uncertainty_m=uncertainty_m, coordinates=exterior, holes=holes, ) ) if len(features) >= limit: return features return features @staticmethod def _compose_features( buildings: list[BuildingFeature], parts: list[BuildingFeature], limit: int, ) -> list[BuildingFeature]: parts_by_parent: dict[str, list[BuildingFeature]] = {} orphan_parts: list[BuildingFeature] = [] for part in parts: if part.parent_building_id: parts_by_parent.setdefault(part.parent_building_id, []).append(part) else: orphan_parts.append(part) result: list[BuildingFeature] = [] expanded_parents: set[str] = set() for building in buildings: dataset_id = building.source_dataset_id children = parts_by_parent.get(dataset_id or "", []) if children: if dataset_id not in expanded_parents: result.extend(children) expanded_parents.add(dataset_id) else: result.append(building) if len(result) >= limit: return result[:limit] for parent_id, children in parts_by_parent.items(): if parent_id not in expanded_parents: result.extend(children) result.extend(orphan_parts) return result[:limit] @staticmethod def _ring(points: object) -> list[list[float]]: return [[float(longitude), float(latitude)] for longitude, latitude, *_ in points] @staticmethod def _finite_number(value: object) -> float | None: if value is None: return None number = float(value) return number if math.isfinite(number) else None @classmethod def _minimum_height(cls, row: dict[str, object]) -> float: explicit = cls._finite_number(row.get("min_height")) if explicit is not None and 0 <= explicit <= 1000: return explicit floor = cls._finite_number(row.get("min_floor")) if floor is not None and 0 < floor <= 300: return floor * 3.2 return 0.0 @classmethod def _height( cls, row: dict[str, object], fallback_roof: tuple[float, float | None] | None = None, minimum_height_m: float = 0.0, ) -> tuple[float, str, float | None]: explicit = cls._finite_number(row.get("height")) if explicit is not None and minimum_height_m < explicit <= 1000: return explicit - minimum_height_m, "overture_height", None floors = cls._finite_number(row.get("num_floors")) if floors is not None and 0 < floors <= 300: return floors * 3.2, "estimated_from_overture_floors", 1.6 if fallback_roof is not None and fallback_roof[0] > minimum_height_m: return ( fallback_roof[0] - minimum_height_m, "inherited_parent_roof_height", fallback_roof[1], ) return 10.0, "default_assumption", 5.0