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
7 changes: 7 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,13 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

## [Unreleased]

## [2.4.3] - 2026-10-08

### Fixed
- Preserve accepted update handler failures as `UpdateFailed`, including the
failure message, HTTP status, Server response and workflow, run, update and
failure identities. Existing `InvalidArgument` catches remain compatible.

## [2.4.2] - 2026-10-07

### Fixed
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.4.2"
version = "2.4.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.4.2"
registry-version = "2.4.2"
product-train = "2.4.3"
registry-version = "2.4.3"
supported-server-versions = "2.5.0"
worker-protocol-version = "1.19"
control-plane-version = "2"
Expand Down
2 changes: 2 additions & 0 deletions src/durable_workflow/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -95,6 +95,7 @@
ServerError,
SignalFailed,
Unauthorized,
UpdateFailed,
UpdateRejected,
UpdateValidationFailed,
WorkflowAlreadyStarted,
Expand Down Expand Up @@ -385,6 +386,7 @@
"SignalFailed",
"TransportRetryPolicy",
"Unauthorized",
"UpdateFailed",
"UpdateRejected",
"UpdateValidationFailed",
"WorkflowAlreadyStarted",
Expand Down
4 changes: 3 additions & 1 deletion src/durable_workflow/client.py
Original file line number Diff line number Diff line change
Expand Up @@ -4349,7 +4349,9 @@ async def update_workflow(
:class:`~durable_workflow.errors.UpdateRejected` when the workflow's
validator rejects the update, or
:class:`~durable_workflow.errors.UpdateValidationFailed` when the
declared validation boundary cannot be enforced.
declared validation boundary cannot be enforced. An accepted update
whose handler fails raises :class:`~durable_workflow.errors.UpdateFailed`
with the failure message, Server response and durable identities.
"""
await self._require_update_wait_stage(wait_for or "accepted")
body: dict[str, Any] = {}
Expand Down
30 changes: 30 additions & 0 deletions src/durable_workflow/errors.py
Original file line number Diff line number Diff line change
Expand Up @@ -453,6 +453,29 @@ def __init__(
self.recorded_event_types = list(recorded_event_types)


class UpdateFailed(InvalidArgument):
"""An accepted workflow update failed during handler execution.

``body`` preserves the Server response, including its durable failure
identity. This inherits from :class:`InvalidArgument` to preserve catches
written for earlier SDK versions that mapped all HTTP 422 responses there.
"""

def __init__(self, message: str, *, status: int, body: dict[str, Any]) -> None:
super().__init__(message)
self.status = status
self.body = body
self.workflow_id = self._identity(body, "workflow_id")
self.run_id = self._identity(body, "run_id")
self.update_id = self._identity(body, "update_id")
self.failure_id = self._identity(body, "failure_id")

@staticmethod
def _identity(body: dict[str, Any], name: str) -> str | None:
value = body.get(name)
return value if isinstance(value, str) and value.strip() else None


class UpdateRejected(DurableWorkflowError):
"""A workflow update was rejected by the workflow's validator."""

Expand Down Expand Up @@ -859,6 +882,13 @@ def update_validation_failed(default: str) -> UpdateValidationFailed:
raise signal_failed("signal argument validation failed")
if reason == "invalid_query_arguments":
raise query_failed("query argument validation failed")
if isinstance(body, dict) and body.get("update_status") == "failed":
failure_message = body.get("failure_message")
if not isinstance(failure_message, str) or not failure_message.strip():
failure_message = message
if not isinstance(failure_message, str) or not failure_message.strip():
failure_message = "workflow update failed"
raise UpdateFailed(failure_message, status=status, body=body)
errors = None
if isinstance(body, dict):
errors = body.get("errors") or body.get("validation_errors")
Expand Down
29 changes: 29 additions & 0 deletions tests/test_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,7 @@
ServerError,
SignalFailed,
Unauthorized,
UpdateFailed,
UpdateRejected,
WorkflowAlreadyStarted,
WorkflowNotFound,
Expand Down Expand Up @@ -3090,6 +3091,34 @@ async def test_heartbeat(self, client: Client) -> None:


class TestUpdateWorkflow:
@pytest.mark.asyncio
async def test_failed_update_retries_preserve_failure_identity(self, client: Client) -> None:
body = {
"update_status": "failed", "accepted": True,
"failure_message": "inventory unavailable", "workflow_id": "wf-1",
"run_id": "run-1", "update_id": "update-1", "failure_id": "failure-1",
}
with patch.object(
client._http, "request", new_callable=AsyncMock, return_value=_mock_response(422, body),
) as request:
failures = []
for _ in range(2):
with pytest.raises(UpdateFailed) as exc_info:
await client.update_workflow(
"wf-1", "reserve", wait_for="completed", request_id="request-1",
)
failures.append(exc_info.value)
assert request.await_count == 2
assert all(call.kwargs["json"]["request_id"] == "request-1"
for call in request.await_args_list)
for error in failures:
assert str(error) == "inventory unavailable"
assert error.status == 422
assert error.body == body
assert (error.workflow_id, error.run_id, error.update_id, error.failure_id) == (
"wf-1", "run-1", "update-1", "failure-1",
)

@pytest.mark.asyncio
async def test_unsupported_wait_stage_raises_typed_error_before_update(self) -> None:
client = Client("http://localhost:8080", token="test-token")
Expand Down
53 changes: 53 additions & 0 deletions tests/test_errors.py
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@
ServerError,
SignalFailed,
Unauthorized,
UpdateFailed,
UpdateRejected,
UpdateValidationFailed,
WorkflowAlreadyStarted,
Expand Down Expand Up @@ -102,13 +103,64 @@ def test_422_invalid(self) -> None:
with pytest.raises(InvalidArgument) as exc_info:
_raise_for_status(422, {"message": "bad", "errors": {"f": ["req"]}})
assert exc_info.value.errors == {"f": ["req"]}
assert type(exc_info.value) is InvalidArgument

def test_422_accepted_update_failure_preserves_diagnostics_and_legacy_catch(self) -> None:
body = {
"update_status": "failed", "accepted": True,
"failure_message": "inventory unavailable", "message": "generic failure",
"workflow_id": "wf-1", "run_id": "run-1", "update_id": "update-1",
"failure_id": "failure-1",
}
with pytest.raises(InvalidArgument) as exc_info:
_raise_for_status(422, body)
error = exc_info.value
assert isinstance(error, UpdateFailed)
assert str(error) == "inventory unavailable"
assert error.status == 422
assert error.body is body
assert error.errors is None
assert (error.workflow_id, error.run_id, error.update_id, error.failure_id) == (
"wf-1", "run-1", "update-1", "failure-1",
)

@pytest.mark.parametrize("failure_message", [None, "", " ", 42])
@pytest.mark.parametrize("message", ["handler failed", None, "", " ", 42])
def test_update_failure_message_fallback(self, failure_message: object, message: object) -> None:
with pytest.raises(UpdateFailed) as exc_info:
_raise_for_status(422, {
"update_status": "failed", "failure_message": failure_message, "message": message,
})
assert str(exc_info.value) == (
message if message == "handler failed" else "workflow update failed"
)

def test_update_failure_missing_or_malformed_identities_remain_absent(self) -> None:
with pytest.raises(UpdateFailed) as exc_info:
_raise_for_status(422, {
"update_status": "failed", "workflow_id": 42, "run_id": " ", "update_id": None,
})
error = exc_info.value
assert (error.workflow_id, error.run_id, error.update_id, error.failure_id) == (
None, None, None, None,
)

def test_rejected_update_arguments_are_not_handler_failure(self) -> None:
with pytest.raises(InvalidArgument) as exc_info:
_raise_for_status(422, {
"update_status": "rejected", "message": "invalid arguments",
"validation_errors": {"quantity": ["must be positive"]},
})
assert type(exc_info.value) is InvalidArgument
assert exc_info.value.errors == {"quantity": ["must be positive"]}

def test_422_update_validator_rejected_is_typed(self) -> None:
with pytest.raises(UpdateRejected) as exc_info:
_raise_for_status(
422,
{
"reason": "update_validator_rejected",
"update_status": "failed",
"message": "approval required",
"validation_errors": {"approved": ["must be true"]},
},
Expand All @@ -132,6 +184,7 @@ def test_update_validation_infrastructure_failures_are_typed(
status,
{
"reason": reason,
"update_status": "failed",
"message": "validation could not complete",
"retryable": retryable,
},
Expand Down
30 changes: 29 additions & 1 deletion tests/test_sync.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,7 @@
import httpx
import pytest

from durable_workflow import serializer
from durable_workflow import InvalidArgument, UpdateFailed, serializer
from durable_workflow.client import WorkflowHandle
from durable_workflow.external_storage import ExternalPayloadCache, LocalFilesystemExternalStorage
from durable_workflow.sync import Client, SyncStandaloneActivityHandle, SyncWorkflowHandle
Expand Down Expand Up @@ -581,6 +581,34 @@ def test_archive_workflow(self) -> None:


class TestSyncClientUpdate:
def test_failed_update_preserves_diagnostics(self) -> None:
body = {
"update_status": "failed", "failure_message": "inventory unavailable",
"workflow_id": "wf-1", "run_id": "run-1", "update_id": "update-1",
"failure_id": "failure-1",
}
with Client("http://localhost:8080") as client:
client._async._cluster_info = {
"control_plane": {"request_contract": {"operations": {"update": {
"fields": {"wait_for": {"canonical_values": ["accepted", "completed"]}},
}}}},
}
with (
patch.object(client._async._http, "request", new_callable=AsyncMock,
return_value=_mock_response(422, body)) as request,
pytest.raises(InvalidArgument) as exc_info,
):
client.update_workflow("wf-1", "reserve", request_id="request-1")
error = exc_info.value
assert isinstance(error, UpdateFailed)
assert str(error) == "inventory unavailable"
assert error.status == 422
assert error.body == body
assert (error.workflow_id, error.run_id, error.update_id, error.failure_id) == (
"wf-1", "run-1", "update-1", "failure-1",
)
assert request.call_args.kwargs["json"]["request_id"] == "request-1"

def test_update(self) -> None:
client = Client("http://localhost:8080")
client._async._cluster_info = {
Expand Down
Loading