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