Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
10 changes: 9 additions & 1 deletion temporalio/worker/_worker.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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")
Expand Down
32 changes: 32 additions & 0 deletions tests/worker/test_worker.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down
Loading