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

## [Unreleased]

## [2.0.4] - 2026-09-08

### Fixed
- On supporting Servers, retry a draining refusal for a completion payload
with its exact activity, workflow, or query lease and immutable payload slot.
Client uploads, unknown capabilities, and hard storage fences remain blocked.
- Preserve computed outcomes during late upload pressure. Existing storage
retries stay interruptible and do not execute handlers again; Server upload
allowances, namespace quotas, and stale-lease errors remain authoritative.

## [2.0.3] - 2026-09-08

### 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.0.3"
version = "2.0.4"
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.0.3"
registry-version = "2.0.3"
product-train = "2.0.4"
registry-version = "2.0.4"
supported-server-versions = "2.0.0"
worker-protocol-version = "1.19"
control-plane-version = "2"
Expand Down
63 changes: 60 additions & 3 deletions src/durable_workflow/client.py
Original file line number Diff line number Diff line change
Expand Up @@ -20,8 +20,10 @@

import asyncio
import hashlib
import json as json_module
import math
import os
import re
import time
import uuid
import warnings
Expand All @@ -30,7 +32,7 @@
from importlib.metadata import PackageNotFoundError
from importlib.metadata import version as _pkg_version
from typing import Any
from urllib.parse import quote, urlencode, urlsplit
from urllib.parse import quote, unquote, urlencode, urlsplit

import httpx

Expand Down Expand Up @@ -101,6 +103,26 @@
"/external-payloads/v1/{referenceId}"
)
_RUNTIME_EXTERNAL_PAYLOAD_ERROR_BODY_LIMIT = 64 * 1024
_PAYLOAD_COMPLETION_SCHEMA = "durable-workflow.v2.payload-completion-context.v1"
_PAYLOAD_COMPLETION_HEADER = "X-Durable-Workflow-Payload-Completion"


def _payload_completion_context(path: str, body: Any) -> dict[str, Any] | None:
match = re.fullmatch(r"/worker/(activity|workflow|query)-tasks/([^/]+)/(complete|fail)", path.split("?")[0])
if match is None or not isinstance(body, dict):
return None
kind, task_id, operation = match.groups()
attempt = body.get("activity_attempt_id" if kind == "activity" else f"{kind}_task_attempt")
owner = body.get("lease_owner")
if not isinstance(owner, str) or not owner:
return None
if kind == "activity":
if not isinstance(attempt, str) or not attempt:
return None
elif type(attempt) is not int or attempt < 1:
return None
return {"schema": _PAYLOAD_COMPLETION_SCHEMA, "kind": kind, "task_id": unquote(task_id),
"attempt": attempt, "lease_owner": owner, "operation": operation}


def _default_sdk_version() -> str:
Expand Down Expand Up @@ -277,6 +299,7 @@ class _RuntimeExternalPayloadTransport:
max_payload_bytes: int
request_timeout_seconds: float
status: str
completion_context: bool = False


@dataclass
Expand Down Expand Up @@ -1734,6 +1757,9 @@ async def _request(
worker=worker,
transport=transport,
uploaded={},
completion=(
_payload_completion_context(path, json) if worker and transport.completion_context else None
),
)

start = time.perf_counter()
Expand Down Expand Up @@ -1904,11 +1930,15 @@ async def _runtime_external_payload_transport(
)

status = policy.get("status")
completion = upload.get("completion_context")
transport = _RuntimeExternalPayloadTransport(
threshold_bytes=threshold_bytes,
max_payload_bytes=max_payload_bytes,
request_timeout_seconds=float(request_timeout_seconds),
status=status if isinstance(status, str) else "unknown",
completion_context=isinstance(completion, dict)
and completion.get("schema") == _PAYLOAD_COMPLETION_SCHEMA
and completion.get("header") == _PAYLOAD_COMPLETION_HEADER,
)
self._runtime_external_payload_transport_cache = transport
self._runtime_external_payload_transport_resolved = True
Expand Down Expand Up @@ -1972,6 +2002,8 @@ async def _externalize_runtime_payloads(
worker: bool,
transport: _RuntimeExternalPayloadTransport,
uploaded: dict[tuple[str, str, int], RuntimeExternalPayloadReference],
completion: dict[str, Any] | None = None,
slot: tuple[str | int, ...] = (),
) -> Any:
if isinstance(value, dict):
if (
Expand All @@ -1984,6 +2016,8 @@ async def _externalize_runtime_payloads(
worker=worker,
transport=transport,
uploaded=uploaded,
completion=completion,
slot=(*slot, "result"),
)
if "external_payload" in externalized_result:
normalized_command["result"] = externalized_result
Expand Down Expand Up @@ -2014,6 +2048,7 @@ async def _externalize_runtime_payloads(
sha256=sha256,
worker=worker,
transport=transport,
completion={**completion, "slot": list(slot)} if completion is not None else None,
)
uploaded[identity] = reference
return {"codec": codec, "external_payload": reference.to_dict()}
Expand All @@ -2024,6 +2059,8 @@ async def _externalize_runtime_payloads(
worker=worker,
transport=transport,
uploaded=uploaded,
completion=completion,
slot=(*slot, key),
)
for key, item in value.items()
}
Expand All @@ -2045,8 +2082,10 @@ async def _externalize_runtime_payloads(
worker=worker,
transport=transport,
uploaded=uploaded,
completion=completion,
slot=(*slot, index),
)
for item in value
for index, item in enumerate(value)
]
return value

Expand All @@ -2058,6 +2097,7 @@ async def _upload_runtime_payload(
sha256: str,
worker: bool,
transport: _RuntimeExternalPayloadTransport,
completion: dict[str, Any] | None = None,
) -> RuntimeExternalPayloadReference:
headers = self._headers(worker=worker)
headers.update({
Expand All @@ -2069,13 +2109,30 @@ async def _upload_runtime_payload(
})

async def _do_request() -> httpx.Response:
attempt_headers = dict(headers)
response = await self._http.request(
"POST",
f"/api{_RUNTIME_EXTERNAL_PAYLOAD_UPLOAD_PATH}",
headers=headers,
headers=attempt_headers,
content=data,
timeout=transport.request_timeout_seconds,
)
if (response.status_code == 503 and completion is not None
and len(response.content) <= _RUNTIME_EXTERNAL_PAYLOAD_ERROR_BODY_LIMIT):
try:
refusal = response.json()
except ValueError:
refusal = None
context = json_module.dumps(completion, separators=(",", ":"))
if (isinstance(refusal, dict) and refusal.get("reason") == "storage_pressure"
and refusal.get("storage_state") == "draining" and len(context.encode("utf-8")) <= 4096):
# Keep the same identity if the ordinary transport policy retries
# an ambiguous response; never rerun application activity code.
attempt_headers[_PAYLOAD_COMPLETION_HEADER] = context
response = await self._http.request(
"POST", f"/api{_RUNTIME_EXTERNAL_PAYLOAD_UPLOAD_PATH}", headers=attempt_headers,
content=data, timeout=transport.request_timeout_seconds,
)
response.raise_for_status()
return response

Expand Down
6 changes: 6 additions & 0 deletions src/durable_workflow/retry_policy.py
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,12 @@ def _storage_refusal(exc: Exception) -> tuple[ServerError, str | None] | None:
body = exc.response.json()
except ValueError:
return None
# A payload upload is content-addressed and precedes completion submission.
# Even a late pressure refusal can retry those same bytes. This local retry
# classification does not alter the original response exposed to callers.
if (isinstance(body, dict) and "request_admitted" not in body and exc.request.method == "POST"
and exc.request.url.path.endswith("/api/external-payloads/v1")):
body = {**body, "request_admitted": False}
error = ServerError(exc.response.status_code, body)
if error.reason() not in ("storage_pressure", "storage_admission_unavailable"):
return None
Expand Down
4 changes: 2 additions & 2 deletions tests/test_release_metadata.py
Original file line number Diff line number Diff line change
Expand Up @@ -41,9 +41,9 @@ def test_worker_release_identity_matches_supported_server_and_protocol() -> None
project = manifest["project"]
release = manifest["tool"]["durable-workflow"]

assert project["version"] == "2.0.3"
assert project["version"] == "2.0.4"
assert release["product-train"] == project["version"]
assert release["registry-version"] == "2.0.3"
assert release["registry-version"] == "2.0.4"
assert release["supported-server-versions"] == "2.0.0"
assert release["worker-protocol-version"] == PROTOCOL_VERSION == "1.19"
assert release["durable-selection"] is True
Expand Down
110 changes: 110 additions & 0 deletions tests/test_runtime_external_payload_transport.py
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@
ExternalPayloadUnavailable,
ExternalPayloadUnsupported,
RuntimeCapabilityUnsupported,
ServerError,
)
from durable_workflow.external_storage import (
RUNTIME_EXTERNAL_PAYLOAD_REFERENCE_SCHEMA,
Expand Down Expand Up @@ -174,6 +175,115 @@ def handler(self, request: httpx.Request) -> httpx.Response:
return httpx.Response(200, json=response)


class CompletionPayloadServer(FakeRuntimePayloadServer):
def __init__(self, *, supported: bool = True, state: str = "draining", reject_bound: bool = False) -> None:
super().__init__()
self.supported = supported
self.state = state
self.reject_bound = reject_bound
self.upload_requests: list[httpx.Request] = []

def cluster_info(self) -> dict[str, Any]:
info = super().cluster_info()
if self.supported:
info["namespace"]["external_payload_storage"]["transport"]["upload"]["completion_context"] = {
"schema": "durable-workflow.v2.payload-completion-context.v1",
"header": "X-Durable-Workflow-Payload-Completion",
}
return info

def handler(self, request: httpx.Request) -> httpx.Response:
if request.method == "POST" and request.url.path == "/api/external-payloads/v1":
self.upload_requests.append(request)
if self.state != "normal" and "X-Durable-Workflow-Payload-Completion" not in request.headers:
return httpx.Response(503, json={"reason": "storage_pressure", "storage_state": self.state,
"request_admitted": False, "retryable": True, "retry_after_seconds": 1})
if self.reject_bound:
return httpx.Response(409, json={"reason": "external_payload_completion_lease_rejected",
"retryable": False, "message": "Lease rejected."})
return super().handler(request)


def completion_cases() -> list[tuple[str, str, dict[str, Any], list[str | int]]]:
envelope = serializer.envelope("x" * 100)
activity = {"lease_owner": "worker", "activity_attempt_id": "attempt"}
cases = [
("activity", "complete", {**activity, "result": envelope}, ["result"]),
("activity", "fail", {**activity, "failure": {"details": envelope}}, ["failure", "details"]),
("query", "complete", {"lease_owner": "worker", "query_task_attempt": 2, "result_envelope": envelope},
["result_envelope"]),
]
for kind, field in [("complete_workflow", "result"), ("schedule_activity", "arguments"),
("upsert_memo", "entries"), ("start_service_operation", "request_payload")]:
cases.append(("workflow", "complete", {"lease_owner": "worker", "workflow_task_attempt": 2,
"commands": [{"type": kind, field: envelope}]}, ["commands", 0, field]))
cases.append(("workflow", "complete", {"lease_owner": "worker", "workflow_task_attempt": 2,
"commands": [{"type": "fail_workflow", "exception": {"details": envelope}}]},
["commands", 0, "exception", "details"]))
cases.append(("workflow", "complete", {"lease_owner": "worker", "workflow_task_attempt": 2,
"commands": [{"type": "record_side_effect", "result": serializer.encode("x" * 100)}]},
["commands", 0, "result"]))
cases.append(("workflow", "complete", {"lease_owner": "worker", "workflow_task_attempt": 2,
"commands": [{"type": "record_side_effect", "workflow_stream": {"items": [
{"payload": envelope, "payload_codec": "avro"}]}}]},
["commands", 0, "workflow_stream", "items", 0, "payload"]))
return cases


@pytest.mark.parametrize("kind,operation,body,slot", completion_cases())
async def test_draining_upload_retry_carries_current_completion_identity(
kind: str, operation: str, body: dict[str, Any], slot: list[str | int],
) -> None:
server = CompletionPayloadServer()
async with runtime_client(server, retry_policy=TransportRetryPolicy(max_attempts=1)) as client:
await client._request("POST", f"/worker/{kind}-tasks/task/{operation}", worker=True, json=body)
first, bound = server.upload_requests
assert "X-Durable-Workflow-Payload-Completion" not in first.headers
assert json.loads(bound.headers["X-Durable-Workflow-Payload-Completion"]) == {
"schema": "durable-workflow.v2.payload-completion-context.v1", "kind": kind,
"task_id": "task", "attempt": "attempt" if kind == "activity" else 2,
"lease_owner": "worker", "operation": operation, "slot": slot,
}
assert first.content == bound.content
assert first.headers["authorization"] == bound.headers["authorization"]
assert first.headers["x-namespace"] == bound.headers["x-namespace"]
assert len(server.requests) == 1


@pytest.mark.parametrize("supported,worker,state", [(False, True, "draining"), (True, False, "draining"),
(True, True, "fenced")])
async def test_completion_upload_does_not_bypass_unsupported_client_or_fenced_admission(
supported: bool, worker: bool, state: str,
) -> None:
server = CompletionPayloadServer(supported=supported, state=state)
async with runtime_client(server, retry_policy=TransportRetryPolicy(max_attempts=1)) as client:
with pytest.raises((ServerError, ExternalPayloadError)):
await client._request("POST", "/worker/activity-tasks/task/complete" if worker else "/workflows",
worker=worker, json={"lease_owner": "worker", "activity_attempt_id": "attempt",
"result" if worker else "input": serializer.envelope("x" * 100)})
assert len(server.upload_requests) == 1
assert server.requests == []


async def test_rejected_bound_upload_does_not_loop_or_submit_completion() -> None:
server = CompletionPayloadServer(reject_bound=True)
async with runtime_client(server, retry_policy=TransportRetryPolicy(max_attempts=3)) as client:
with pytest.raises((ServerError, ExternalPayloadError)):
await client.complete_activity_task(task_id="task", activity_attempt_id="attempt", lease_owner="worker",
result="x" * 100)
assert len(server.upload_requests) == 2
assert server.requests == []


async def test_normal_payload_upload_is_unchanged_when_completion_capability_exists() -> None:
server = CompletionPayloadServer(state="normal")
async with runtime_client(server) as client:
await client.complete_activity_task(task_id="task", activity_attempt_id="attempt", lease_owner="worker",
result="x" * 100)
assert len(server.upload_requests) == 1
assert "X-Durable-Workflow-Payload-Completion" not in server.upload_requests[0].headers


def runtime_client(
server: FakeRuntimePayloadServer,
*,
Expand Down
8 changes: 6 additions & 2 deletions tests/test_storage_admission.py
Original file line number Diff line number Diff line change
Expand Up @@ -225,8 +225,9 @@ async def stop_on_sleep(delay: float) -> None:


@pytest.mark.parametrize("kind", ["activity", "query", "workflow"])
@pytest.mark.parametrize("late_upload_refusal", [False, True])
async def test_runtime_payload_upload_and_acknowledgement_are_not_repeated(
kind: str, retry_sleeps: list[float],
kind: str, retry_sleeps: list[float], late_upload_refusal: bool,
) -> None:
server = FakeRuntimePayloadServer()
uploads: list[bytes] = []
Expand All @@ -236,7 +237,10 @@ def handler(request: httpx.Request) -> httpx.Response:
if request.method == "POST" and request.url.path == "/api/external-payloads/v1":
uploads.append(request.content)
if len(uploads) <= 3:
return httpx.Response(503, json=pressure())
refusal = pressure()
if late_upload_refusal:
refusal.pop("request_admitted")
return httpx.Response(503, json=refusal)
if request.url.path.endswith("/complete"):
acknowledgements.append(request.content)
if len(acknowledgements) <= 4:
Expand Down