diff --git a/CHANGELOG.md b/CHANGELOG.md index 3f6f779..bb1d20c 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -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 diff --git a/pyproject.toml b/pyproject.toml index 8f4060d..2ce3e73 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -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" @@ -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" diff --git a/src/durable_workflow/client.py b/src/durable_workflow/client.py index 715e3b2..3b2dd9a 100644 --- a/src/durable_workflow/client.py +++ b/src/durable_workflow/client.py @@ -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 @@ -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 @@ -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: @@ -277,6 +299,7 @@ class _RuntimeExternalPayloadTransport: max_payload_bytes: int request_timeout_seconds: float status: str + completion_context: bool = False @dataclass @@ -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() @@ -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 @@ -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 ( @@ -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 @@ -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()} @@ -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() } @@ -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 @@ -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({ @@ -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 diff --git a/src/durable_workflow/retry_policy.py b/src/durable_workflow/retry_policy.py index e90abc5..dd685fc 100644 --- a/src/durable_workflow/retry_policy.py +++ b/src/durable_workflow/retry_policy.py @@ -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 diff --git a/tests/test_release_metadata.py b/tests/test_release_metadata.py index 87a7dd4..cf09012 100644 --- a/tests/test_release_metadata.py +++ b/tests/test_release_metadata.py @@ -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 diff --git a/tests/test_runtime_external_payload_transport.py b/tests/test_runtime_external_payload_transport.py index 99a18de..9983f7e 100644 --- a/tests/test_runtime_external_payload_transport.py +++ b/tests/test_runtime_external_payload_transport.py @@ -19,6 +19,7 @@ ExternalPayloadUnavailable, ExternalPayloadUnsupported, RuntimeCapabilityUnsupported, + ServerError, ) from durable_workflow.external_storage import ( RUNTIME_EXTERNAL_PAYLOAD_REFERENCE_SCHEMA, @@ -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, *, diff --git a/tests/test_storage_admission.py b/tests/test_storage_admission.py index 0aa138a..dcc6925 100644 --- a/tests/test_storage_admission.py +++ b/tests/test_storage_admission.py @@ -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] = [] @@ -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: