Skip to content
Merged
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: 1 addition & 1 deletion CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
6 changes: 3 additions & 3 deletions pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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"
Expand Down
6 changes: 3 additions & 3 deletions src/durable_workflow/retry_policy.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down
6 changes: 3 additions & 3 deletions src/durable_workflow/worker.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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)
Expand Down
4 changes: 2 additions & 2 deletions tests/test_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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:
Expand Down
15 changes: 11 additions & 4 deletions tests/test_retry_policy.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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,
})
Expand Down
2 changes: 1 addition & 1 deletion tests/test_worker.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
Loading