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
40 changes: 30 additions & 10 deletions polyglot/python_worker/scripts/sdk_updates.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down Expand Up @@ -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):
Expand Down
47 changes: 40 additions & 7 deletions polyglot/python_worker/tests/test_sdk_updates.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)


Expand Down Expand Up @@ -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):
Expand Down
2 changes: 1 addition & 1 deletion polyglot/qualified-artifact-tuple.json
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
4 changes: 4 additions & 0 deletions polyglot/updates/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Loading