Source code for openworld_radio_twin.batch

import asyncio
import fcntl
import json
import math
import re
import time
import unicodedata
from collections.abc import Callable
from concurrent.futures import Executor
from contextlib import contextmanager
from datetime import datetime, timezone
from functools import partial
from pathlib import Path
from typing import Literal

import numpy as np
from pydantic import BaseModel, Field, field_validator, model_validator

from openworld_radio_twin import __version__
from openworld_radio_twin.artifacts import write_radio_products
from openworld_radio_twin.compiler import load_building_source
from openworld_radio_twin.geodesy import LocalFrame
from openworld_radio_twin.models import (
    MAX_PREVIEW_GRID_CELLS,
    MAX_SIONNA_GRID_CELLS,
    MAX_TERRAIN_GRID_CELLS,
    Coordinate,
    EngineName,
    SimulationRequest,
    SimulationResponse,
    SionnaEngineConfig,
    TransmitterConfig,
)
from openworld_radio_twin.providers.base import SCENE_BOUNDARY_POLICY, BuildingQueryResult
from openworld_radio_twin.providers.building_sources import resolve_building_source
from openworld_radio_twin.scene_format import (
    SCENE_CASES_DIRECTORY,
    SCENE_ENVIRONMENT_ARRAY_PATHS,
    SCENE_HEIGHT_MAP_PATH,
    SCENE_INFO_PATH,
    SCENE_MESH_DIRECTORY,
    SCENE_PROVENANCE_PATH,
    SCENE_SOURCE_BUILDINGS_PATH,
    SCENE_TERRAIN_INFO_PATH,
    SCENE_XML_PATH,
)
from openworld_radio_twin.simulation.preview import simulate_preview
from openworld_radio_twin.simulation.radio_map import build_coverage_products
from openworld_radio_twin.simulation.scene_builder import (
    build_scene_assets,
    load_environment_assets,
    measurement_surface_path,
    write_measurement_surface,
    write_voxel_assets,
)
from openworld_radio_twin.simulation.sionna_engine import (
    SIONNA_RADIO_MAP_CONTRACT_VERSION,
    simulate_sionna_scene,
)

ProgressCallback = Callable[[int, int, str], None]
MAX_CASES_PER_SCENE = 10_000
RADIO_METRICS = ("path_gain", "rss", "sinr")
RadioMetricName = Literal["path_gain", "rss", "sinr"]
RADIO_METRIC_ARTIFACTS = {
    "path_gain": {
        "array": "arrays/path_gain/path_gain.npz",
        "images": ("images/path_gain_db.png", "images/path_gain_db_masked.png"),
        "unit": "linear ratio",
    },
    "rss": {
        "array": "arrays/rss/rss.npz",
        "images": ("images/rss_db.png", "images/rss_db_masked.png"),
        "unit": "W",
    },
    "sinr": {
        "array": "arrays/sinr/sinr.npz",
        "images": ("images/sinr.png", "images/sinr_masked.png"),
        "unit": "linear ratio",
    },
}


[docs] class FloatRange(BaseModel): minimum: float maximum: float
[docs] @model_validator(mode="after") def ordered(self) -> "FloatRange": if self.maximum < self.minimum: raise ValueError("maximum must be greater than or equal to minimum") return self
[docs] class IntegerRange(BaseModel): minimum: int maximum: int
[docs] @model_validator(mode="after") def ordered(self) -> "IntegerRange": if self.maximum < self.minimum: raise ValueError("maximum must be greater than or equal to minimum") return self
[docs] class BatchOutputSpec(BaseModel): """Case artifacts retained after a successful simulation. Case metadata is mandatory because it is the provenance and resume-completion record. Empty metric lists therefore provide a metadata-only output profile. """ array_metrics: list[RadioMetricName] = Field(default_factory=lambda: list(RADIO_METRICS)) image_metrics: list[RadioMetricName] = Field(default_factory=lambda: list(RADIO_METRICS)) association_image: bool = True
[docs] @field_validator("array_metrics", "image_metrics") @classmethod def unique_ordered_metrics(cls, metrics: list[RadioMetricName]) -> list[RadioMetricName]: if len(set(metrics)) != len(metrics): raise ValueError("metrics must be unique") selected = set(metrics) return [metric for metric in RADIO_METRICS if metric in selected]
[docs] class BatchCaseSpace(BaseModel): strategy: Literal["fixed", "random"] = "random" count: int = Field(default=8, ge=1, le=MAX_CASES_PER_SCENE) engine: EngineName = EngineName.preview engine_config: dict[str, str | int | float | bool] = Field(default_factory=dict) transmitter_count: IntegerRange = IntegerRange(minimum=1, maximum=3) frequency_ghz: FloatRange = FloatRange(minimum=3.5, maximum=3.5) power_dbm: FloatRange = FloatRange(minimum=24, maximum=36) altitude_m: FloatRange = FloatRange(minimum=15, maximum=35) offset_radius_m: FloatRange = FloatRange(minimum=0, maximum=250) azimuth_deg: FloatRange = FloatRange(minimum=0, maximum=359) downtilt_deg: FloatRange = FloatRange(minimum=0, maximum=12) antenna_patterns: list[Literal["sector", "isotropic"]] = Field( default_factory=lambda: ["sector"] ) resolutions_m: list[int] = Field(default_factory=lambda: [15]) max_depths: list[int] = Field(default_factory=lambda: [3]) samples_per_tx: list[int] = Field(default_factory=lambda: [1_000_000]) association_metrics: list[Literal["path_gain", "rss", "sinr"]] = Field( default_factory=lambda: ["sinr"] )
[docs] @model_validator(mode="after") def validate_space(self) -> "BatchCaseSpace": if not 1 <= self.transmitter_count.minimum <= self.transmitter_count.maximum <= 8: raise ValueError("transmitter count must stay between 1 and 8") if not self.antenna_patterns: raise ValueError("at least one antenna pattern is required") if not self.resolutions_m or any(not 1 <= value <= 200 for value in self.resolutions_m): raise ValueError("resolutions must stay between 1 and 200 metres") if not self.max_depths or any(not 0 <= value <= 8 for value in self.max_depths): raise ValueError("reflection depths must stay between 0 and 8") if not self.samples_per_tx or any( not 10_000 <= value <= 100_000_000 for value in self.samples_per_tx ): raise ValueError("ray samples must stay between 10,000 and 100,000,000") if not self.association_metrics: raise ValueError("at least one association metric is required") ranges = { "frequency": (self.frequency_ghz, 0.1, 100), "power": (self.power_dbm, -20, 80), "altitude": (self.altitude_m, 0, 1000), "TX offset": (self.offset_radius_m, 0, 5000), "azimuth": (self.azimuth_deg, 0, 359.999999), "downtilt": (self.downtilt_deg, -30, 90), } for name, (values, lower, upper) in ranges.items(): if values.minimum < lower or values.maximum > upper: raise ValueError(f"{name} range must stay between {lower:g} and {upper:g}") if self.engine == EngineName.sionna and ( self.frequency_ghz.minimum < 1 or self.frequency_ghz.maximum > 10 ): raise ValueError("Sionna frequency range must stay between 1 and 10 GHz") return self
[docs] class BatchSceneSpec(BaseModel): label: str = Field(min_length=1, max_length=80) scene_mode: Literal["geospatial", "empty"] = "geospatial" latitude: float = Field(ge=-90, le=90) longitude: float = Field(ge=-180, le=180) radius_m: int = Field(default=750, ge=100, le=5000) receiver_height_m: float = Field(default=1.5, ge=0.1, le=100) voxel_pitch_m: float = Field(default=3.0, ge=0.5, le=25) include_buildings: bool = True building_source: str = Field(default="auto", pattern=r"^(?:[a-z][a-z0-9_-]*|3dbag)$") include_terrain: bool = False material_profile: Literal["itu", "uniform"] = "itu" terrain_resolution_m: float = Field(default=5.0, ge=1.0, le=100.0) cases: BatchCaseSpace = Field(default_factory=BatchCaseSpace)
[docs] @model_validator(mode="after") def fit_transmitters_in_scene(self) -> "BatchSceneSpec": if self.scene_mode == "empty": self.include_buildings = False self.include_terrain = False self.material_profile = "uniform" terrain_side = math.ceil(2 * self.radius_m / self.terrain_resolution_m) if self.include_terrain and terrain_side**2 > MAX_TERRAIN_GRID_CELLS: raise ValueError( f"terrain grid exceeds {MAX_TERRAIN_GRID_CELLS:,} cells; " "increase terrain resolution" ) if self.cases.offset_radius_m.maximum > self.radius_m: raise ValueError("maximum TX offset must not exceed the scene radius") maximum = ( MAX_SIONNA_GRID_CELLS if self.cases.engine == EngineName.sionna else MAX_PREVIEW_GRID_CELLS ) for resolution in self.cases.resolutions_m: side = math.ceil(2 * self.radius_m / resolution) if side**2 > maximum: raise ValueError( f"{resolution} m resolution creates a {side}x{side} grid, exceeding " f"the {maximum:,}-cell {self.cases.engine.value} limit" ) return self
@property def include_surface_materials(self) -> bool: """Return whether semantic surface classes are required by this scene.""" return self.material_profile == "itu" @property def include_environment_grid(self) -> bool: """Return whether a gridded ground (DEM or flat) is compiled for this scene.""" return self.include_terrain or self.include_surface_materials
[docs] class BatchDatasetRequest(BaseModel): dataset_name: str = Field(min_length=1, max_length=80) random_seed: int = Field(default=42, ge=0, le=4_294_967_295) continue_on_error: bool = False output: BatchOutputSpec = Field(default_factory=BatchOutputSpec) scenes: list[BatchSceneSpec] = Field(min_length=1, max_length=20)
[docs] class ExpandedCase(BaseModel): name: str seed: int request: SimulationRequest
[docs] class PlannedScene(BaseModel): name: str label: str specification: BatchSceneSpec cases: list[ExpandedCase]
[docs] class BatchPlan(BaseModel): schema_version: int = 1 dataset_name: str dataset_slug: str random_seed: int generator: str = "numpy.random.PCG64" total_scenes: int total_cases: int total_receiver_cells: int total_initial_ray_samples: int scenes: list[PlannedScene]
[docs] def slugify(value: str, fallback: str) -> str: ascii_value = unicodedata.normalize("NFKD", value).encode("ascii", "ignore").decode() slug = re.sub(r"[^a-zA-Z0-9]+", "-", ascii_value).strip("-").lower() return (slug or fallback)[:48]
def _sample_float(rng: np.random.Generator, values: FloatRange, random: bool) -> float: if not random or values.minimum == values.maximum: return values.minimum return float(rng.uniform(values.minimum, values.maximum)) def _sample_int(rng: np.random.Generator, values: IntegerRange, random: bool) -> int: if not random or values.minimum == values.maximum: return values.minimum return int(rng.integers(values.minimum, values.maximum + 1)) def _choice(rng: np.random.Generator, values: list, random: bool): return values[int(rng.integers(0, len(values)))] if random else values[0] def _expand_case( dataset_seed: int, scene_index: int, case_index: int, scene: BatchSceneSpec, ) -> ExpandedCase: space = scene.cases engine_config = dict(space.engine_config) if scene.scene_mode == "empty" and space.engine == EngineName.sionna: engine_config.update( { "include_ground": False, "include_vegetation": False, "include_paved": False, "include_water": False, "include_building_roofs": False, "include_building_walls": False, } ) random = space.strategy == "random" rng = np.random.default_rng(np.random.SeedSequence([dataset_seed, scene_index, case_index])) case_seed = int(rng.integers(0, 2**32, dtype=np.uint64)) tx_count = _sample_int(rng, space.transmitter_count, random) frequency = round(_sample_float(rng, space.frequency_ghz, random), 6) pattern = _choice(rng, space.antenna_patterns, random) frame = LocalFrame.at(scene.longitude, scene.latitude) transmitters = [] for tx_index in range(tx_count): east = north = 0.0 if tx_index: inner = space.offset_radius_m.minimum outer = space.offset_radius_m.maximum radial = inner if not random else math.sqrt(rng.uniform(inner**2, outer**2)) angle = 0.0 if not random else rng.uniform(0.0, 2.0 * math.pi) east, north = radial * math.sin(angle), radial * math.cos(angle) longitude, latitude, _ = frame.to_geographic(east, north, 0.0) transmitters.append( TransmitterConfig( position=Coordinate( latitude=latitude, longitude=longitude, altitude_m=round(_sample_float(rng, space.altitude_m, random), 6), ), frequency_ghz=frequency, power_dbm=round(_sample_float(rng, space.power_dbm, random), 6), azimuth_deg=round(_sample_float(rng, space.azimuth_deg, random), 6), downtilt_deg=round(_sample_float(rng, space.downtilt_deg, random), 6), antenna_pattern=pattern, ) ) request = SimulationRequest( transmitters=transmitters, engine=space.engine, engine_config=engine_config, radius_m=scene.radius_m, resolution_m=int(_choice(rng, space.resolutions_m, random)), receiver_height_m=scene.receiver_height_m, max_depth=int(_choice(rng, space.max_depths, random)), samples_per_tx=int(_choice(rng, space.samples_per_tx, random)), seed=case_seed, association_metric=_choice(rng, space.association_metrics, random), include_buildings=scene.include_buildings, building_source=scene.building_source, include_terrain=scene.include_terrain, material_profile=scene.material_profile, terrain_resolution_m=scene.terrain_resolution_m, ) frequency_token = f"{frequency:g}".replace(".", "p") return ExpandedCase( name=f"case_{case_index:04d}_{pattern}_{tx_count}tx_{frequency_token}ghz", seed=case_seed, request=request, )
[docs] def build_batch_plan(payload: BatchDatasetRequest) -> BatchPlan: scenes = [] total_cells = 0 total_samples = 0 for scene_index, specification in enumerate(payload.scenes, start=1): cases = [ _expand_case(payload.random_seed, scene_index, index, specification) for index in range(1, specification.cases.count + 1) ] name = f"scene_{scene_index:03d}_{slugify(specification.label, 'scene')}" scenes.append( PlannedScene( name=name, label=specification.label, specification=specification, cases=cases, ) ) for case in cases: side = math.ceil(2 * specification.radius_m / case.request.resolution_m) total_cells += side**2 if case.request.engine == EngineName.sionna.value: total_samples += case.request.samples_per_tx * len(case.request.transmitters) return BatchPlan( dataset_name=payload.dataset_name, dataset_slug=slugify(payload.dataset_name, "dataset"), random_seed=payload.random_seed, total_scenes=len(scenes), total_cases=sum(len(scene.cases) for scene in scenes), total_receiver_cells=total_cells, total_initial_ray_samples=total_samples, scenes=scenes, )
def _write_json(path: Path, value: object) -> None: temporary = path.with_suffix(f"{path.suffix}.tmp") temporary.write_text(json.dumps(value, indent=2), encoding="utf-8") temporary.replace(path) def _read_json(path: Path) -> dict: 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 _write_building_source(path: Path, result: BuildingQueryResult) -> None: _write_json( 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], }, ) def _resume_contract(payload: BatchDatasetRequest) -> dict: value = payload.model_dump(mode="json") value.pop("continue_on_error", None) for scene in value["scenes"]: scene["cases"].pop("count", None) if scene["cases"]["engine"] == EngineName.sionna.value: scene["cases"]["engine_config"] = SionnaEngineConfig.model_validate( scene["cases"]["engine_config"] ).model_dump(mode="json") return value def _validate_resume_manifest( manifest: dict, payload: BatchDatasetRequest, plan: BatchPlan, ) -> BatchDatasetRequest: if manifest.get("dataset_slug") != plan.dataset_slug: raise ValueError("Existing manifest dataset identity does not match the requested dataset") software_version = manifest.get("software", {}).get("openworld_radio_twin") preview_patch_compatible = ( software_version == "0.2.0" and __version__ == "0.2.1" and all(scene.cases.engine == EngineName.preview.value for scene in payload.scenes) ) if software_version != __version__ and not preview_patch_compatible: raise ValueError( "Resume requires the same OpenWorld Radio Twin version that created the dataset " f"({software_version!r} != {__version__!r})" ) try: previous = BatchDatasetRequest.model_validate(manifest["request"]) except Exception as exc: raise ValueError("Existing manifest does not contain a valid resumable request") from exc if _resume_contract(previous) != _resume_contract(payload): raise ValueError( "Resume may only increase per-scene case counts; seeds, scenes, engines, ranges, " "geometry, and solver settings must match the existing dataset" ) if any(scene.cases.engine == EngineName.sionna.value for scene in payload.scenes): previous_contract = manifest.get("execution_contracts", {}).get("sionna_radio_map", 1) if previous_contract != SIONNA_RADIO_MAP_CONTRACT_VERSION: raise ValueError( "Cannot resume Sionna data with a different radio-map contract " f"({previous_contract} != {SIONNA_RADIO_MAP_CONTRACT_VERSION}). " "SINR is now computed from cell-averaged RSS, not averaged triangle SINR. " "Use a new dataset directory or the original generator for this dataset." ) for index, (old_scene, new_scene) in enumerate( zip(previous.scenes, payload.scenes, strict=True), start=1 ): if new_scene.cases.count < old_scene.cases.count: raise ValueError( f"Scene {index} case count cannot decrease from " f"{old_scene.cases.count} to {new_scene.cases.count} during resume" ) return previous def _required_completed_case_files(output: BatchOutputSpec) -> tuple[str, ...]: required = [str(RADIO_METRIC_ARTIFACTS[metric]["array"]) for metric in output.array_metrics] for metric in output.image_metrics: required.extend(RADIO_METRIC_ARTIFACTS[metric]["images"]) if output.association_image: required.append("images/association.png") return tuple(required) def _completed_case_record( case_directory: Path, case: ExpandedCase, output: BatchOutputSpec, ) -> dict | None: info_path = case_directory / "case_info.json" if not info_path.is_file(): return None try: record = _read_json(info_path) except (OSError, ValueError, json.JSONDecodeError): return None recorded_request = record.get("request") expected_request = case.request.model_dump(mode="json") if recorded_request is None: return None if recorded_request != expected_request: raise ValueError( f"Existing case {case_directory.name} was generated from an incompatible request" ) if record.get("status") != "completed": return None required_files = _required_completed_case_files(output) if any(not (case_directory / relative).is_file() for relative in required_files): return None return record def _scan_completed_cases( dataset_directory: Path, plan: BatchPlan, output: BatchOutputSpec, ) -> dict[tuple[str, str], dict]: completed: dict[tuple[str, str], dict] = {} for scene in plan.scenes: cases_directory = dataset_directory / scene.name / SCENE_CASES_DIRECTORY for case in scene.cases: record = _completed_case_record(cases_directory / case.name, case, output) if record is not None: completed[(scene.name, case.name)] = record return completed @contextmanager def _exclusive_dataset_lock(dataset_directory: Path): lock_path = dataset_directory / ".generation.lock" with lock_path.open("a+", encoding="utf-8") as handle: try: fcntl.flock(handle.fileno(), fcntl.LOCK_EX | fcntl.LOCK_NB) except BlockingIOError as exc: raise RuntimeError( f"Dataset generation is already active for {dataset_directory}" ) from exc try: yield finally: fcntl.flock(handle.fileno(), fcntl.LOCK_UN) def _preview_details(request: SimulationRequest, building_count: int, cells: int) -> dict: return { "device": "CPU", "model": "FSPL with directional and building obstruction penalties", "buildings_considered": building_count, "transmitter_count": len(request.transmitters), "array_contract": "linear_float32_(n_tx,rows,columns)", "cell_size_m": request.resolution_m, "receiver_cells": cells, "solver_seed": request.seed, "association_metric": request.association_metric, "terrain_propagation": False, } async def _run_in_worker(executor: Executor, function, *args): future = executor.submit(partial(function, *args)) while not future.done(): await asyncio.sleep(0.05) return future.result()
[docs] async def generate_batch_dataset( payload: BatchDatasetRequest, plan: BatchPlan, output_root: Path, provider, building_limit: int, sionna_lock: asyncio.Lock, worker: Executor, progress: ProgressCallback, resume: bool = False, ) -> Path: dataset_directory = output_root.expanduser().resolve() / plan.dataset_slug output_root.expanduser().resolve().mkdir(parents=True, exist_ok=True) dataset_directory.mkdir(parents=False, exist_ok=resume) with _exclusive_dataset_lock(dataset_directory): try: return await _generate_batch_dataset_locked( payload, plan, dataset_directory, provider, building_limit, sionna_lock, worker, progress, resume, ) except Exception as exc: manifest_path = dataset_directory / "dataset_manifest.json" if manifest_path.is_file(): manifest = _read_json(manifest_path) runs = manifest.get("generation_runs", []) if manifest.get("status") == "running" and runs and runs[-1]["status"] == "running": completed_at = datetime.now(timezone.utc).isoformat() manifest["status"] = "failed" manifest["completed_at"] = completed_at runs[-1].update( status="failed", completed_at=completed_at, error=str(exc), ) _write_json(manifest_path, manifest) raise
async def _generate_batch_dataset_locked( payload: BatchDatasetRequest, plan: BatchPlan, dataset_directory: Path, provider, building_limit: int, sionna_lock: asyncio.Lock, worker: Executor, progress: ProgressCallback, resume: bool, ) -> Path: started_at = datetime.now(timezone.utc).isoformat() manifest_path = dataset_directory / "dataset_manifest.json" existing_manifest = manifest_path.is_file() if existing_manifest: if not resume: raise FileExistsError(f"Dataset already exists at {dataset_directory}") manifest = _read_json(manifest_path) _validate_resume_manifest(manifest, payload, plan) else: unexpected = [ path.name for path in dataset_directory.iterdir() if path.name != ".generation.lock" ] if unexpected: raise ValueError( f"Cannot resume {dataset_directory}: dataset_manifest.json is missing from a " "non-empty directory" ) manifest = { "schema_version": plan.schema_version, "dataset_name": plan.dataset_name, "dataset_slug": plan.dataset_slug, "status": "running", "created_at": started_at, "completed_at": None, "software": {"openworld_radio_twin": __version__}, "execution_contracts": {"sionna_radio_map": SIONNA_RADIO_MAP_CONTRACT_VERSION}, "randomness": { "dataset_seed": plan.random_seed, "generator": plan.generator, "derivation": "SeedSequence([dataset_seed, one_based_scene, one_based_case])", "expanded_parameters_recorded": True, }, "request": payload.model_dump(mode="json"), "summary": {}, "scenes": [], "generation_runs": [], } completed_records = _scan_completed_cases(dataset_directory, plan, payload.output) existing_scenes = { item["name"]: item for item in manifest.get("scenes", []) if isinstance(item, dict) and "name" in item } manifest["scenes"] = [ existing_scenes.get(scene.name, {"name": scene.name, "cases": []}) for scene in plan.scenes ] manifest["request"] = payload.model_dump(mode="json") manifest["status"] = "running" manifest["completed_at"] = None manifest["summary"] = { "scenes": plan.total_scenes, "cases": plan.total_cases, "completed_cases": len(completed_records), "failed_cases": 0, } generation_runs = manifest.setdefault("generation_runs", []) run_record = { "started_at": started_at, "software": {"openworld_radio_twin": __version__}, "completed_at": None, "resume": existing_manifest, "requested_cases": plan.total_cases, "initial_completed_cases": len(completed_records), "generated_cases": 0, "failed_cases": 0, "status": "running", } generation_runs.append(run_record) _write_json(manifest_path, manifest) progress( len(completed_records), plan.total_cases, f"Found {len(completed_records)} completed case(s); preparing shared scene assets", ) processed = len(completed_records) for scene_index, scene in enumerate(plan.scenes): specification = scene.specification scene_directory = dataset_directory / scene.name scene_directory.mkdir(exist_ok=True) scene_has_completed_cases = any(key[0] == scene.name for key in completed_records) building_source_path = scene_directory / SCENE_SOURCE_BUILDINGS_PATH if building_source_path.is_file(): try: building_result = load_building_source(building_source_path) except (OSError, KeyError, TypeError, ValueError, json.JSONDecodeError): if scene_has_completed_cases: raise ValueError( f"Scene {scene.name} has invalid persisted building inputs required " "for a deterministic resume" ) from None building_result = None else: building_result = None if building_result is not None and specification.include_buildings: expected_source = resolve_building_source( specification.building_source, specification.latitude, specification.longitude, provider.settings, provider.building_sources, ) if building_result.source_selection != expected_source: raise ValueError( f"Scene {scene.name} has a different or unrecorded building-source " "selection/version; use its original software and source configuration " "to resume, or choose a new dataset directory" ) if building_result is None and specification.include_buildings: if scene_has_completed_cases: raise ValueError( f"Scene {scene.name} predates resumable source persistence; " "its completed cases cannot be extended safely" ) building_result = await provider.buildings( specification.latitude, specification.longitude, specification.radius_m, building_limit, source=specification.building_source, ) _write_building_source(building_source_path, building_result) elif building_result is None: building_result = BuildingQueryResult(buildings=[], provider="disabled") _write_building_source(building_source_path, building_result) environment = None environment_warnings = [] if specification.include_environment_grid: try: environment = load_environment_assets(scene_directory) except (OSError, KeyError, TypeError, ValueError, json.JSONDecodeError): if scene_has_completed_cases: raise ValueError( f"Scene {scene.name} is missing persisted terrain inputs required " "for a deterministic resume" ) from None environment, environment_warnings = await provider.environment( specification.latitude, specification.longitude, specification.radius_m, specification.terrain_resolution_m, specification.include_surface_materials, terrain_mode="elevation" if specification.include_terrain else "flat", ) assets = await _run_in_worker( worker, build_scene_assets, scene_directory, building_result.buildings, specification.longitude, specification.latitude, specification.radius_m, environment, specification.scene_mode != "empty", specification.material_profile, ) voxels = await _run_in_worker( worker, write_voxel_assets, assets, specification.receiver_height_m, specification.voxel_pitch_m, ) measurement_surfaces = [] receiver_grids = sorted( { (case.request.resolution_m, case.request.receiver_height_m) for case in scene.cases } ) for resolution_m, receiver_height_m in receiver_grids: surface = await _run_in_worker( worker, write_measurement_surface, measurement_surface_path(assets, resolution_m, receiver_height_m), assets, resolution_m, receiver_height_m, ) measurement_surfaces.append( { "path": str(surface.path.relative_to(scene_directory)), "resolution_m": resolution_m, "receiver_height_agl_m": receiver_height_m, "cell_shape": list(surface.cell_shape), "height_reference": "terrain_agl", } ) cases_directory = scene_directory / SCENE_CASES_DIRECTORY cases_directory.mkdir(exist_ok=True) scene_record = { "name": scene.name, "label": scene.label, "status": "running", "specification": specification.model_dump(mode="json", exclude={"cases"}), "provenance": { "building_provider": building_result.provider, "building_release": building_result.release, "scene_boundary_policy": SCENE_BOUNDARY_POLICY, "source_building_count": building_result.source_count, "boundary_excluded_building_count": building_result.boundary_excluded_count, "exported_building_count": assets.building_count, "triangle_count": assets.triangle_count, "terrain_provider": ( environment.terrain.source.provider if environment is not None else "disabled" ), "terrain_spacing_m": ( environment.terrain.spacing_m if environment is not None else 0.0 ), "surface_provider": ( environment.surface_provenance.get("provider", "disabled") if environment is not None else "disabled" ), }, "shared_assets": { "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(), "voxels": str(voxels.voxel_path.relative_to(scene_directory)), "receiver_slice": str(voxels.slice_path.relative_to(scene_directory)), "voxel_depth": str(voxels.depth_array_path.relative_to(scene_directory)), "measurement_surfaces": measurement_surfaces, **( { "environment": SCENE_TERRAIN_INFO_PATH.as_posix(), "environment_arrays": { name: path.as_posix() for name, path in SCENE_ENVIRONMENT_ARRAY_PATHS.items() }, } if environment is not None else {} ), }, "cases": [], } _write_json(scene_directory / SCENE_INFO_PATH, scene_record) manifest["scenes"][scene_index] = scene_record _write_json(manifest_path, manifest) for case in scene.cases: case_directory = cases_directory / case.name existing_record = completed_records.get((scene.name, case.name)) if existing_record is not None: scene_record["cases"].append(existing_record) continue case_directory.mkdir(exist_ok=True) progress(processed, plan.total_cases, f"Running {scene.name}/{case.name}") case_started = time.perf_counter() _write_json( case_directory / "case_info.json", { "name": case.name, "status": "running", "seed": case.seed, "scene": "../../scene_info.json", "request": case.request.model_dump(mode="json"), "started_at": datetime.now(timezone.utc).isoformat(), }, ) try: if case.request.engine == EngineName.sionna.value: async with sionna_lock: radio_data, engine_details, warnings = await _run_in_worker( worker, simulate_sionna_scene, case.request, assets, len(building_result.buildings), ) else: radio_data = await _run_in_worker( worker, simulate_preview, case.request, building_result.buildings, ) warnings = [ "Preview mode is analytical and is not a calibrated ray-tracing result." ] if case.request.include_terrain: warnings.append( "Terrain and surface classes are exported but are not propagation " "inputs to the planar CPU preview; use Sionna RT for terrain-aware " "cases." ) engine_details = _preview_details( case.request, len(building_result.buildings), radio_data.rows * radio_data.columns, ) warnings.extend(environment_warnings) products = build_coverage_products( case.request, radio_data, building_result.buildings, environment if case.request.engine == EngineName.sionna.value else None, ) response = SimulationResponse( simulation_id=case.name, engine=case.request.engine, elapsed_ms=round((time.perf_counter() - case_started) * 1000), coverage=products.coverage, buildings=[], warnings=warnings, engine_details=engine_details, data_provenance={ "scene": scene.name, "shared_geometry": "../../scene_info.json", "building_provider": building_result.provider, **building_result.source_provenance, "building_count": assets.building_count, }, ) image_products = await _run_in_worker( worker, write_radio_products, case_directory, response, radio_data, products.building_mask, payload.output.array_metrics, payload.output.image_metrics, payload.output.association_image, ) case_record = { "name": case.name, "status": "completed", "seed": case.seed, "scene": "../../scene_info.json", "request": case.request.model_dump(mode="json"), "grid_mapping": response.coverage.grid_mapping.model_dump(mode="json"), "radio_layers": { name: layer.model_dump(mode="json") for name, layer in response.coverage.layers.items() }, "array_contract": { "shape": list(radio_data.shape), "axis_order": [ "transmitter", "south_to_north_row", "west_to_east_column", ], **{ metric: { "path": RADIO_METRIC_ARTIFACTS[metric]["array"], "unit": RADIO_METRIC_ARTIFACTS[metric]["unit"], } for metric in payload.output.array_metrics }, }, "image_products": image_products, "engine_details": engine_details, "elapsed_ms": response.elapsed_ms, "warnings": warnings, } _write_json(case_directory / "case_info.json", case_record) except Exception as exc: case_record = { "name": case.name, "status": "failed", "seed": case.seed, "scene": "../../scene_info.json", "request": case.request.model_dump(mode="json"), "error": str(exc), } _write_json(case_directory / "case_info.json", case_record) manifest["summary"]["failed_cases"] += 1 run_record["failed_cases"] += 1 if not payload.continue_on_error: scene_record["cases"].append(case_record) scene_record["status"] = "failed" manifest["status"] = "failed" run_record["status"] = "failed" run_record["completed_at"] = datetime.now(timezone.utc).isoformat() _write_json(scene_directory / SCENE_INFO_PATH, scene_record) _write_json(manifest_path, manifest) raise else: manifest["summary"]["completed_cases"] += 1 run_record["generated_cases"] += 1 scene_record["cases"].append(case_record) processed += 1 _write_json(scene_directory / SCENE_INFO_PATH, scene_record) _write_json(manifest_path, manifest) progress(processed, plan.total_cases, f"Completed {scene.name}/{case.name}") scene_record["status"] = ( "completed_with_errors" if any(case["status"] == "failed" for case in scene_record["cases"]) else "completed" ) _write_json(scene_directory / SCENE_INFO_PATH, scene_record) manifest["status"] = ( "completed_with_errors" if manifest["summary"]["failed_cases"] else "completed" ) manifest["completed_at"] = datetime.now(timezone.utc).isoformat() run_record["status"] = manifest["status"] run_record["completed_at"] = manifest["completed_at"] _write_json(manifest_path, manifest) return dataset_directory