diff --git a/polyglot/python_worker/scripts/sdk_updates.py b/polyglot/python_worker/scripts/sdk_updates.py index 026b395..77158c4 100644 --- a/polyglot/python_worker/scripts/sdk_updates.py +++ b/polyglot/python_worker/scripts/sdk_updates.py @@ -7,7 +7,7 @@ import json import os -from durable_workflow import Client, InvalidArgument, serializer +from durable_workflow import Client, UpdateFailed, serializer RUNTIMES = ("php", "python", "rust") DIRECTIONS = tuple((caller, runtime) for caller in RUNTIMES for runtime in RUNTIMES) @@ -175,24 +175,44 @@ async def replacement(client): async def failure(client): request_id = f"{required('DURABLE_WORKFLOW_UPDATE_ID')}-failure" - try: - await client.update_workflow(workflow_id("rust"), "fail", - args=[request("python", request_id)], wait_for="completed", request_id=request_id) - except InvalidArgument as error: - sdk_error = {"type": type(error).__name__, "message": str(error)} - else: - raise RuntimeError("Expected the published Python SDK's exception for a failed update.") + errors = [] + for _ in range(2): + try: + await client.update_workflow(workflow_id("rust"), "fail", + args=[request("python", request_id)], wait_for="completed", request_id=request_id) + except UpdateFailed as error: + errors.append(error) + else: + raise RuntimeError("Expected the published Python SDK's typed failed update.") execution, history = await observe(client, "rust") update_id, related = update_events(history, request_id) failed = [event for event in related if event["event_type"] == "UpdateCompleted" and event["payload"].get("failure_id")] if (len(failed) != 1 or "update-probe-failure" not in json.dumps(failed[0]["payload"]) or sum(event["event_type"] == "UpdateCompleted" for event in related) != 1): - raise RuntimeError(f"Missing durable handler failure: {sdk_error!r}, {related!r}") + raise RuntimeError(f"Missing durable handler failure: {related!r}") + payload = failed[0]["payload"] + expected = {"workflow_id": execution.workflow_id, "run_id": execution.run_id, + "update_id": update_id, "failure_id": payload["failure_id"]} + for error in errors: + if (error.status != 422 or str(error) != payload.get("message") + or "update-probe-failure" not in str(error) + or error.body.get("update_status") != "failed" + or error.body.get("command_status") != "accepted" or error.body.get("accepted") is False + or any(getattr(error, name) != value or error.body.get(name) != value + for name, value in expected.items())): + actual = {name: getattr(error, name) for name in expected} + actual.update(message=str(error), status=error.status, + update_status=error.body.get("update_status"), + command_status=error.body.get("command_status")) + raise RuntimeError(f"SDK failed-update diagnostics differ: expected={expected!r}, actual={actual!r}") if execution.status in ("failed", "terminated", "completed"): raise RuntimeError("An update failure unexpectedly terminated the workflow.") emit(scenario="handler-failure", runtime="rust", update_id=update_id, - sdk_error=sdk_error, events=related) + run_id=execution.run_id, sdk_error={"type": type(errors[0]).__name__, "message": str(errors[0]), + "status": errors[0].status, **expected}, + duplicate_sdk_error={"type": type(errors[1]).__name__, "message": str(errors[1]), + "status": errors[1].status, **expected}, events=related) async def finish(client): diff --git a/polyglot/python_worker/tests/test_sdk_updates.py b/polyglot/python_worker/tests/test_sdk_updates.py index 64933c0..7b73674 100644 --- a/polyglot/python_worker/tests/test_sdk_updates.py +++ b/polyglot/python_worker/tests/test_sdk_updates.py @@ -18,12 +18,17 @@ SCRIPT = pathlib.Path(__file__).parents[1] / "scripts" / "sdk_updates.py" spec = importlib.util.spec_from_file_location("sdk_updates", SCRIPT) updates = importlib.util.module_from_spec(spec) -class InvalidArgument(Exception): - pass +class UpdateFailed(Exception): + def __init__(self, message, *, status, body): + super().__init__(message) + self.status = status + self.body = body + for name in ("workflow_id", "run_id", "update_id", "failure_id"): + setattr(self, name, body.get(name)) with patch.dict(sys.modules, {"durable_workflow": types.SimpleNamespace( - Client=object, InvalidArgument=InvalidArgument, serializer=serializer)}): + Client=object, UpdateFailed=UpdateFailed, serializer=serializer)}): spec.loader.exec_module(updates) @@ -107,23 +112,51 @@ async def test_client_decodes_the_envelope_field(self): async def test_sdk_failure_requires_a_matching_durable_handler_failure(self): request_id = "example-failure" - client = types.SimpleNamespace(update_workflow=AsyncMock(side_effect=InvalidArgument("failed"))) - execution = types.SimpleNamespace(status="waiting") + body = {"update_status": "failed", "command_status": "accepted", "workflow_id": "example-rust", + "run_id": "run", "update_id": "original", "failure_id": "failure"} + client = types.SimpleNamespace(update_workflow=AsyncMock( + side_effect=UpdateFailed("update-probe-failure", status=422, body=body))) + execution = types.SimpleNamespace(status="waiting", workflow_id="example-rust", run_id="run") history = {"events": [ {"event_type": "UpdateAccepted", "payload": {"update_id": "original", "arguments": serializer.envelope([updates.request("python", request_id)])}}, {"event_type": "UpdateCompleted", "payload": {"update_id": "original", - "failure_id": "failure", "failure_message": "update-probe-failure"}}, + "failure_id": "failure", "message": "update-probe-failure"}}, ]} with patch.dict(os.environ, {"DURABLE_WORKFLOW_UPDATE_ID": "example"}): with patch.object(updates, "observe", AsyncMock(return_value=(execution, history))): with patch.object(updates, "emit") as emit: await updates.failure(client) self.assertEqual(emit.call_args.kwargs["update_id"], "original") - history["events"][-1]["payload"]["failure_message"] = "unrelated validation error" + self.assertEqual(client.update_workflow.await_count, 2) + self.assertEqual(emit.call_args.kwargs["sdk_error"], + emit.call_args.kwargs["duplicate_sdk_error"]) + history["events"][-1]["payload"]["message"] = "unrelated validation error" with self.assertRaisesRegex(RuntimeError, "durable handler failure"): await updates.failure(client) + async def test_sdk_failure_rejects_wrong_diagnostics(self): + body = {"update_status": "failed", "command_status": "accepted", "workflow_id": "example-rust", + "run_id": "run", "update_id": "original", "failure_id": "failure"} + execution = types.SimpleNamespace(status="waiting", workflow_id="example-rust", run_id="run") + history = {"events": [ + {"event_type": "UpdateAccepted", "payload": {"update_id": "original", + "arguments": serializer.envelope([updates.request("python", "example-failure")])}}, + {"event_type": "UpdateCompleted", "payload": {"update_id": "original", + "failure_id": "failure", "message": "update-probe-failure"}}, + ]} + for change in ({"message": ""}, {"status": 409}, {"failure_id": "other"}, + {"run_id": "replacement"}, {"workflow_id": "other"}, {"update_id": "other"}, + {"accepted": False}, {"update_status": "rejected"}, {"command_status": "rejected"}): + with self.subTest(change=change): + error = UpdateFailed(change.get("message", "update-probe-failure"), + status=change.get("status", 422), body={**body, **change}) + client = types.SimpleNamespace(update_workflow=AsyncMock(side_effect=error)) + with (patch.dict(os.environ, {"DURABLE_WORKFLOW_UPDATE_ID": "example"}), + patch.object(updates, "observe", AsyncMock(return_value=(execution, history)))): + with self.assertRaisesRegex(RuntimeError, "diagnostics"): + await updates.failure(client) + class HistoryPaginationTest(unittest.IsolatedAsyncioTestCase): async def test_observer_reads_later_pages(self): diff --git a/polyglot/qualified-artifact-tuple.json b/polyglot/qualified-artifact-tuple.json index 722a620..58aa071 100644 --- a/polyglot/qualified-artifact-tuple.json +++ b/polyglot/qualified-artifact-tuple.json @@ -5,7 +5,7 @@ "artifacts": { "cli": "2.1.3", "sdk-php": "2.2.3", - "sdk-python": "2.4.2", + "sdk-python": "2.4.3", "sdk-rust": "3.3.0", "server": "2.5.9", "waterline": "2.3.1", diff --git a/polyglot/updates/README.md b/polyglot/updates/README.md index 2e5e24c..bb0cff1 100644 --- a/polyglot/updates/README.md +++ b/polyglot/updates/README.md @@ -29,6 +29,10 @@ request returns that original completion. A failed Rust handler must persist a failed update while leaving its workflow live. All three original workflows finish once after a signal. +The Python client must raise `UpdateFailed` with the handler's message and +matching workflow, run, update and failure IDs. Repeating that failed request +must return the same error identities with one durable failed completion. + Rust does not support synchronous pre-accept update validators. The installed crate must return `UnsupportedUpdateValidators` for a contract claiming one. This check does not claim validator execution, process loss during an external