Source code for openworld_radio_twin.compiler

"""Compile geographic inputs into the canonical reusable scene directory."""

from __future__ import annotations

import asyncio
import errno
import hashlib
import json
import os
import shutil
import tempfile
from collections.abc import Coroutine
from concurrent.futures import ThreadPoolExecutor
from contextlib import asynccontextmanager, contextmanager
from dataclasses import replace
from pathlib import Path
from typing import Any, Literal, TypeVar

import httpx

try:
    import fcntl
except ImportError:  # pragma: no cover - Windows only
    fcntl = None
    import msvcrt

from openworld_radio_twin import __version__
from openworld_radio_twin.config import Settings, get_settings
from openworld_radio_twin.environment import EnvironmentData
from openworld_radio_twin.jsonio import dump_json_with_compact_rows
from openworld_radio_twin.models import BuildingFeature
from openworld_radio_twin.providers.base import SCENE_BOUNDARY_POLICY, BuildingQueryResult
from openworld_radio_twin.providers.building_sources import (
    default_building_sources,
    resolve_building_source,
)
from openworld_radio_twin.providers.map_data import MapDataProvider
from openworld_radio_twin.scene_format import (
    SCENE_CASES_DIRECTORY,
    SCENE_FORMAT,
    SCENE_HEIGHT_MAP_PATH,
    SCENE_INFO_PATH,
    SCENE_LAYOUT_VERSION,
    SCENE_MESH_DIRECTORY,
    SCENE_PROVENANCE_PATH,
    SCENE_SOURCE_BUILDINGS_PATH,
    SCENE_TERRAIN_INFO_PATH,
    SCENE_VOXEL_DEPTH_DIRECTORY,
    SCENE_VOXEL_SLICES_DIRECTORY,
    SCENE_XML_PATH,
)
from openworld_radio_twin.simulation.scene_builder import (
    BUILDING_MESH_SCHEMA_VERSION,
    MEASUREMENT_SURFACE_SCHEMA_VERSION,
    MaterialProfileName,
    SceneAssets,
    build_scene_assets,
    load_scene_assets,
    measurement_surface_path,
    write_measurement_surface,
    write_voxel_assets,
)

BuildingMode = Literal["auto", "none"]
TerrainMode = Literal["elevation", "flat", "none"]
ExistingPolicy = Literal["reuse", "error", "replace"]
CellSize = float | tuple[float, float]
T = TypeVar("T")

__all__ = [
    "BuildingMode",
    "CellSize",
    "ExistingPolicy",
    "TerrainMode",
    "compile_cached_scene",
    "compile_scene",
    "compile_scene_sync",
    "load_building_source",
]


def _geometry_request(
    latitude: float,
    longitude: float,
    radius_m: float,
    buildings: BuildingMode,
    terrain: TerrainMode,
    material_profile: MaterialProfileName,
    terrain_resolution_m: float,
    building_source: str = "auto",
    settings: Settings | None = None,
) -> dict[str, object]:
    if not -90 <= latitude <= 90:
        raise ValueError("latitude must be between -90 and 90 degrees")
    if not -180 <= longitude <= 180:
        raise ValueError("longitude must be between -180 and 180 degrees")
    if radius_m <= 0:
        raise ValueError("radius_m must be positive")
    if terrain_resolution_m <= 0:
        raise ValueError("terrain_resolution_m must be positive")
    if buildings not in {"auto", "none"}:
        raise ValueError(f"Unknown buildings mode: {buildings}")
    if terrain not in {"elevation", "flat", "none"}:
        raise ValueError(f"Unknown terrain mode: {terrain}")
    if material_profile not in {"itu", "uniform"}:
        raise ValueError(f"Unknown material profile: {material_profile}")
    return {
        "latitude": float(latitude),
        "longitude": float(longitude),
        "radius_m": float(radius_m),
        "buildings": buildings,
        "building_source": building_source,
        "source_selection": resolve_building_source(
            building_source if buildings != "none" else "none",
            latitude,
            longitude,
            settings or get_settings(),
            default_building_sources(),
        ),
        "terrain": terrain,
        "material_profile": material_profile,
        "terrain_resolution_m": float(terrain_resolution_m),
    }


def scene_cache_key(request: dict[str, object]) -> str:
    """Return the deterministic cache key for geometry-affecting inputs."""
    payload = {
        "compiler_version": __version__,
        "building_mesh_schema_version": BUILDING_MESH_SCHEMA_VERSION,
        "scene_format": SCENE_FORMAT,
        "layout_version": SCENE_LAYOUT_VERSION,
        # Bumped when the geometry compiled for an unchanged request changes.
        # 2: surface classes are compiled on flat ground as well.
        "surface_layer_revision": 2,
        "request": request,
    }
    encoded = json.dumps(payload, sort_keys=True, separators=(",", ":")).encode()
    return hashlib.sha256(encoded).hexdigest()[:20]


def default_scene_cache_root(settings: Settings | None = None) -> Path:
    runtime_settings = settings or get_settings()
    configured = getattr(runtime_settings, "scene_cache_root", None)
    if configured is not None:
        return Path(configured).expanduser()
    return runtime_settings.cache_root / "scenes"


def cached_scene_directory(
    request: dict[str, object],
    cache_dir: str | Path | None = None,
    settings: Settings | None = None,
) -> Path:
    root = (
        Path(cache_dir).expanduser()
        if cache_dir is not None
        else default_scene_cache_root(settings)
    )
    latitude = float(request["latitude"])
    longitude = float(request["longitude"])
    prefix = f"scene_{latitude:+.5f}_{longitude:+.5f}".replace("+", "p").replace("-", "m")
    return root / f"{prefix}_{scene_cache_key(request)}"


def _write_json_atomic(path: Path, value: object) -> None:
    temporary = path.with_name(f".{path.name}.tmp-{os.getpid()}")
    temporary.write_text(json.dumps(value, indent=2), encoding="utf-8")
    temporary.replace(path)


@contextmanager
def _scene_lock(target: Path):
    target.parent.mkdir(parents=True, exist_ok=True)
    lock_path = target.parent / f".{target.name}.compile.lock"
    with lock_path.open("a+b") as handle:
        _acquire_file_lock(handle, blocking=True)
        try:
            yield
        finally:
            _release_file_lock(handle)


def _acquire_file_lock(handle, *, blocking: bool) -> None:
    if fcntl is not None:
        operation = fcntl.LOCK_EX if blocking else fcntl.LOCK_EX | fcntl.LOCK_NB
        fcntl.flock(handle.fileno(), operation)
        return
    handle.seek(0, os.SEEK_END)  # pragma: no cover - Windows only
    if handle.tell() == 0:  # pragma: no cover - Windows only
        handle.write(b"\0")
        handle.flush()
    handle.seek(0)
    mode = msvcrt.LK_LOCK if blocking else msvcrt.LK_NBLCK
    msvcrt.locking(handle.fileno(), mode, 1)


def _release_file_lock(handle) -> None:
    if fcntl is not None:
        fcntl.flock(handle.fileno(), fcntl.LOCK_UN)
        return
    handle.seek(0)  # pragma: no cover - Windows only
    msvcrt.locking(handle.fileno(), msvcrt.LK_UNLCK, 1)


@asynccontextmanager
async def _async_scene_lock(target: Path):
    target.parent.mkdir(parents=True, exist_ok=True)
    lock_path = target.parent / f".{target.name}.compile.lock"
    with lock_path.open("a+b") as handle:
        while True:
            try:
                _acquire_file_lock(handle, blocking=False)
                break
            except OSError as exc:
                if exc.errno not in {errno.EACCES, errno.EAGAIN}:
                    raise
                await asyncio.sleep(0.05)
        try:
            yield
        finally:
            _release_file_lock(handle)


def _write_building_source(path: Path, result: BuildingQueryResult) -> None:
    dump_json_with_compact_rows(
        path,
        {
            "provider": result.provider,
            "release": result.release,
            "warnings": result.warnings,
            "boundary_excluded_count": result.boundary_excluded_count,
            "source_selection": result.source_selection,
            "buildings": [building.model_dump(mode="json") for building in result.buildings],
        },
        "buildings",
    )


def _read_json(path: Path) -> dict[str, Any]:
    value = json.loads(path.read_text(encoding="utf-8"))
    if not isinstance(value, dict):
        raise ValueError(f"Expected a JSON object in {path}")
    return value


def resolve_scene_directory(path: str | Path) -> Path:
    candidate = Path(path).expanduser().resolve()
    directory = candidate if candidate.is_dir() else candidate.parent
    if not (directory / SCENE_XML_PATH).is_file():
        raise ValueError(f"No {SCENE_XML_PATH} found in {directory}")
    if not (directory / SCENE_PROVENANCE_PATH).is_file():
        raise ValueError(f"No {SCENE_PROVENANCE_PATH} found in {directory}")
    return directory


def _validate_scene_tree(directory: Path) -> None:
    required = (
        SCENE_XML_PATH,
        SCENE_PROVENANCE_PATH,
        SCENE_INFO_PATH,
        SCENE_SOURCE_BUILDINGS_PATH,
        SCENE_HEIGHT_MAP_PATH,
        SCENE_MESH_DIRECTORY,
        SCENE_CASES_DIRECTORY,
    )
    missing = [path.as_posix() for path in required if not (directory / path).exists()]
    if missing:
        raise RuntimeError(f"Scene cache is incomplete at {directory}: {', '.join(missing)}")
    provenance = _read_json(directory / SCENE_PROVENANCE_PATH)
    layout = provenance.get("scene_layout", {})
    if layout.get("format") != SCENE_FORMAT or layout.get("version") != SCENE_LAYOUT_VERSION:
        raise RuntimeError(f"Unsupported scene format in {directory / SCENE_PROVENANCE_PATH}")


def _atomic_install(staging: Path, target: Path) -> None:
    if not target.exists():
        staging.replace(target)
        return
    backup = target.parent / f".{target.name}.backup-{os.getpid()}"
    if backup.exists():
        shutil.rmtree(backup)
    target.replace(backup)
    try:
        staging.replace(target)
    except Exception:
        backup.replace(target)
        raise
    shutil.rmtree(backup)


def _write_compilation_metadata(
    directory: Path,
    scene_name: str,
    request: dict[str, object],
    buildings: BuildingQueryResult,
    environment_warnings: list[str],
    assets: SceneAssets,
) -> dict[str, Any]:
    provenance_path = directory / SCENE_PROVENANCE_PATH
    provenance = _read_json(provenance_path)
    provenance["compilation"] = {
        "compiler_version": __version__,
        "request": request,
        "cache_key": scene_cache_key(request),
        "building_provider": buildings.provider,
        "building_release": buildings.release,
        "source_building_count": buildings.source_count,
        "boundary_excluded_building_count": buildings.boundary_excluded_count,
        "exported_building_count": assets.building_count,
        "scene_boundary_policy": SCENE_BOUNDARY_POLICY,
        "warnings": [*buildings.warnings, *environment_warnings],
    }
    dump_json_with_compact_rows(provenance_path, provenance, "local_buildings")
    shared_assets: dict[str, object] = {
        "scene_xml": SCENE_XML_PATH.as_posix(),
        "coordinate_provenance": SCENE_PROVENANCE_PATH.as_posix(),
        "mesh": f"{SCENE_MESH_DIRECTORY.as_posix()}/",
        "source_buildings": SCENE_SOURCE_BUILDINGS_PATH.as_posix(),
        "height_map": SCENE_HEIGHT_MAP_PATH.as_posix(),
        "measurement_surfaces": [],
    }
    if assets.environment is not None:
        shared_assets["environment"] = SCENE_TERRAIN_INFO_PATH.as_posix()
    scene_info: dict[str, Any] = {
        "schema_version": 1,
        "name": scene_name,
        "label": scene_name,
        "status": "compiled",
        "specification": request,
        "provenance": provenance["compilation"],
        "shared_assets": shared_assets,
        "derived_assets": {
            "measurement_surfaces_complete": False,
            "voxels_complete": False,
        },
        "cases": [],
    }
    _write_json_atomic(directory / SCENE_INFO_PATH, scene_info)
    return scene_info


def _measurement_record(
    path: Path,
    cell_shape: tuple[int, int],
    directory: Path,
    cell_size: CellSize,
    height_m: float,
) -> dict:
    east_m, north_m = _normalize_cell_size(cell_size)
    return {
        "path": str(path.relative_to(directory)),
        "metadata": str(path.with_suffix(".json").relative_to(directory)),
        "arrays": str(path.with_suffix(".npz").relative_to(directory)),
        "cell_size_m": [east_m, north_m],
        "receiver_height_agl_m": height_m,
        "cell_shape": list(cell_shape),
        "height_reference": "terrain_agl",
    }


def _normalize_cell_size(value: CellSize) -> tuple[float, float]:
    if isinstance(value, tuple):
        if len(value) != 2:
            raise ValueError("receiver_cell_size must contain east and north sizes")
        east_m, north_m = value
    else:
        east_m = north_m = value
    if east_m <= 0 or north_m <= 0:
        raise ValueError("receiver_cell_size must be positive")
    return float(east_m), float(north_m)


def _add_derived_assets(
    directory: Path,
    assets: SceneAssets,
    receiver_height_m: float | None,
    receiver_cell_size: CellSize | None,
    voxel_pitch_m: float | None,
) -> None:
    scene_info_path = directory / SCENE_INFO_PATH
    scene_info = _read_json(scene_info_path)
    derived = scene_info.setdefault("derived_assets", {})
    shared = scene_info.setdefault("shared_assets", {})
    if receiver_height_m is not None and receiver_cell_size is not None:
        surface_path = measurement_surface_path(assets, receiver_cell_size, receiver_height_m)
        required = (
            surface_path,
            surface_path.with_suffix(".json"),
            surface_path.with_suffix(".npz"),
        )
        current_surface = all(path.is_file() for path in required)
        if current_surface:
            metadata = _read_json(surface_path.with_suffix(".json"))
            current_surface = metadata.get("schema_version") == MEASUREMENT_SURFACE_SCHEMA_VERSION
        if not current_surface:
            temporary_root = Path(tempfile.mkdtemp(prefix=".measurement-", dir=directory))
            try:
                temporary_assets = replace(assets, directory=temporary_root)
                temporary_surface = write_measurement_surface(
                    measurement_surface_path(
                        temporary_assets,
                        receiver_cell_size,
                        receiver_height_m,
                    ),
                    temporary_assets,
                    receiver_cell_size,
                    receiver_height_m,
                )
                surface_path.parent.mkdir(exist_ok=True)
                for temporary_path in (
                    temporary_surface.path,
                    temporary_surface.path.with_suffix(".json"),
                    temporary_surface.path.with_suffix(".npz"),
                ):
                    temporary_path.replace(surface_path.parent / temporary_path.name)
            finally:
                shutil.rmtree(temporary_root, ignore_errors=True)
        metadata = _read_json(surface_path.with_suffix(".json"))
        cell_shape = tuple(int(value) for value in metadata["cell_shape"])
        if len(cell_shape) != 2:
            metadata_path = surface_path.with_suffix(".json")
            raise RuntimeError(f"Invalid measurement cell shape in {metadata_path}")
        records = list(shared.get("measurement_surfaces", []))
        record = _measurement_record(
            surface_path,
            cell_shape,
            directory,
            receiver_cell_size,
            receiver_height_m,
        )
        records = [item for item in records if item.get("path") != record["path"]]
        shared["measurement_surfaces"] = [*records, record]
        derived["measurement_surfaces_complete"] = True
    if voxel_pitch_m is not None and receiver_height_m is not None:
        pitch_token = f"{voxel_pitch_m:.1f}"
        receiver_token = str(receiver_height_m).rstrip("0").rstrip(".")
        voxel_directory = directory / SCENE_VOXEL_SLICES_DIRECTORY
        depth_directory = directory / SCENE_VOXEL_DEPTH_DIRECTORY
        voxel_path = voxel_directory / f"scene_voxels_pitch{pitch_token}.npz"
        slice_path = voxel_directory / f"scene_slice_{receiver_token}m_pitch{pitch_token}.png"
        depth_path = depth_directory / f"scene_depth_max_z_pitch{pitch_token}.npy"
        required = (
            voxel_path,
            slice_path,
            voxel_directory / SCENE_INFO_PATH,
            depth_path,
            depth_directory / SCENE_INFO_PATH,
        )
        if not all(path.is_file() for path in required):
            temporary_root = Path(tempfile.mkdtemp(prefix=".voxels-", dir=directory))
            try:
                temporary_mesh = temporary_root / SCENE_MESH_DIRECTORY
                temporary_assets = replace(
                    assets,
                    directory=temporary_root,
                    building_assets=tuple(
                        replace(
                            item,
                            rooftop_path=temporary_mesh / item.rooftop_path.name,
                            wall_path=temporary_mesh / item.wall_path.name,
                        )
                        for item in assets.building_assets
                    ),
                )
                generated = write_voxel_assets(
                    temporary_assets,
                    receiver_height_m,
                    voxel_pitch_m,
                )
                generated_files = (
                    generated.voxel_path,
                    generated.slice_path,
                    generated.scene_info_path,
                    generated.depth_array_path,
                    generated.depth_image_path,
                    generated.depth_preview_path,
                    generated.depth_info_path,
                )
                for source in generated_files:
                    target = directory / source.relative_to(temporary_root)
                    target.parent.mkdir(parents=True, exist_ok=True)
                    source.replace(target)
            finally:
                shutil.rmtree(temporary_root, ignore_errors=True)
        shared.update(
            {
                "voxels": str(voxel_path.relative_to(directory)),
                "receiver_slice": str(slice_path.relative_to(directory)),
                "voxel_depth": str(depth_path.relative_to(directory)),
            }
        )
        derived["voxels_complete"] = True
    _write_json_atomic(scene_info_path, scene_info)


async def _acquire_inputs(
    request: dict[str, object],
    provider: MapDataProvider,
    settings: Settings,
) -> tuple[BuildingQueryResult, EnvironmentData | None, list[str]]:
    """Query the geometry inputs of a scene through one shared provider.

    Buildings and the environment grid come from independent services, so both
    queries run at the same time.
    """

    async def buildings() -> BuildingQueryResult:
        if request["buildings"] != "auto":
            return BuildingQueryResult([], "disabled")
        return await provider.buildings(
            float(request["latitude"]),
            float(request["longitude"]),
            int(float(request["radius_m"])),
            settings.max_buildings,
            source=str(request["building_source"]),
        )

    async def environment() -> tuple[EnvironmentData | None, list[str]]:
        include_surfaces = request["material_profile"] == "itu"
        if request["terrain"] != "elevation" and not (
            request["terrain"] == "flat" and include_surfaces
        ):
            return None, []
        return await provider.environment(
            float(request["latitude"]),
            float(request["longitude"]),
            int(float(request["radius_m"])),
            float(request["terrain_resolution_m"]),
            include_surfaces,
            terrain_mode=str(request["terrain"]),
        )

    building_result, (environment_data, environment_warnings) = await asyncio.gather(
        buildings(), environment()
    )
    return building_result, environment_data, environment_warnings


[docs] def load_building_source(path: Path) -> BuildingQueryResult: """Read the provider records persisted next to a compiled scene.""" value = _read_json(path) return BuildingQueryResult( buildings=[BuildingFeature.model_validate(item) for item in value["buildings"]], provider=str(value["provider"]), release=value.get("release"), warnings=[str(item) for item in value.get("warnings", [])], boundary_excluded_count=int(value.get("boundary_excluded_count", 0)), source_selection=value.get("source_selection", {}), )
async def _compile_scene( *, request: dict[str, object], output: Path, receiver_height_m: float | None, receiver_cell_size: CellSize | None, voxel_pitch_m: float | None, if_exists: ExistingPolicy, settings: Settings, provider: MapDataProvider | None = None, ) -> Path: if if_exists not in {"reuse", "error", "replace"}: raise ValueError(f"Unknown existing-scene policy: {if_exists}") output = output.expanduser().resolve() async with _async_scene_lock(output): if output.exists() and if_exists == "error": raise FileExistsError(f"Scene already exists at {output}") if output.exists() and if_exists == "reuse": try: _validate_scene_tree(output) except RuntimeError as exc: raise RuntimeError( f"Cannot reuse the incomplete scene at {output}; rebuild this exact target " "with if_exists='replace'" ) from exc info = _read_json(output / SCENE_INFO_PATH) if info.get("specification") != request: raise ValueError( f"Existing scene at {output} has different geometry inputs; " "use if_exists='replace' or choose another output" ) if receiver_height_m is not None or voxel_pitch_m is not None: assets = load_scene_assets(output) _add_derived_assets( output, assets, receiver_height_m, receiver_cell_size, voxel_pitch_m, ) return output / SCENE_XML_PATH staging = Path(tempfile.mkdtemp(prefix=f".{output.name}.tmp-", dir=output.parent)) try: if provider is not None: building_result, environment, environment_warnings = await _acquire_inputs( request, provider, settings ) else: async with httpx.AsyncClient( follow_redirects=True, timeout=settings.http_request_timeout_seconds, ) as client: building_result, environment, environment_warnings = await _acquire_inputs( request, MapDataProvider(settings, client), settings ) # Meshing and metadata are CPU work; keep the event loop free for the service. await asyncio.to_thread( _build_scene_tree, staging, output.name, request, building_result, environment, environment_warnings, receiver_height_m, receiver_cell_size, voxel_pitch_m, ) _atomic_install(staging, output) finally: if staging.exists(): shutil.rmtree(staging, ignore_errors=True) return output / SCENE_XML_PATH def _build_scene_tree( staging: Path, scene_name: str, request: dict[str, object], building_result: BuildingQueryResult, environment: EnvironmentData | None, environment_warnings: list[str], receiver_height_m: float | None, receiver_cell_size: CellSize | None, voxel_pitch_m: float | None, ) -> None: """Write the complete scene directory from acquired inputs.""" _write_building_source(staging / SCENE_SOURCE_BUILDINGS_PATH, building_result) assets = build_scene_assets( staging, building_result.buildings, float(request["longitude"]), float(request["latitude"]), float(request["radius_m"]), environment, request["terrain"] != "none", request["material_profile"], ) (staging / SCENE_CASES_DIRECTORY).mkdir(exist_ok=True) _write_compilation_metadata( staging, scene_name, request, building_result, environment_warnings, assets ) _add_derived_assets(staging, assets, receiver_height_m, receiver_cell_size, voxel_pitch_m) _validate_scene_tree(staging)
[docs] async def compile_scene( *, latitude: float, longitude: float, radius_m: float, output: str | Path, buildings: BuildingMode = "auto", building_source: str = "auto", terrain: TerrainMode = "elevation", material_profile: MaterialProfileName = "itu", terrain_resolution_m: float = 5.0, receiver_height_m: float = 1.5, receiver_cell_size: CellSize = (10.0, 10.0), voxel_pitch_m: float = 3.0, if_exists: ExistingPolicy = "reuse", settings: Settings | None = None, ) -> Path: """Compile a portable scene and return its Sionna XML path.""" runtime_settings = settings or get_settings() request = _geometry_request( latitude, longitude, radius_m, buildings, terrain, material_profile, terrain_resolution_m, building_source, runtime_settings, ) return await _compile_scene( request=request, output=Path(output), receiver_height_m=receiver_height_m, receiver_cell_size=receiver_cell_size, voxel_pitch_m=voxel_pitch_m, if_exists=if_exists, settings=runtime_settings, )
[docs] async def compile_cached_scene( *, latitude: float, longitude: float, radius_m: float, cache_dir: str | Path | None, buildings: BuildingMode, terrain: TerrainMode, material_profile: MaterialProfileName, terrain_resolution_m: float, building_source: str = "auto", settings: Settings | None = None, provider: MapDataProvider | None = None, ) -> Path: """Compile into the shared cache, or reuse the directory for identical geometry inputs. ``provider`` lets a running service reuse its own HTTP client and in-memory caches; the loader and CLI leave it unset and open a client for the compilation. """ runtime_settings = settings or get_settings() request = _geometry_request( latitude, longitude, radius_m, buildings, terrain, material_profile, terrain_resolution_m, building_source, runtime_settings, ) output = cached_scene_directory(request, cache_dir, runtime_settings) return await _compile_scene( request=request, output=output, receiver_height_m=None, receiver_cell_size=None, voxel_pitch_m=None, if_exists="reuse", settings=runtime_settings, provider=provider, )
def _run_sync(coroutine: Coroutine[Any, Any, T]) -> T: try: asyncio.get_running_loop() except RuntimeError: return asyncio.run(coroutine) with ThreadPoolExecutor(max_workers=1, thread_name_prefix="owrt-compiler") as worker: return worker.submit(asyncio.run, coroutine).result()
[docs] def compile_scene_sync(**kwargs: Any) -> Path: """Synchronous adapter for :func:`compile_scene`, including notebooks.""" return _run_sync(compile_scene(**kwargs))