From 05f3870281df5c35b8afa90824a4fd1785dcf2ff Mon Sep 17 00:00:00 2001 From: Durable Workflow Date: Sat, 3 Oct 2026 08:29:26 +0000 Subject: [PATCH 1/3] Preserve configured worker polls after Server discovery --- src/durable_workflow/worker.py | 38 +++++----------------------------- tests/test_worker.py | 15 ++++++++++---- 2 files changed, 16 insertions(+), 37 deletions(-) diff --git a/src/durable_workflow/worker.py b/src/durable_workflow/worker.py index e639a3d..50f63fa 100644 --- a/src/durable_workflow/worker.py +++ b/src/durable_workflow/worker.py @@ -951,30 +951,6 @@ def _server_supports_update_validation_tasks(info: dict[str, Any]) -> bool: ) -def _server_long_poll_timeout(info: dict[str, Any]) -> float | None: - worker_protocol = info.get("worker_protocol") - if not isinstance(worker_protocol, dict): - return None - - capabilities = worker_protocol.get("server_capabilities") - if not isinstance(capabilities, dict): - return None - - timeout = capabilities.get("long_poll_timeout") - if isinstance(timeout, bool): - return None - if isinstance(timeout, int | float): - return float(timeout) if timeout > 0 else None - if isinstance(timeout, str): - try: - parsed = float(timeout) - except ValueError: - return None - return parsed if parsed > 0 else None - - return None - - def _contract_version_matches(value: Any, expected: int) -> bool: if isinstance(value, int): return value == expected @@ -1041,7 +1017,6 @@ def __init__( raise ValueError("heartbeat_interval must be positive") self._poll_timeout = poll_timeout - self._poll_http_timeout = poll_timeout self.max_concurrent_workflow_tasks = max_concurrent_workflow_tasks self.max_concurrent_activity_tasks = max_concurrent_activity_tasks self.max_concurrent_worker_sessions = max_concurrent_worker_sessions @@ -1191,9 +1166,6 @@ async def _register(self) -> None: "multiplexed workflow/update-validation polling. Refusing registration so validated " "updates cannot be accepted without validator approval or exceed worker capacity." ) - server_long_poll_timeout = _server_long_poll_timeout(info) - if server_long_poll_timeout is not None: - self._poll_http_timeout = max(self._poll_http_timeout, server_long_poll_timeout + 5.0) log.debug( "server compatibility accepted: app_version=%s control_plane=%s worker_protocol=%s", info.get("version", "unknown"), @@ -2530,7 +2502,7 @@ async def _poll_workflow_tasks(self) -> None: task = await self.client.poll_workflow_task( worker_id=self.worker_id, task_queue=self.task_queue, - timeout=self._poll_http_timeout, + timeout=self._poll_timeout, build_id=self.build_id, task_kinds=self._workflow_poll_task_kinds(), history_page_size=WORKFLOW_HISTORY_PAGE_SIZE, @@ -2636,7 +2608,7 @@ async def _poll_activity_tasks(self) -> None: task = await self.client.poll_activity_task( worker_id=self.worker_id, task_queue=self.task_queue, - timeout=self._poll_http_timeout, + timeout=self._poll_timeout, build_id=self.build_id, ) except asyncio.CancelledError: @@ -2695,7 +2667,7 @@ async def _poll_query_tasks(self, *, client: Client | None = None, track_tasks: task = await client.poll_query_task( worker_id=self.worker_id, task_queue=self.task_queue, - timeout=self._poll_http_timeout, + timeout=self._poll_timeout, build_id=self.build_id, ) except Exception as e: @@ -3220,7 +3192,7 @@ async def _run_until_loop( task = await self.client.poll_workflow_task( worker_id=self.worker_id, task_queue=self.task_queue, - timeout=self._poll_http_timeout, + timeout=self._poll_timeout, build_id=self.build_id, task_kinds=self._workflow_poll_task_kinds(), history_page_size=WORKFLOW_HISTORY_PAGE_SIZE, @@ -3269,7 +3241,7 @@ async def _run_until_loop( task = await self.client.poll_activity_task( worker_id=self.worker_id, task_queue=self.task_queue, - timeout=self._poll_http_timeout, + timeout=self._poll_timeout, build_id=self.build_id, ) if self._stop.is_set(): diff --git a/tests/test_worker.py b/tests/test_worker.py index 293b144..5176277 100644 --- a/tests/test_worker.py +++ b/tests/test_worker.py @@ -711,7 +711,14 @@ async def test_worker_without_validators_does_not_poll_validation_queue(self, mo ) @pytest.mark.asyncio - async def test_register_keeps_http_timeout_above_server_long_poll(self, mock_client: AsyncMock) -> None: + @pytest.mark.parametrize(("poll_loop", "poll_method"), [ + ("_poll_workflow_tasks", "poll_workflow_task"), + ("_poll_activity_tasks", "poll_activity_task"), + ("_poll_query_tasks", "poll_query_task"), + ]) + async def test_registration_preserves_configured_poll_window( + self, mock_client: AsyncMock, poll_loop: str, poll_method: str, + ) -> None: mock_client.get_cluster_info = AsyncMock( return_value=compatible_cluster_info( worker_protocol={ @@ -736,12 +743,12 @@ async def poll_once(**_: object) -> None: worker._stop.set() return None - mock_client.poll_workflow_task.side_effect = poll_once + getattr(mock_client, poll_method).side_effect = poll_once await worker._register() - await worker._poll_workflow_tasks() + await getattr(worker, poll_loop)() - assert mock_client.poll_workflow_task.call_args.kwargs["timeout"] == 17.0 + assert getattr(mock_client, poll_method).call_args.kwargs["timeout"] == 0.01 @pytest.mark.asyncio async def test_register_keeps_baseline_capabilities_when_server_does_not_support_query_tasks( From 1cf31a2b5e861982fbf7e4d5849380021bca53a3 Mon Sep 17 00:00:00 2001 From: Durable Workflow Date: Sat, 3 Oct 2026 08:32:18 +0000 Subject: [PATCH 2/3] Keep the stable transport fix separate from poll-loop refactoring --- src/durable_workflow/worker.py | 12 +++++++----- 1 file changed, 7 insertions(+), 5 deletions(-) diff --git a/src/durable_workflow/worker.py b/src/durable_workflow/worker.py index 50f63fa..8046e82 100644 --- a/src/durable_workflow/worker.py +++ b/src/durable_workflow/worker.py @@ -1017,6 +1017,8 @@ def __init__( raise ValueError("heartbeat_interval must be positive") self._poll_timeout = poll_timeout + # Client supplies HTTP grace separately from this requested poll window. + self._poll_http_timeout = poll_timeout self.max_concurrent_workflow_tasks = max_concurrent_workflow_tasks self.max_concurrent_activity_tasks = max_concurrent_activity_tasks self.max_concurrent_worker_sessions = max_concurrent_worker_sessions @@ -2502,7 +2504,7 @@ async def _poll_workflow_tasks(self) -> None: task = await self.client.poll_workflow_task( worker_id=self.worker_id, task_queue=self.task_queue, - timeout=self._poll_timeout, + timeout=self._poll_http_timeout, build_id=self.build_id, task_kinds=self._workflow_poll_task_kinds(), history_page_size=WORKFLOW_HISTORY_PAGE_SIZE, @@ -2608,7 +2610,7 @@ async def _poll_activity_tasks(self) -> None: task = await self.client.poll_activity_task( worker_id=self.worker_id, task_queue=self.task_queue, - timeout=self._poll_timeout, + timeout=self._poll_http_timeout, build_id=self.build_id, ) except asyncio.CancelledError: @@ -2667,7 +2669,7 @@ async def _poll_query_tasks(self, *, client: Client | None = None, track_tasks: task = await client.poll_query_task( worker_id=self.worker_id, task_queue=self.task_queue, - timeout=self._poll_timeout, + timeout=self._poll_http_timeout, build_id=self.build_id, ) except Exception as e: @@ -3192,7 +3194,7 @@ async def _run_until_loop( task = await self.client.poll_workflow_task( worker_id=self.worker_id, task_queue=self.task_queue, - timeout=self._poll_timeout, + timeout=self._poll_http_timeout, build_id=self.build_id, task_kinds=self._workflow_poll_task_kinds(), history_page_size=WORKFLOW_HISTORY_PAGE_SIZE, @@ -3241,7 +3243,7 @@ async def _run_until_loop( task = await self.client.poll_activity_task( worker_id=self.worker_id, task_queue=self.task_queue, - timeout=self._poll_timeout, + timeout=self._poll_http_timeout, build_id=self.build_id, ) if self._stop.is_set(): From 257995dab46ae75654c27989f4cb372666107611 Mon Sep 17 00:00:00 2001 From: Durable Workflow Date: Sat, 3 Oct 2026 08:38:59 +0000 Subject: [PATCH 3/3] Prepare Python SDK 2.3.8 poll-window release --- CHANGELOG.md | 7 +++++++ pyproject.toml | 6 +++--- 2 files changed, 10 insertions(+), 3 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 912ea70..ac1c48b 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -6,6 +6,13 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] +## [2.3.8] - 2026-10-03 + +### Fixed +- Preserve the configured worker poll window after Server discovery for workflow, + activity and query tasks. HTTP grace remains separate, so short polls release + Server admission slots when requested. + ## [2.3.7] - 2026-09-29 ### Fixed diff --git a/pyproject.toml b/pyproject.toml index 47b1ad5..692e769 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta" [project] name = "durable-workflow" -version = "2.3.7" +version = "2.3.8" description = "Python client and worker SDK for Durable Workflow Cloud and self-hosted Server" readme = "README.md" requires-python = ">=3.10" @@ -71,8 +71,8 @@ durable-workflow-replay-conformance = "durable_workflow.replay_conformance:main" durable-workflow-workflow-updates-conformance = "durable_workflow.workflow_updates_conformance:main" [tool.durable-workflow] -product-train = "2.3.7" -registry-version = "2.3.7" +product-train = "2.3.8" +registry-version = "2.3.8" supported-server-versions = "2.4.8" worker-protocol-version = "1.19" control-plane-version = "2"