diff --git a/model/common/src/icon4py/model/common/backend_configuration.py b/model/common/src/icon4py/model/common/backend_configuration.py new file mode 100644 index 0000000000..f5803d0b99 --- /dev/null +++ b/model/common/src/icon4py/model/common/backend_configuration.py @@ -0,0 +1,122 @@ +# ICON4Py - ICON inspired code in Python and GT4Py +# +# Copyright (c) 2022-2024, ETH Zurich and MeteoSwiss +# All rights reserved. +# +# Please, refer to the LICENSE file in the root directory. +# SPDX-License-Identifier: BSD-3-Clause +"""External workspace allocation for the DaCe backend. + +The DaCe backend in GT4Py can be configured with an +:class:`~gt4py.next.program_processors.runners.dace.workflow.common.ExternalWorkspace` +to provide workspace memory for transient SDFG arrays +(``transient_memory_mode = EXTERNAL``). When a :class:`BackendConfig` is +provided, :func:`get_dace_options ` +calls the :class:`IconWorkspaceAllocator` defined here — a process-wide +singleton that caches a single workspace slab per device and reuses it across +every compiled program. + +The size of the workspace is configurable per experiment via :class:`BackendConfig` +(see :func:`backend_config_from_env` for an environment-variable based default). +""" + +from __future__ import annotations + +import dataclasses +import os +import typing +from collections.abc import Iterable +from typing import ClassVar, Final + +import gt4py.next as gtx +from gt4py.next.program_processors.runners.dace.workflow import common as gtx_wfdcommon + +from icon4py.model.common.config import options as common_conf_opt +from icon4py.model.common.utils import data_allocation + + +@dataclasses.dataclass(frozen=True, kw_only=True) +class BackendConfig: + """External DaCe workspace sizing, configurable per experiment.""" + + workspace_size: typing.Annotated[ + int, + common_conf_opt.ConfigOption( + description=( + "Size of the workspace memory (in Bytes) for externally allocated " + "temporary fields. This is a performance feature of the DaCe backend " + "to avoid runtime allocation of temporary fields, for each program " + "call. Note that the memory buffer is allocated once and shared " + "across all compiled programs." + ), + icon_equivalent=None, + ), + ] = 256 * 1024 * 1024 # 256 MiB + + def __post_init__(self) -> None: + if self.workspace_size <= 0: + raise ValueError(f"'workspace_size' must be positive, got {self.workspace_size}.") + + +def backend_config_from_env() -> BackendConfig | None: + """Build a :class:`BackendConfig` from environment variables. + + Reads ``ICON4PY_BACKEND_WORKSPACE_SIZE`` and returns ``None`` when it is + not set. + """ + size = os.environ.get("ICON4PY_BACKEND_WORKSPACE_SIZE") + if size is None: + return None + return BackendConfig(workspace_size=int(size)) + + +def _get_slab(nbytes: int, device: gtx.DeviceType) -> data_allocation.NDArray: + """Allocate a `nbytes`-byte buffer allocated on ``device``.""" + xp = data_allocation.array_ns(use_cupy=(device != gtx.DeviceType.CPU)) + return xp.empty(nbytes, dtype=xp.uint8) + + +class IconWorkspaceAllocator: + """Singleton workspace allocator for the DaCe backend. + + Exactly one instance exists per process (enforced by `__new__`); all DaCe + backends share it via the module-level `ICON_WORKSPACE_ALLOCATOR`. It keeps + a single private workspace slab per device in `_workspace_slabs`, reused + across every compiled program. On a cache hit the slab's size is validated + against the value passed to `allocate`. + """ + + _instance: ClassVar[IconWorkspaceAllocator | None] = None + _workspace_slabs: ClassVar[dict[gtx.DeviceType, data_allocation.NDArray]] = {} + + def __new__(cls) -> IconWorkspaceAllocator: + if cls._instance is None: + cls._instance = super().__new__(cls) + return cls._instance + + def allocate( + self, + devices: gtx.DeviceType | Iterable[gtx.DeviceType], + *, + size: int, + ) -> gtx_wfdcommon.ExternalWorkspace: + if isinstance(devices, gtx.DeviceType): + devices = [devices] + wsp: gtx_wfdcommon.ExternalWorkspace = {} + for dev in devices: + if (cached := self._workspace_slabs.get(dev)) is not None: + if cached.nbytes != size: + raise ValueError( + f"Workspace size mismatch for {dev!s}: cached slab has " + f"{cached.nbytes} bytes but 'allocate' was called with " + f"size={size}." + ) + wsp[dev] = cached + else: + slab = _get_slab(size, dev) + self._workspace_slabs[dev] = slab + wsp[dev] = slab + return wsp + + +ICON_WORKSPACE_ALLOCATOR: Final[IconWorkspaceAllocator] = IconWorkspaceAllocator() diff --git a/model/common/src/icon4py/model/common/model_backends.py b/model/common/src/icon4py/model/common/model_backends.py index 5dce3b926a..e1eede8b0c 100644 --- a/model/common/src/icon4py/model/common/model_backends.py +++ b/model/common/src/icon4py/model/common/model_backends.py @@ -5,12 +5,17 @@ # # Please, refer to the LICENSE file in the root directory. # SPDX-License-Identifier: BSD-3-Clause +from __future__ import annotations + from typing import Any, Final, TypeAlias, TypeGuard import gt4py.next as gtx +import gt4py.next.custom_layout_allocators as gtx_allocators import gt4py.next.typing as gtx_typing -from gt4py.next import backend as gtx_backend, custom_layout_allocators as gtx_allocators +from gt4py.next import backend as gtx_backend from gt4py.next.program_processors.runners import dace as gtx_dace, gtfn +from gt4py.next.program_processors.runners.dace import transformations as gtx_transformations +from gt4py.next.program_processors.runners.dace.workflow import common as gtx_wfdcommon # DeviceType should always be imported from here, as we might replace it by an ICON4Py internal implementation @@ -81,6 +86,7 @@ def make_custom_dace_backend( use_metrics: bool = True, use_zero_origin: bool = False, use_max_domain_range_on_unstructured_shift: bool | None = None, + external_workspace: gtx_wfdcommon.ExternalWorkspace | None = None, **_, ) -> gtx_typing.Backend: """Customize the dace backend with the given configuration parameters. @@ -97,15 +103,34 @@ def make_custom_dace_backend( use_max_domain_range_on_unstructured_shift: When True, compute `as_fieldop` expressions everywhere. Otherwise, when all connectivities are given at compile time, infer the minimal domain of all `as_fieldop` statically. + external_workspace: The external workspace memory to use as storage for + the transient arrays. If `None`, the transient arrays will be allocated + inside the SDFG with scope lifetime. Returns: A dace backend with custom configuration for the target device. """ + if external_workspace is not None: + if optimization_args is None: + optimization_args = { + "transient_memory_mode": gtx_transformations.TransientMemoryMode.EXTERNAL, + } + elif transient_memory_mode := optimization_args.get("transient_memory_mode"): + if transient_memory_mode != gtx_transformations.TransientMemoryMode.EXTERNAL: + raise ValueError( + f"Cannot use external workspace with transient_memory_mode={transient_memory_mode}." + ) + else: + optimization_args["transient_memory_mode"] = ( + gtx_transformations.TransientMemoryMode.EXTERNAL + ) + on_gpu = device == GPU return gtx_dace.make_dace_backend( gpu=on_gpu, auto_optimize=auto_optimize, async_sdfg_call=async_sdfg_call, + external_workspace=external_workspace, optimization_args=optimization_args, unstructured_horizontal_has_unit_stride=True, use_metrics=use_metrics, diff --git a/model/common/src/icon4py/model/common/model_options.py b/model/common/src/icon4py/model/common/model_options.py index 23dd55c9fc..11fdc99f6e 100644 --- a/model/common/src/icon4py/model/common/model_options.py +++ b/model/common/src/icon4py/model/common/model_options.py @@ -16,7 +16,7 @@ from gt4py.next import backend as gtx_backend from gt4py.next.program_processors.runners.dace import transformations as gtx_transformations -from icon4py.model.common import model_backends +from icon4py.model.common import backend_configuration as backend_cfg, model_backends log = logging.getLogger(__name__) @@ -35,11 +35,25 @@ def _dace_remove_access_node_copies(sdfg: dace.SDFG) -> None: def get_dace_options( - program_name: str, **backend_descriptor: Any + program_name: str, + backend_config: backend_cfg.BackendConfig | None, + **backend_descriptor: Any, ) -> model_backends.BackendDescriptor: - is_rocm_device = backend_descriptor.get("device") == model_backends.DeviceType.ROCM + device = backend_descriptor.get("device") or model_backends.CPU optimization_args = backend_descriptor.get("optimization_args", {}) optimization_hooks = optimization_args.get("optimization_hooks", {}) + + if backend_config is not None: + # The workspace memory allows to avoid the overhead of runtime allocations, + # which are expensive in the AMD runtime. + backend_descriptor["external_workspace"] = backend_cfg.ICON_WORKSPACE_ALLOCATOR.allocate( + device, + size=backend_config.workspace_size, + ) + optimization_args["transient_memory_mode"] = ( + gtx_transformations.TransientMemoryMode.EXTERNAL + ) + if program_name in [ "vertically_implicit_solver_at_corrector_step", "vertically_implicit_solver_at_predictor_step", @@ -60,7 +74,7 @@ def get_dace_options( backend_descriptor["use_zero_origin"] = True if program_name == "graupel_run": optimization_args["fuse_tasklets"] = True - if not is_rocm_device: + if device != model_backends.DeviceType.ROCM: optimization_args["gpu_maxnreg"] = 80 optimization_args["gpu_block_size_2d"] = (64, 6) optimization_args["gpu_memory_pool"] = False @@ -78,12 +92,17 @@ def get_gtfn_options( return backend_descriptor -def get_options(program_name: str, **backend_descriptor: Any) -> model_backends.BackendDescriptor: +def get_options( + program_name: str, + *, + backend_config: backend_cfg.BackendConfig | None, + **backend_descriptor: Any, +) -> model_backends.BackendDescriptor: if "backend_factory" not in backend_descriptor: # here we could set a backend_factory per program backend_descriptor["backend_factory"] = model_backends.make_custom_dace_backend if backend_descriptor["backend_factory"] == model_backends.make_custom_dace_backend: - backend_descriptor = get_dace_options(program_name, **backend_descriptor) + backend_descriptor = get_dace_options(program_name, backend_config, **backend_descriptor) if backend_descriptor["backend_factory"] == model_backends.make_custom_gtfn_backend: backend_descriptor = get_gtfn_options(program_name, **backend_descriptor) @@ -96,7 +115,9 @@ def customize_backend( | model_backends.DeviceType | model_backends.BackendDescriptor | None, + backend_config: backend_cfg.BackendConfig | None = None, ) -> gtx_typing.Backend | None: + backend_config = backend_config or backend_cfg.backend_config_from_env() program_name = program.__name__ if program is not None else "" if backend is None or isinstance(backend, gtx_backend.Backend): backend_name = backend.name if backend is not None else "embedded" @@ -106,7 +127,9 @@ def customize_backend( backend_descriptor = ( {"device": backend} if isinstance(backend, model_backends.DeviceType) else backend ) - backend_descriptor = get_options(program_name, **backend_descriptor) + backend_descriptor = get_options( + program_name, backend_config=backend_config, **backend_descriptor + ) backend_descriptor["device"] = backend_descriptor.get( "device", model_backends.CPU ) # set default device @@ -132,10 +155,11 @@ def setup_program( horizontal_sizes: dict[str, gtx.int32] | None = None, vertical_sizes: dict[str, gtx.int32] | None = None, offset_provider: gtx_typing.OffsetProvider | None = None, + backend_config: backend_cfg.BackendConfig | None = None, ) -> Callable[..., None]: """ This function processes arguments to the GT4Py program. It - - binds arguments that don't change during model run ('constant_args', 'horizontal_sizes', "vertical_sizes'); + - binds arguments that don't change during model run ('constant_args', 'horizontal_sizes', 'vertical_sizes'); - inlines scalar arguments into the GT4Py program at compile-time (via GT4Py's 'compile'). Args: - backend: GT4Py backend, @@ -145,6 +169,8 @@ def setup_program( - horizontal_sizes: horizontal domain bounds, - vertical_sizes: vertical domain bounds, - offset_provider: GT4Py offset_provider, + - backend_config: external DaCe workspace sizing, or `None` to fall back + to the 'ICON4PY_BACKEND_WORKSPACE_SIZE' environment variable. """ constant_args = {} if constant_args is None else constant_args variants = {} if variants is None else variants @@ -152,7 +178,7 @@ def setup_program( vertical_sizes = {} if vertical_sizes is None else vertical_sizes offset_provider = {} if offset_provider is None else offset_provider - backend = customize_backend(program, backend) + backend = customize_backend(program, backend, backend_config=backend_config) bound_static_args = {k: v for k, v in constant_args.items() if gtx.is_scalar_type(v)} static_args_program = program.with_backend(backend) diff --git a/model/common/src/icon4py/model/common/utils/data_allocation.py b/model/common/src/icon4py/model/common/utils/data_allocation.py index 23f4071eb3..f2b1014661 100644 --- a/model/common/src/icon4py/model/common/utils/data_allocation.py +++ b/model/common/src/icon4py/model/common/utils/data_allocation.py @@ -8,7 +8,6 @@ from __future__ import annotations -import logging as log import math from types import ModuleType from typing import TYPE_CHECKING, Any, TypeAlias, TypeGuard, TypeVar @@ -71,15 +70,17 @@ def as_numpy(array: NDArrayInterface) -> np.ndarray: return cp.asnumpy(array) -def _array_ns(try_cupy: bool) -> ModuleType: - """CuPy if requested and installed, NumPy otherwise.""" - if try_cupy: +def array_ns(use_cupy: bool) -> ModuleType: + """CuPy if requested, NumPy otherwise. + + Raises RuntimeError if CuPy is requested but not available. + """ + if use_cupy: try: import cupy as cp # noqa: PLC0415 [import-outside-top-level] - - return cp - except ImportError: - log.warning("No cupy installed, falling back to numpy for array_ns") + except ImportError as err: + raise RuntimeError(f"cupy is not available: {err!r}.") from err + return cp import numpy as np # noqa: PLC0415 [import-outside-top-level] return np @@ -87,7 +88,7 @@ def _array_ns(try_cupy: bool) -> ModuleType: def import_array_ns(allocator: gtx_typing.Allocator | None) -> ModuleType: """Import cupy or numpy depending on a chosen GT4Py backend DevicType.""" - return _array_ns(device_utils.is_cupy_device(allocator)) + return array_ns(device_utils.is_cupy_device(allocator)) def scalar_like_array[ScalarT: gtx_typing.Scalar]( diff --git a/model/common/tests/common/test_backend_configuration.py b/model/common/tests/common/test_backend_configuration.py new file mode 100644 index 0000000000..eb8476aa24 --- /dev/null +++ b/model/common/tests/common/test_backend_configuration.py @@ -0,0 +1,95 @@ +# ICON4Py - ICON inspired code in Python and GT4Py +# +# Copyright (c) 2022-2024, ETH Zurich and MeteoSwiss +# All rights reserved. +# +# Please, refer to the LICENSE file in the root directory. +# SPDX-License-Identifier: BSD-3-Clause +from __future__ import annotations + +import dataclasses +import sys +from collections.abc import Iterator + +import numpy as np +import pytest + +from icon4py.model.common import backend_configuration as bc + + +CPU = bc.gtx.DeviceType.CPU + + +@pytest.fixture +def allocator() -> Iterator[bc.IconWorkspaceAllocator]: + bc.IconWorkspaceAllocator._workspace_slabs.clear() + yield bc.ICON_WORKSPACE_ALLOCATOR + bc.IconWorkspaceAllocator._workspace_slabs.clear() + + +class TestBackendConfig: + def test_valid_construction(self) -> None: + config = bc.BackendConfig(workspace_size=1024) + assert config.workspace_size == 1024 + + @pytest.mark.parametrize("size", [0, -1, -1024]) + def test_invalid_size_raises(self, size: int) -> None: + with pytest.raises(ValueError, match="workspace_size"): + bc.BackendConfig(workspace_size=size) + + def test_is_frozen(self) -> None: + config = bc.BackendConfig(workspace_size=1024) + with pytest.raises(dataclasses.FrozenInstanceError): + config.workspace_size = 2048 # type: ignore[misc] + + +class TestBackendConfigFromEnv: + def test_returns_none_when_size_not_set(self, monkeypatch: pytest.MonkeyPatch) -> None: + monkeypatch.delenv("ICON4PY_BACKEND_WORKSPACE_SIZE", raising=False) + assert bc.backend_config_from_env() is None + + def test_reads_size_from_env(self, monkeypatch: pytest.MonkeyPatch) -> None: + monkeypatch.setenv("ICON4PY_BACKEND_WORKSPACE_SIZE", "8192") + config = bc.backend_config_from_env() + assert config is not None + assert config.workspace_size == 8192 + + +class TestGetSlab: + @pytest.mark.parametrize("nbytes", [1, 100, 1000, 4096]) + def test_slab_has_correct_size(self, nbytes: int) -> None: + slab = bc._get_slab(nbytes, CPU) + assert slab.nbytes == nbytes + + def test_raises_runtime_error_when_cupy_not_available( + self, monkeypatch: pytest.MonkeyPatch + ) -> None: + monkeypatch.setitem(sys.modules, "cupy", None) + with pytest.raises(RuntimeError, match="cupy is not available:"): + bc._get_slab(1000, bc.gtx.DeviceType.CUDA) + + +class TestIconWorkspaceAllocator: + def test_is_singleton(self) -> None: + assert bc.IconWorkspaceAllocator() is bc.IconWorkspaceAllocator() + assert bc.ICON_WORKSPACE_ALLOCATOR is bc.IconWorkspaceAllocator() + + def test_allocate_single_device(self, allocator: bc.IconWorkspaceAllocator) -> None: + wsp = allocator.allocate(CPU, size=512) + assert CPU in wsp + assert np.asarray(wsp[CPU]).nbytes == 512 + + def test_allocate_iterable_of_devices(self, allocator: bc.IconWorkspaceAllocator) -> None: + wsp = allocator.allocate([CPU], size=512) + assert CPU in wsp + assert np.asarray(wsp[CPU]).nbytes == 512 + + def test_cache_hit_returns_same_slab(self, allocator: bc.IconWorkspaceAllocator) -> None: + wsp1 = allocator.allocate(CPU, size=512) + wsp2 = allocator.allocate(CPU, size=512) + assert wsp1[CPU] is wsp2[CPU] + + def test_cache_hit_size_mismatch_raises(self, allocator: bc.IconWorkspaceAllocator) -> None: + allocator.allocate(CPU, size=512) + with pytest.raises(ValueError, match="size mismatch"): + allocator.allocate(CPU, size=1024) diff --git a/model/common/tests/common/test_model_options.py b/model/common/tests/common/test_model_options.py index ab5ec10232..2d1e3fb453 100644 --- a/model/common/tests/common/test_model_options.py +++ b/model/common/tests/common/test_model_options.py @@ -12,10 +12,19 @@ import gt4py.next.typing as gtx_typing import pytest -from icon4py.model.common import field_type_aliases as fa, model_backends +from icon4py.model.common import ( + backend_configuration as backend_cfg, + field_type_aliases as fa, + model_backends, +) from icon4py.model.common.model_options import customize_backend, setup_program +@pytest.fixture(autouse=True) +def clear_backend_workspace_env(monkeypatch: pytest.MonkeyPatch) -> None: + monkeypatch.delenv("ICON4PY_BACKEND_WORKSPACE_SIZE", raising=False) + + @gtx.field_operator # type: ignore[call-overload] def field_op_return_field(field: fa.CellKField[float], factor: float) -> fa.CellKField[float]: return field + factor @@ -57,6 +66,17 @@ def test_custom_backend_device() -> None: assert repr(default_backend) == repr(backend) +def test_custom_backend_with_external_workspace_config_and_no_explicit_device() -> None: + backend = customize_backend( + None, + {"backend_factory": model_backends.make_custom_dace_backend, "device": None}, + backend_config=backend_cfg.BackendConfig(workspace_size=4096), + ) + assert backend is not None + assert hasattr(backend, "external_workspace") + assert backend.external_workspace[gtx.DeviceType.CPU].size == 4096 + + @pytest.mark.parametrize( "backend", [ diff --git a/model/driver/src/icon4py/model/driver/config.py b/model/driver/src/icon4py/model/driver/config.py index 5c1c5792dd..348c7a5598 100644 --- a/model/driver/src/icon4py/model/driver/config.py +++ b/model/driver/src/icon4py/model/driver/config.py @@ -26,6 +26,7 @@ from icon4py.model.atmosphere.subgrid_scale_physics.muphys import config as muphys_config from icon4py.model.atmosphere.tracer_advection import tracer_advection from icon4py.model.common import ( + backend_configuration as backend_cfg, initial_condition, prescribed_tendencies, time, @@ -247,6 +248,17 @@ class DriverConfig: icon_equivalent=None, ), ] = False + backend_config: typing.Annotated[ + backend_cfg.BackendConfig | None, + common_conf_opt.ConfigOption( + description=( + "Backend configuration options, which affect performance but not " + "the scientific outcome. `None` falls back to environment variables, " + "if set, otherwise the default configuration is used." + ), + icon_equivalent=None, + ), + ] = dataclasses.field(default_factory=backend_cfg.backend_config_from_env) output_backend: typing.Annotated[ common_io.OutputBackend, common_conf_opt.ConfigOption( diff --git a/model/driver/src/icon4py/model/driver/main.py b/model/driver/src/icon4py/model/driver/main.py index c21c174fdf..0fdfd81e39 100644 --- a/model/driver/src/icon4py/model/driver/main.py +++ b/model/driver/src/icon4py/model/driver/main.py @@ -86,11 +86,6 @@ def main( run. """ - backend = model_options.customize_backend( - program=None, backend=driver_utils.get_backend_from_name(icon4py_backend) - ) - allocator = model_backends.get_allocator(backend) - process_props = decomposition_defs.get_process_properties( decomposition_defs.get_runtype(with_mpi=mpi_decomp.mpi4py is not None) ) @@ -110,6 +105,13 @@ def main( driver_overrides["output_path"] = output_path config = config.with_overrides(driver=driver_overrides) + backend = model_options.customize_backend( + program=None, + backend=driver_utils.get_backend_from_name(icon4py_backend), + backend_config=config.driver.backend_config, + ) + allocator = model_backends.get_allocator(backend) + grid_manager = driver_utils.create_grid_manager( grid_file_path=grid_file_path, vertical_grid_config=config.vertical_grid, diff --git a/model/driver/tests/driver/unit_tests/data/test_config.yml b/model/driver/tests/driver/unit_tests/data/test_config.yml index 11afd1ecdc..067b335646 100644 --- a/model/driver/tests/driver/unit_tests/data/test_config.yml +++ b/model/driver/tests/driver/unit_tests/data/test_config.yml @@ -82,6 +82,7 @@ driver: ndyn_substeps: 5 enable_statistics_logging: false enable_output: false + backend_config: output_backend: zarr output_mode: distributed nonhydrostatic: diff --git a/pyproject.toml b/pyproject.toml index cded4a069c..2d5c7af64b 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -423,15 +423,15 @@ name = 'gridtools' url = 'https://gridtools.github.io/pypi/' [tool.uv.sources] +# dace = {index = "gridtools"} +# gt4py = {git = "https://github.com/GridTools/gt4py", branch = "main"} +# gt4py = {index = "test.pypi"} icon4py-atmosphere-diffusion = {workspace = true} icon4py-atmosphere-dycore = {workspace = true} icon4py-atmosphere-microphysics = {workspace = true} icon4py-atmosphere-muphys = {workspace = true} icon4py-atmosphere-physics-driver = {workspace = true} icon4py-atmosphere-tmx = {workspace = true} -# gt4py = {git = "https://github.com/GridTools/gt4py", branch = "main"} -# gt4py = {index = "test.pypi"} -# dace = {index = "gridtools"} icon4py-atmosphere-tracer_advection = {workspace = true} icon4py-bindings = {workspace = true} icon4py-common = {workspace = true}