diff --git a/sdk/python/feast/feature_server.py b/sdk/python/feast/feature_server.py index 7412931ce8a..e42e28d6db9 100644 --- a/sdk/python/feast/feature_server.py +++ b/sdk/python/feast/feature_server.py @@ -20,6 +20,7 @@ import time import traceback from collections import defaultdict +from concurrent.futures import ThreadPoolExecutor from contextlib import asynccontextmanager from datetime import datetime from importlib import resources as importlib_resources @@ -31,6 +32,7 @@ from fastapi import ( Depends, FastAPI, + Query, Request, Response, WebSocket, @@ -41,7 +43,7 @@ from fastapi.logger import logger from fastapi.responses import JSONResponse from fastapi.staticfiles import StaticFiles -from pydantic import BaseModel, field_validator +from pydantic import BaseModel, Field, field_validator import feast from feast import metrics as feast_metrics @@ -54,6 +56,7 @@ ) from feast.feast_object import FeastObject from feast.feature_server_utils import convert_response_to_dict +from feast.feature_view import FeatureViewState from feast.feature_view_utils import get_feature_view_from_feature_store from feast.filter_models import ComparisonFilter, CompoundFilter from feast.permissions.action import WRITE, AuthzedAction @@ -94,12 +97,26 @@ class MaterializeRequest(BaseModel): feature_views: Optional[List[str]] = None disable_event_timestamp: bool = False full_feature_names: bool = False + version: Optional[str] = Field( + None, + description=( + "Optional version to materialize (e.g. 'v2'). Requires feature_views " + "with exactly one entry and registry.enable_online_feature_view_versioning." + ), + ) class MaterializeIncrementalRequest(BaseModel): end_ts: str feature_views: Optional[List[str]] = None full_feature_names: bool = False + version: Optional[str] = Field( + None, + description=( + "Optional version to materialize (e.g. 'v2'). Requires feature_views " + "with exactly one entry and registry.enable_online_feature_view_versioning." + ), + ) class GetOnlineFeaturesRequest(BaseModel): @@ -333,6 +350,128 @@ async def load_static_artifacts(app: FastAPI, store): logger.warning(f"Failed to load static artifacts: {e}") +def _authorize_materialize_views( + store: "feast.FeatureStore", + feature_view_names: Optional[List[str]], + version: Optional[str] = None, +) -> List[str]: + """Resolve + authorize feature views for materialization. + + Returns the resolved list of FV names (all eligible FVs when + feature_view_names is None). + """ + parsed_version = store._validate_materialize_version(version, feature_view_names) + feature_views_to_materialize = store._get_feature_views_to_materialize( + feature_view_names, version=parsed_version + ) + for fv in feature_views_to_materialize: + assert_permissions( + resource=fv, + actions=[AuthzedAction.WRITE_ONLINE], + ) + return [fv.name for fv in feature_views_to_materialize] + + +def _check_already_materializing( + store: "feast.FeatureStore", + fv_names: List[str], +) -> Optional[JSONResponse]: + """Return a 409 JSONResponse if any requested FV is already MATERIALIZING.""" + conflicting: List[str] = [] + for fv_name in fv_names: + try: + fv = store.registry.get_feature_view( + fv_name, store.project, allow_cache=False + ) + if getattr(fv, "state", None) == FeatureViewState.MATERIALIZING: + conflicting.append(fv_name) + except (FeatureViewNotFoundException, KeyError): + pass + except Exception as e: + logger.warning( + f"Unexpected error checking MATERIALIZING state for {fv_name}: {e}" + ) + if conflicting: + return JSONResponse( + status_code=409, + content={ + "error": ( + f"Cannot start async materialization — the following feature " + f"views are already in MATERIALIZING state: {conflicting}. " + f"Use ?force=true to override." + ), + "feature_views": conflicting, + }, + ) + return None + + +def _update_fv_state( + store: "feast.FeatureStore", + fv_names: List[str], + state: FeatureViewState, +) -> None: + """Set FV state in the registry for each named feature view.""" + for fv_name in fv_names: + try: + fv = store.registry.get_feature_view( + fv_name, store.project, allow_cache=False + ) + fv.state = state + store.registry.apply_feature_view(fv, store.project) + except (FeatureViewNotFoundException, KeyError): + logger.warning(f"Feature view {fv_name} not found; skip state={state}") + except Exception as e: + logger.warning(f"Failed to set state={state} for {fv_name}: {e}") + + +def _reset_stuck_materializing_to_generated( + store: "feast.FeatureStore", + fv_names: List[str], +) -> None: + """Reset FVs currently in MATERIALIZING to GENERATED (force override).""" + stuck: List[str] = [] + for fv_name in fv_names: + try: + fv = store.registry.get_feature_view( + fv_name, store.project, allow_cache=False + ) + if getattr(fv, "state", None) == FeatureViewState.MATERIALIZING: + stuck.append(fv_name) + except (FeatureViewNotFoundException, KeyError): + pass + except Exception as e: + logger.warning(f"Unexpected error while force-resetting {fv_name}: {e}") + if stuck: + _update_fv_state(store, stuck, FeatureViewState.GENERATED) + logger.info( + "Force reset MATERIALIZING → GENERATED for feature views: %s", stuck + ) + + +def _parse_materialize_timestamps( + request: "MaterializeRequest", +) -> tuple: + """Parse and validate start/end timestamps from a MaterializeRequest.""" + if request.disable_event_timestamp: + now = datetime.now() + return datetime(1970, 1, 1), now + + if not request.start_ts or not request.end_ts: + raise ValueError( + "start_ts and end_ts are required when disable_event_timestamp is False" + ) + try: + start_date = utils.make_tzaware(parser.parse(request.start_ts)) + end_date = utils.make_tzaware(parser.parse(request.end_ts)) + except (ValueError, TypeError) as e: + raise ValueError(f"Invalid timestamp format: {e}") from e + + if start_date >= end_date: + raise ValueError(f"start_ts ({start_date}) must be before end_ts ({end_date})") + return start_date, end_date + + def get_app( store: "feast.FeatureStore", registry_ttl_sec: int = DEFAULT_FEATURE_SERVER_REGISTRY_TTL, @@ -412,6 +551,22 @@ def get_app( else: logger.debug("Offline write batching is DISABLED") + # Dedicated pool for async materialize so long Spark/offline waits do not + # starve the default executor used by online serving and run_in_threadpool. + _mat_workers_raw = os.environ.get("FEAST_MATERIALIZE_MAX_WORKERS", "2") + try: + materialize_max_workers = max(1, int(_mat_workers_raw)) + except ValueError: + logger.warning( + "Invalid FEAST_MATERIALIZE_MAX_WORKERS=%r; using default 2", + _mat_workers_raw, + ) + materialize_max_workers = 2 + materialize_executor = ThreadPoolExecutor( + max_workers=materialize_max_workers, + thread_name_prefix="feast-materialize", + ) + def stop_refresh(): nonlocal shutting_down shutting_down = True @@ -445,6 +600,9 @@ async def lifespan(app: FastAPI): stop_refresh() if offline_batcher is not None: offline_batcher.shutdown() + # wait=False: do not block process exit on in-flight materialize + # (same fire-and-forget contract as returning 202 mid-job). + materialize_executor.shutdown(wait=False) await store.close() app = FastAPI(lifespan=lifespan) @@ -798,70 +956,115 @@ async def chat_ui(): return Response(content=content, media_type="text/html") @app.post("/materialize", dependencies=[Depends(inject_user_details)]) - async def materialize(request: MaterializeRequest) -> None: + async def materialize( + request: MaterializeRequest, + async_mode: bool = Query(False, alias="async"), + force: bool = Query(False), + ): with feast_metrics.track_request_latency("/materialize"): - if request.feature_views: - for feature_view in request.feature_views: - resource = await _get_feast_object(feature_view, True) - assert_permissions( - resource=resource, - actions=[AuthzedAction.WRITE_ONLINE], - ) - else: - feature_views_to_materialize = store._get_feature_views_to_materialize( - None - ) - for fv in feature_views_to_materialize: - assert_permissions( - resource=fv, - actions=[AuthzedAction.WRITE_ONLINE], - ) + fv_names = _authorize_materialize_views( + store, request.feature_views, version=request.version + ) + start_date, end_date = _parse_materialize_timestamps(request) - if request.disable_event_timestamp: - now = datetime.now() - start_date = datetime(1970, 1, 1) - end_date = now - else: - if not request.start_ts or not request.end_ts: - raise ValueError( - "start_ts and end_ts are required when disable_event_timestamp is False" - ) - start_date = utils.make_tzaware(parser.parse(request.start_ts)) - end_date = utils.make_tzaware(parser.parse(request.end_ts)) + if async_mode: + if force: + _reset_stuck_materializing_to_generated(store, fv_names) + else: + conflict = _check_already_materializing(store, fv_names) + if conflict: + return conflict + + # Reserve MATERIALIZING before 202 so concurrent requests hit 409. + # store.materialize() treats already-MATERIALIZING as a no-op. + _update_fv_state(store, fv_names, FeatureViewState.MATERIALIZING) + + def _run_materialize(): + try: + store.materialize( + start_date, + end_date, + fv_names, + disable_event_timestamp=request.disable_event_timestamp, + full_feature_names=request.full_feature_names, + version=request.version, + ) + except Exception as e: + logger.error( + f"Async materialization failed for {fv_names}: {e}", + exc_info=True, + ) + _reset_stuck_materializing_to_generated(store, fv_names) + + loop = asyncio.get_running_loop() + loop.run_in_executor(materialize_executor, _run_materialize) + + return JSONResponse( + status_code=202, + content={"status": "accepted", "feature_views": fv_names}, + ) await run_in_threadpool( store.materialize, start_date, end_date, - request.feature_views, + fv_names, disable_event_timestamp=request.disable_event_timestamp, full_feature_names=request.full_feature_names, + version=request.version, ) @app.post("/materialize-incremental", dependencies=[Depends(inject_user_details)]) - async def materialize_incremental(request: MaterializeIncrementalRequest) -> None: + async def materialize_incremental( + request: MaterializeIncrementalRequest, + async_mode: bool = Query(False, alias="async"), + force: bool = Query(False), + ): with feast_metrics.track_request_latency("/materialize-incremental"): - if request.feature_views: - for feature_view in request.feature_views: - resource = await _get_feast_object(feature_view, True) - assert_permissions( - resource=resource, - actions=[AuthzedAction.WRITE_ONLINE], - ) - else: - feature_views_to_materialize = store._get_feature_views_to_materialize( - None + fv_names = _authorize_materialize_views( + store, request.feature_views, version=request.version + ) + end_date = utils.make_tzaware(parser.parse(request.end_ts)) + + if async_mode: + if force: + _reset_stuck_materializing_to_generated(store, fv_names) + else: + conflict = _check_already_materializing(store, fv_names) + if conflict: + return conflict + + _update_fv_state(store, fv_names, FeatureViewState.MATERIALIZING) + + def _run_materialize_incremental(): + try: + store.materialize_incremental( + end_date, + fv_names, + full_feature_names=request.full_feature_names, + version=request.version, + ) + except Exception as e: + logger.error( + f"Async materialize-incremental failed for {fv_names}: {e}", + exc_info=True, + ) + _reset_stuck_materializing_to_generated(store, fv_names) + + loop = asyncio.get_running_loop() + loop.run_in_executor(materialize_executor, _run_materialize_incremental) + + return JSONResponse( + status_code=202, + content={"status": "accepted", "feature_views": fv_names}, ) - for fv in feature_views_to_materialize: - assert_permissions( - resource=fv, - actions=[AuthzedAction.WRITE_ONLINE], - ) + await run_in_threadpool( store.materialize_incremental, - utils.make_tzaware(parser.parse(request.end_ts)), - request.feature_views, + end_date, + fv_names, full_feature_names=request.full_feature_names, + version=request.version, ) @app.exception_handler(Exception) diff --git a/sdk/python/feast/feature_store.py b/sdk/python/feast/feature_store.py index 5e6f520032a..546f6f7680d 100644 --- a/sdk/python/feast/feature_store.py +++ b/sdk/python/feast/feature_store.py @@ -509,8 +509,15 @@ def _transition_fv_to_materializing( Transition a feature view to MATERIALIZING state. Rolls back all already-transitioned FVs if this one can't transition. + Already MATERIALIZING is a no-op (async server may have reserved the state + before returning 202); rollback target is GENERATED in that case. """ - previous_state = getattr(feature_view, "state", None) + current = getattr(feature_view, "state", None) + if current == FeatureViewState.MATERIALIZING: + previous_states[feature_view.name] = FeatureViewState.GENERATED + return + + previous_states[feature_view.name] = current if ( hasattr(feature_view, "state") and feature_view.state != FeatureViewState.STATE_UNSPECIFIED @@ -523,7 +530,6 @@ def _transition_fv_to_materializing( ) feature_view.state = FeatureViewState.MATERIALIZING self.registry.apply_feature_view(feature_view, self.project, commit=True) - previous_states[feature_view.name] = previous_state def _submit_and_process_materialization_jobs( self, @@ -2435,12 +2441,116 @@ def _materialize_odfv( ) self.write_to_online_store(feature_view.name, df=transformed_df) + def _get_remote_materialize_url(self) -> str: + """Get the feature server URL from online_store.path for remote materialization.""" + online_cfg = self.config.online_store + url = getattr(online_cfg, "path", None) + if not url: + raise ValueError( + "online_store.path must be set to use remote materialization. " + "Configure online_store with type: remote and a valid path." + ) + return url.rstrip("/") + + def _get_remote_http_session(self): + """Get an HTTP session with auth configured for the feature server.""" + import requests + + auth_config = getattr(self.config, "auth_config", None) + if auth_config and getattr(auth_config, "type", "no_auth") != "no_auth": + from feast.permissions.client.http_auth_requests_wrapper import ( + get_http_auth_requests_session, + ) + + return get_http_auth_requests_session(auth_config) + + return requests.Session() + + def _is_remote_topology(self) -> bool: + """True when this client talks to a remote feature server for online ops.""" + return getattr(self.config.online_store, "type", None) == "remote" + + def _post_to_feature_server( + self, + endpoint: str, + payload: Dict[str, Any], + query_params: Optional[Dict[str, str]] = None, + ) -> Dict[str, Any]: + """POST JSON to the feature server; raise on 409 / 4xx / 5xx.""" + url = f"{self._get_remote_materialize_url()}{endpoint}" + session = self._get_remote_http_session() + cert = getattr(self.config.online_store, "cert", "") or "" + verify: Any = cert if cert else True + try: + response = session.post( + url, + json=payload, + params=query_params or {}, + verify=verify, + ) + if response.status_code == 409: + try: + detail = response.json() + message = detail.get("error", response.text) + except Exception: + message = response.text + raise RuntimeError( + f"Remote materialization conflict (409): {message}" + ) from None + if response.status_code >= 400: + raise RuntimeError( + f"Remote materialization failed " + f"({response.status_code}): {response.text}" + ) + if not response.content: + return {} + try: + return response.json() + except Exception: + return {"status": "accepted", "raw": response.text} + finally: + session.close() + + def _delegate_remote_materialize( + self, + endpoint: str, + payload: Dict[str, Any], + force: bool = False, + run_async: bool = False, + ) -> None: + """POST materialize to the feature server. + + When run_async=False (default), omits async and blocks until the server + finishes synchronous materialization (HTTP response). + When run_async=True, sends ?async=true and returns after 202. + force=True is only valid with run_async=True (server force applies to async). + """ + if force and not run_async: + raise ValueError( + "force=True requires run_async=True. " + "force only overrides stuck MATERIALIZING on the async path." + ) + + query_params: Dict[str, str] = {} + if run_async: + query_params["async"] = "true" + if force: + query_params["force"] = "true" + + result = self._post_to_feature_server(endpoint, payload, query_params or None) + if run_async: + _logger.info("Remote materialization accepted (%s): %s", endpoint, result) + else: + _logger.info("Remote materialization completed (%s): %s", endpoint, result) + def materialize_incremental( self, end_date: datetime, feature_views: Optional[List[str]] = None, full_feature_names: bool = False, version: Optional[str] = None, + force: bool = False, + run_async: bool = False, ) -> None: """ Materialize incremental new data from the offline store into the online store. @@ -2459,6 +2569,11 @@ def materialize_incremental( feature view name. version (str): Optional version to materialize (e.g., 'v2'). Requires feature_views with exactly one entry and enable_online_feature_view_versioning to be enabled. + force (bool): When using remote topology with run_async=True, pass force=true to + override stuck MATERIALIZING state on the feature server. Ignored for local topology. + run_async (bool): When using remote topology, if False (default) POST without async and + block until the server finishes sync materialization. If True, POST with ?async=true + and return after 202. Ignored for local topology. Raises: Exception: A feature view being materialized does not have a TTL set. @@ -2474,6 +2589,22 @@ def materialize_incremental( ... """ + if self._is_remote_topology(): + payload: Dict[str, Any] = { + "end_ts": end_date.isoformat(), + "feature_views": feature_views, + "full_feature_names": full_feature_names, + } + if version is not None: + payload["version"] = version + self._delegate_remote_materialize( + "/materialize-incremental", + payload, + force=force, + run_async=run_async, + ) + return + parsed_version = self._validate_materialize_version(version, feature_views) feature_views_to_materialize = self._get_feature_views_to_materialize( feature_views, version=parsed_version @@ -2568,10 +2699,10 @@ def tqdm_builder(length): ) else: for feature_view, start_date in regular_fvs_with_dates: - # Transition state to MATERIALIZING before starting. - # Only enforce when the state machine is active (not STATE_UNSPECIFIED). previous_state = getattr(feature_view, "state", None) - if ( + if previous_state == FeatureViewState.MATERIALIZING: + previous_state = FeatureViewState.GENERATED + elif ( hasattr(feature_view, "state") and feature_view.state != FeatureViewState.STATE_UNSPECIFIED ): @@ -2601,7 +2732,6 @@ def tqdm_builder(length): ) except Exception: fv_success = False - # Roll back state to previous value on failure. if ( hasattr(feature_view, "state") and previous_state is not None @@ -2663,6 +2793,8 @@ def materialize( disable_event_timestamp: bool = False, full_feature_names: bool = False, version: Optional[str] = None, + force: bool = False, + run_async: bool = False, ) -> None: """ Materialize data from the offline store into the online store. @@ -2681,6 +2813,11 @@ def materialize( feature view name. version (str): Optional version to materialize (e.g., 'v2'). Requires feature_views with exactly one entry and enable_online_feature_view_versioning to be enabled. + force (bool): When using remote topology with run_async=True, pass force=true to + override stuck MATERIALIZING state on the feature server. Ignored for local topology. + run_async (bool): When using remote topology, if False (default) POST without async and + block until the server finishes sync materialization. If True, POST with ?async=true + and return after 202. Ignored for local topology. Examples: Materialize all features into the online store over the interval @@ -2695,6 +2832,21 @@ def materialize( ... """ + if self._is_remote_topology(): + payload: Dict[str, Any] = { + "start_ts": start_date.isoformat(), + "end_ts": end_date.isoformat(), + "feature_views": feature_views, + "disable_event_timestamp": disable_event_timestamp, + "full_feature_names": full_feature_names, + } + if version is not None: + payload["version"] = version + self._delegate_remote_materialize( + "/materialize", payload, force=force, run_async=run_async + ) + return + if utils.make_tzaware(start_date) > utils.make_tzaware(end_date): raise ValueError( f"The given start_date {start_date} is greater than the given end_date {end_date}." @@ -2760,10 +2912,10 @@ def tqdm_builder(length): ) else: for feature_view, fv_start in regular_fvs_with_dates: - # Transition state to MATERIALIZING before starting. - # Only enforce when the state machine is active (not STATE_UNSPECIFIED). previous_state = getattr(feature_view, "state", None) - if ( + if previous_state == FeatureViewState.MATERIALIZING: + previous_state = FeatureViewState.GENERATED + elif ( hasattr(feature_view, "state") and feature_view.state != FeatureViewState.STATE_UNSPECIFIED ): @@ -2794,7 +2946,6 @@ def tqdm_builder(length): ) except Exception: fv_success = False - # Roll back state to previous value on failure. if ( hasattr(feature_view, "state") and previous_state is not None diff --git a/sdk/python/feast/infra/registry/snowflake.py b/sdk/python/feast/infra/registry/snowflake.py index 0438e997ab1..fc75db20e9c 100644 --- a/sdk/python/feast/infra/registry/snowflake.py +++ b/sdk/python/feast/infra/registry/snowflake.py @@ -1147,6 +1147,10 @@ def apply_materialization( FeatureViewNotFoundException, ) fv.materialization_intervals.append((start_date, end_date)) + if hasattr(fv, "state"): + from feast.feature_view import FeatureViewState + + fv.state = FeatureViewState.AVAILABLE_ONLINE self._apply_object( fv_table_str, project, diff --git a/sdk/python/feast/infra/registry/sql.py b/sdk/python/feast/infra/registry/sql.py index 8d2c88fba57..75d01b585cf 100644 --- a/sdk/python/feast/infra/registry/sql.py +++ b/sdk/python/feast/infra/registry/sql.py @@ -1237,6 +1237,10 @@ def apply_materialization( FeatureViewNotFoundException, ) fv.materialization_intervals.append((start_date, end_date)) + if hasattr(fv, "state"): + from feast.feature_view import FeatureViewState + + fv.state = FeatureViewState.AVAILABLE_ONLINE self._apply_object( table, project, "feature_view_name", fv, "feature_view_proto" ) diff --git a/sdk/python/tests/unit/test_materialize_version_forwarding.py b/sdk/python/tests/unit/test_materialize_version_forwarding.py new file mode 100644 index 00000000000..cd7f6b40d61 --- /dev/null +++ b/sdk/python/tests/unit/test_materialize_version_forwarding.py @@ -0,0 +1,199 @@ +# Copyright 2026 The Feast Authors +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# https://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. +"""Server tests: version is forwarded on materialize sync and async paths.""" + +from concurrent.futures import Future, ThreadPoolExecutor +from inspect import signature +from unittest.mock import AsyncMock, MagicMock, patch + +from fastapi.testclient import TestClient + +from feast.feature_server import _authorize_materialize_views, get_app +from feast.feature_store import FeatureStore +from feast.feature_view import FeatureViewState + + +def _mock_store_for_materialize(): + fs = MagicMock() + fs.project = "test_project" + fs.initialize = AsyncMock() + fs.close = AsyncMock() + fs._validate_materialize_version.return_value = "v2" + fv = MagicMock() + fv.name = "driver_hourly_stats" + fv.state = FeatureViewState.GENERATED + fs._get_feature_views_to_materialize.return_value = [fv] + fs.registry.get_feature_view.return_value = fv + fs.materialize = MagicMock() + fs.materialize_incremental = MagicMock() + return fs + + +def _run_executor_inline(loop, executor, func, *args): + """Make asyncio run_in_executor execute synchronously for tests.""" + fut: Future = Future() + try: + fut.set_result(func(*args) if args else func()) + except Exception as exc: + fut.set_exception(exc) + return fut + + +def test_run_async_defaults_to_false(): + for method_name in ( + "materialize", + "materialize_incremental", + "_delegate_remote_materialize", + ): + param = signature(getattr(FeatureStore, method_name)).parameters["run_async"] + assert param.default is False, ( + f"{method_name}.run_async default should be False" + ) + + +@patch("feast.feature_server.assert_permissions") +def test_authorize_materialize_views_forwards_version(_mock_perms): + fs = _mock_store_for_materialize() + names = _authorize_materialize_views(fs, ["driver_hourly_stats"], version="v2") + assert names == ["driver_hourly_stats"] + fs._validate_materialize_version.assert_called_once_with( + "v2", ["driver_hourly_stats"] + ) + fs._get_feature_views_to_materialize.assert_called_once_with( + ["driver_hourly_stats"], version="v2" + ) + + +@patch("feast.feature_server.assert_permissions") +def test_materialize_sync_forwards_version(_mock_perms): + fs = _mock_store_for_materialize() + client = TestClient(get_app(fs)) + + response = client.post( + "/materialize", + json={ + "start_ts": "2021-01-01T00:00:00", + "end_ts": "2021-01-02T00:00:00", + "feature_views": ["driver_hourly_stats"], + "version": "v2", + }, + ) + + assert response.status_code == 200 + fs._validate_materialize_version.assert_called_with("v2", ["driver_hourly_stats"]) + fs.materialize.assert_called_once() + _, kwargs = fs.materialize.call_args + assert kwargs.get("version") == "v2" + assert fs.materialize.call_args.args[2] == ["driver_hourly_stats"] + + +@patch("feast.feature_server.assert_permissions") +def test_materialize_incremental_sync_forwards_version(_mock_perms): + fs = _mock_store_for_materialize() + client = TestClient(get_app(fs)) + + response = client.post( + "/materialize-incremental", + json={ + "end_ts": "2021-01-02T00:00:00", + "feature_views": ["driver_hourly_stats"], + "version": "v2", + }, + ) + + assert response.status_code == 200 + fs._validate_materialize_version.assert_called_with("v2", ["driver_hourly_stats"]) + fs.materialize_incremental.assert_called_once() + _, kwargs = fs.materialize_incremental.call_args + assert kwargs.get("version") == "v2" + assert fs.materialize_incremental.call_args.args[1] == ["driver_hourly_stats"] + + +@patch("feast.feature_server.assert_permissions") +def test_materialize_async_uses_dedicated_executor(_mock_perms): + """Async materialize must not use the shared default pool (executor=None).""" + fs = _mock_store_for_materialize() + seen: dict = {} + + def capture_executor(loop, executor, func, *args): + seen["executor"] = executor + return _run_executor_inline(loop, executor, func, *args) + + with patch( + "asyncio.BaseEventLoop.run_in_executor", + new=capture_executor, + ): + client = TestClient(get_app(fs)) + response = client.post( + "/materialize?async=true", + json={ + "start_ts": "2021-01-01T00:00:00", + "end_ts": "2021-01-02T00:00:00", + "feature_views": ["driver_hourly_stats"], + }, + ) + + assert response.status_code == 202 + assert seen.get("executor") is not None + assert isinstance(seen["executor"], ThreadPoolExecutor) + assert "feast-materialize" in seen["executor"]._thread_name_prefix + + +@patch( + "asyncio.BaseEventLoop.run_in_executor", + new=_run_executor_inline, +) +@patch("feast.feature_server.assert_permissions") +def test_materialize_async_forwards_version(_mock_perms): + fs = _mock_store_for_materialize() + client = TestClient(get_app(fs)) + + response = client.post( + "/materialize?async=true", + json={ + "start_ts": "2021-01-01T00:00:00", + "end_ts": "2021-01-02T00:00:00", + "feature_views": ["driver_hourly_stats"], + "version": "v2", + }, + ) + + assert response.status_code == 202 + fs.materialize.assert_called_once() + _, kwargs = fs.materialize.call_args + assert kwargs.get("version") == "v2" + + +@patch( + "asyncio.BaseEventLoop.run_in_executor", + new=_run_executor_inline, +) +@patch("feast.feature_server.assert_permissions") +def test_materialize_incremental_async_forwards_version(_mock_perms): + fs = _mock_store_for_materialize() + client = TestClient(get_app(fs)) + + response = client.post( + "/materialize-incremental?async=true", + json={ + "end_ts": "2021-01-02T00:00:00", + "feature_views": ["driver_hourly_stats"], + "version": "v2", + }, + ) + + assert response.status_code == 202 + fs.materialize_incremental.assert_called_once() + _, kwargs = fs.materialize_incremental.call_args + assert kwargs.get("version") == "v2"