diff --git a/CHANGELOG.md b/CHANGELOG.md index 7d68251aa..ab5728c76 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -42,6 +42,8 @@ to include examples, links to docs, or any other relevant information. - `temporalio.contrib.strands` activity and MCP tools now give the model the Activity's failure message and expose its exception to after-tool hooks. +- Reject unsupported `InterpreterPoolExecutor` activity executors when creating + a worker instead of failing during activity execution. - Encoding a datetime search attribute without a timezone now raises `ValueError("Timezone must be present on all search attribute dates")` on diff --git a/temporalio/worker/_worker.py b/temporalio/worker/_worker.py index 60f824c4d..39faf615a 100644 --- a/temporalio/worker/_worker.py +++ b/temporalio/worker/_worker.py @@ -172,7 +172,8 @@ def __init__( activity_executor: Concurrent executor to use for non-async activities. This is required if any activities are non-async. :py:class:`concurrent.futures.ThreadPoolExecutor` is - recommended. If this is a + recommended. ``concurrent.futures.InterpreterPoolExecutor`` + is not supported. If this is a :py:class:`concurrent.futures.ProcessPoolExecutor`, all non-async activities must be picklable. ``max_workers`` on the executor should at least be ``max_concurrent_activities`` or a @@ -417,6 +418,13 @@ def _init_from_config(self, client: temporalio.client.Client, config: WorkerConf Client is safe to take separately since it can't be modified by worker plugins. """ self._config = config + if sys.version_info >= (3, 14) and isinstance( + config.get("activity_executor"), concurrent.futures.InterpreterPoolExecutor + ): + raise ValueError( # pyright: ignore[reportUnreachable] + "InterpreterPoolExecutor is not supported as an activity_executor. " + "Use ThreadPoolExecutor or ProcessPoolExecutor instead." + ) if not ( config.get("activities") or config.get("nexus_service_handlers") diff --git a/tests/worker/test_worker.py b/tests/worker/test_worker.py index 7aa666cd4..1c796cd6d 100644 --- a/tests/worker/test_worker.py +++ b/tests/worker/test_worker.py @@ -5,11 +5,13 @@ import multiprocessing import multiprocessing.context import os +import sys import uuid from collections.abc import Awaitable, Callable, Sequence from contextlib import contextmanager from datetime import timedelta from typing import Any +from unittest.mock import Mock from urllib.request import urlopen import nexusrpc @@ -77,6 +79,36 @@ def test_load_default_worker_binary_id(): assert val1 == val2 +@pytest.mark.skipif( + sys.version_info < (3, 14), + reason="InterpreterPoolExecutor requires Python 3.14 or newer", +) +@pytest.mark.parametrize("subclass", [False, True]) +def test_activity_executor_rejects_interpreter_pool(subclass: bool): + if sys.version_info >= (3, 14): + + class CustomInterpreterPoolExecutor(concurrent.futures.InterpreterPoolExecutor): # pyright: ignore[reportUnreachable] + pass + + executor_class = concurrent.futures.InterpreterPoolExecutor + if subclass: + executor_class = CustomInterpreterPoolExecutor + + client = Mock(spec=Client) + client.config.return_value = {"plugins": []} + with executor_class(max_workers=1) as executor: + with pytest.raises( + ValueError, + match="InterpreterPoolExecutor is not supported as an activity_executor", + ): + Worker( + client, + task_queue="test-interpreter-pool", + activities=[never_run_activity], + activity_executor=executor, + ) + + @activity.defn async def never_run_activity() -> None: raise NotImplementedError