From 2458e2bef29a45c099866170f524fe8c71ab4a17 Mon Sep 17 00:00:00 2001 From: Durable Workflow Date: Sat, 26 Sep 2026 09:51:55 +0000 Subject: [PATCH] Match Server poll task kinds and prepare Python SDK 2.3.3 --- CHANGELOG.md | 2 +- pyproject.toml | 6 +++--- src/durable_workflow/retry_policy.py | 6 +++--- src/durable_workflow/worker.py | 6 +++--- tests/test_client.py | 4 ++-- tests/test_retry_policy.py | 15 +++++++++++---- tests/test_worker.py | 2 +- 7 files changed, 24 insertions(+), 17 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 24b4dda..d20e0c9 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -6,7 +6,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] -## [2.3.2] - 2026-09-26 +## [2.3.3] - 2026-09-26 ### Fixed - Workers recognize Server's typed long-poll capacity response for workflow, diff --git a/pyproject.toml b/pyproject.toml index ad46e17..fea3f45 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta" [project] name = "durable-workflow" -version = "2.3.2" +version = "2.3.3" 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.2" -registry-version = "2.3.2" +product-train = "2.3.3" +registry-version = "2.3.3" supported-server-versions = "2.4.8" worker-protocol-version = "1.19" control-plane-version = "2" diff --git a/src/durable_workflow/retry_policy.py b/src/durable_workflow/retry_policy.py index 30df3c3..05f96c1 100644 --- a/src/durable_workflow/retry_policy.py +++ b/src/durable_workflow/retry_policy.py @@ -128,9 +128,9 @@ def _poll_capacity_refusal(exc: Exception) -> bool: if request.method != "POST" or "X-Durable-Workflow-Protocol-Version" not in request.headers: return False task_kinds = { - "/api/worker/workflow-tasks/poll": "workflow", - "/api/worker/activity-tasks/poll": "activity", - "/api/worker/query-tasks/poll": "query", + "/api/worker/workflow-tasks/poll": "workflow_task", + "/api/worker/activity-tasks/poll": "activity_task", + "/api/worker/query-tasks/poll": "query_task", } task_kind = next((kind for path, kind in task_kinds.items() if request.url.path.endswith(path)), None) if task_kind is None: diff --git a/src/durable_workflow/worker.py b/src/durable_workflow/worker.py index c1debbe..a29fbfb 100644 --- a/src/durable_workflow/worker.py +++ b/src/durable_workflow/worker.py @@ -2545,7 +2545,7 @@ async def _poll_workflow_tasks(self) -> None: return if _is_storage_admission_error(e): raise - delay = _poll_capacity_delay(e, "workflow", self.task_queue) + delay = _poll_capacity_delay(e, "workflow_task", self.task_queue) if delay is not None: self._record_poll_metrics("workflow", "backpressure", time.perf_counter() - poll_start) log.debug("workflow poll capacity backpressure; retrying in %d s", delay) @@ -2649,7 +2649,7 @@ async def _poll_activity_tasks(self) -> None: return if _is_storage_admission_error(e): raise - delay = _poll_capacity_delay(e, "activity", self.task_queue) + delay = _poll_capacity_delay(e, "activity_task", self.task_queue) if delay is not None: self._record_poll_metrics("activity", "backpressure", time.perf_counter() - poll_start) log.debug("activity poll capacity backpressure; retrying in %d s", delay) @@ -2706,7 +2706,7 @@ async def _poll_query_tasks(self, *, client: Client | None = None, track_tasks: return if _is_storage_admission_error(e): raise - delay = _poll_capacity_delay(e, "query", self.task_queue) + delay = _poll_capacity_delay(e, "query_task", self.task_queue) if delay is not None: self._record_poll_metrics("query", "backpressure", time.perf_counter() - poll_start) log.debug("query poll capacity backpressure; retrying in %d s", delay) diff --git a/tests/test_client.py b/tests/test_client.py index 1e54efd..6b2f9a3 100644 --- a/tests/test_client.py +++ b/tests/test_client.py @@ -2698,7 +2698,7 @@ async def refused(method: str, path: str, **kwargs: object) -> httpx.Response: "task": None, "poll_status": "long_poll_capacity_exhausted", "reason": "long_poll_capacity_exhausted", - "task_kind": "workflow", + "task_kind": "workflow_task", "task_queue": "q1", "retryable": True, "retry_after_seconds": 2, @@ -2708,7 +2708,7 @@ async def refused(method: str, path: str, **kwargs: object) -> httpx.Response: await client.poll_workflow_task(worker_id="worker-1", task_queue="q1") assert len(requests) == 1 - assert refusal.value.poll_capacity_backpressure_delay("workflow", "q1") == 2 + assert refusal.value.poll_capacity_backpressure_delay("workflow_task", "q1") == 2 @pytest.mark.asyncio async def test_poll_workflow_task_response_preserves_no_compatible_status(self, client: Client) -> None: diff --git a/tests/test_retry_policy.py b/tests/test_retry_policy.py index 25e6730..5fe9e6e 100644 --- a/tests/test_retry_policy.py +++ b/tests/test_retry_policy.py @@ -43,10 +43,10 @@ def test_should_retry_429_rate_limit(self) -> None: @pytest.mark.parametrize( ("path", "task_kind"), [ - ("/api/worker/workflow-tasks/poll", "workflow"), - ("/api/runtime/v1/namespaces/acme/api/worker/workflow-tasks/poll", "workflow"), - ("/api/worker/activity-tasks/poll", "activity"), - ("/api/worker/query-tasks/poll", "query"), + ("/api/worker/workflow-tasks/poll", "workflow_task"), + ("/api/runtime/v1/namespaces/acme/api/worker/workflow-tasks/poll", "workflow_task"), + ("/api/worker/activity-tasks/poll", "activity_task"), + ("/api/worker/query-tasks/poll", "query_task"), ], ) @pytest.mark.asyncio @@ -80,6 +80,13 @@ async def refused() -> None: await policy.execute(refused) assert calls == 1 + wrong_kind = httpx.Response(429, request=request, json={ + **response.json(), "task_kind": task_kind.removesuffix("_task"), + }) + assert policy.should_retry( + httpx.HTTPStatusError("wrong kind", request=request, response=wrong_kind), attempt=0, + ) is True + malformed = httpx.Response(429, request=request, json={ "reason": "long_poll_capacity_exhausted", "retry_after_seconds": 2, }) diff --git a/tests/test_worker.py b/tests/test_worker.py index 283342d..e6450c7 100644 --- a/tests/test_worker.py +++ b/tests/test_worker.py @@ -888,7 +888,7 @@ async def test_typed_poll_backpressure_waits_without_warning( "task": None, "poll_status": "long_poll_capacity_exhausted", "reason": "long_poll_capacity_exhausted", - "task_kind": task_kind, + "task_kind": f"{task_kind}_task", "task_queue": "q1", "retryable": True, "retry_after_seconds": 2,