From ec0470ee0116576ea5e3a817ed0f188343173e87 Mon Sep 17 00:00:00 2001 From: Durable Workflow Date: Wed, 7 Oct 2026 23:23:38 +0000 Subject: [PATCH 1/7] Exercise published Rust workflow update clients and handlers --- .github/workflows/sdk-updates.yml | 29 +++ polyglot/php_worker/Dockerfile | 1 + polyglot/php_worker/update_client.php | 19 ++ polyglot/php_worker/worker.php | 11 + polyglot/python_worker/scripts/sdk_updates.py | 208 ++++++++++++++++++ .../python_worker/tests/test_sdk_updates.py | 84 +++++++ polyglot/python_workflow/workflow.py | 19 ++ polyglot/rust_worker/src/main.rs | 6 + polyglot/rust_worker/src/updates.rs | 41 ++++ polyglot/updates/README.md | 42 ++++ scripts/sdk-updates.sh | 86 ++++++++ 11 files changed, 546 insertions(+) create mode 100644 .github/workflows/sdk-updates.yml create mode 100644 polyglot/php_worker/update_client.php create mode 100644 polyglot/python_worker/scripts/sdk_updates.py create mode 100644 polyglot/python_worker/tests/test_sdk_updates.py create mode 100644 polyglot/rust_worker/src/updates.rs create mode 100644 polyglot/updates/README.md create mode 100755 scripts/sdk-updates.sh diff --git a/.github/workflows/sdk-updates.yml b/.github/workflows/sdk-updates.yml new file mode 100644 index 00000000..930d5ffb --- /dev/null +++ b/.github/workflows/sdk-updates.yml @@ -0,0 +1,29 @@ +name: published SDK updates + +on: + push: + branches: [main] + pull_request: + branches: [main] + workflow_dispatch: + +permissions: + contents: read + +concurrency: + group: sdk-updates-${{ github.event_name }}-${{ github.event_name == 'pull_request' && github.ref || github.sha }} + cancel-in-progress: ${{ github.event_name == 'pull_request' }} + +jobs: + updates: + name: SDK updates (Rust directions) + runs-on: ubuntu-latest + timeout-minutes: 25 + env: + SDK_UPDATES_COMPOSE_PROJECT_NAME: sample-app-sdk-updates-${{ github.run_id }}-${{ github.run_attempt }} + steps: + - uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1 + - name: Resolve the exact published tuple + run: scripts/resolve-current-artifacts.sh >> "$GITHUB_ENV" + - name: Run the documented update experiment + run: scripts/sdk-updates.sh diff --git a/polyglot/php_worker/Dockerfile b/polyglot/php_worker/Dockerfile index 2894e08e..fd512ca2 100644 --- a/polyglot/php_worker/Dockerfile +++ b/polyglot/php_worker/Dockerfile @@ -29,6 +29,7 @@ RUN test -n "$DURABLE_WORKFLOW_PHP_SDK_VERSION" \ COPY worker.php ./ COPY php_created_rust_schedule.php ./ +COPY update_client.php ./ COPY task_codec_rejection_probe.php ./ COPY replay_fixtures ./replay_fixtures diff --git a/polyglot/php_worker/update_client.php b/polyglot/php_worker/update_client.php new file mode 100644 index 00000000..2ff07d73 --- /dev/null +++ b/polyglot/php_worker/update_client.php @@ -0,0 +1,19 @@ + 'php', 'request_id' => $argv[2], 'value' => 'hello', 'nested' => ['enabled' => true, 'count' => 42]]; +$result = $client->updateWorkflow($argv[1], $argv[3], [$request], requestId: $argv[2]); +fwrite(STDOUT, json_encode(['caller' => 'php', 'request_id' => $argv[2], 'result' => $result], JSON_THROW_ON_ERROR).PHP_EOL); diff --git a/polyglot/php_worker/worker.php b/polyglot/php_worker/worker.php index a5a16c81..25d2ac5f 100644 --- a/polyglot/php_worker/worker.php +++ b/polyglot/php_worker/worker.php @@ -24,6 +24,7 @@ const WORKFLOW_TYPES = [ 'polyglot.php.greeter', 'polyglot.php.timer', + 'polyglot.php.updates', 'polyglot.PolyglotWorkflow', 'polyglot.php-to-python.greeter', 'polyglot.php-to-python.type-roundtrip', @@ -445,6 +446,16 @@ function signalQueryWorkflow(): Closure function configureWorkflows(Worker $worker, PayloadCodec $codec): void { + $worker->registerWorkflow('polyglot.php.updates', static function (WorkflowContext $context, string $request): array { + $context->waitCondition(static fn (): bool => $context->signals('updates-finish') !== [], 'updates-finish'); + + return ['workflow_runtime' => 'php', 'request' => $request]; + }); + $worker->declareSignal('polyglot.php.updates', 'updates-finish', static function (): void {}); + $worker->registerUpdate('polyglot.php.updates', 'echo', static function (QueryContext $context, array $request): array { + return ['handler_runtime' => 'php', 'request' => $request]; + }); + $worker->registerWorkflow('polyglot.php.timer', static function (WorkflowContext $context, string $request): array { $context->sleep(30); diff --git a/polyglot/python_worker/scripts/sdk_updates.py b/polyglot/python_worker/scripts/sdk_updates.py new file mode 100644 index 00000000..f19271ef --- /dev/null +++ b/polyglot/python_worker/scripts/sdk_updates.py @@ -0,0 +1,208 @@ +"""Observe actual published SDK update clients and workers in the Compose stack.""" + +from __future__ import annotations + +import argparse +import asyncio +import json +import os + +from durable_workflow import Client, serializer + +RUNTIMES = ("php", "python", "rust") +DIRECTIONS = (("php", "rust"), ("python", "rust"), ("rust", "rust"), + ("rust", "php"), ("rust", "python")) + + +def required(name): + value = os.environ.get(name, "").strip() + if not value: + raise RuntimeError(f"Set {name} before running SDK updates.") + return value + + +def emit(**record): + print(json.dumps(record, sort_keys=True), flush=True) + + +def records(name): + return [json.loads(line) for line in required(name).splitlines()] + + +def workflow_id(runtime): + return f"{required('DURABLE_WORKFLOW_UPDATE_ID')}-{runtime}" + + +def request(caller, request_id): + return {"caller": caller, "request_id": request_id, "value": "hello", + "nested": {"enabled": True, "count": 42}} + + +def update_events(history, request_id): + accepted = [] + for event in history["events"]: + if event["event_type"] != "UpdateAccepted": + continue + args = serializer.decode_envelope(event["payload"]["arguments"]) + if args and isinstance(args[0], dict) and args[0].get("request_id") == request_id: + accepted.append(event) + if len(accepted) != 1: + raise RuntimeError(f"Expected one accepted update for {request_id}: {accepted!r}") + update_id = accepted[0]["payload"]["update_id"] + related = [event for event in history["events"] if event["payload"].get("update_id") == update_id] + return update_id, related + + +def verify_completed(history, request_id, runtime, caller): + update_id, related = update_events(history, request_id) + completed = [event for event in related if event["event_type"] == "UpdateCompleted"] + if len(completed) != 1 or completed[0]["payload"].get("failure_id"): + raise RuntimeError(f"Expected one completed update: {related!r}") + types = [event["event_type"] for event in related] + if types.index("UpdateAccepted") >= types.index("UpdateCompleted"): + raise RuntimeError("Update completion precedes acceptance.") + expected = {"handler_runtime": runtime, "request": request(caller, request_id)} + result = serializer.decode_envelope(completed[0]["payload"]["result"]) + if result != expected: + raise RuntimeError(f"Persisted result changed: {result!r}") + return update_id, related, expected + + +async def observe(client, runtime): + execution = await client.describe_workflow(workflow_id(runtime)) + if not execution.run_id: + raise RuntimeError("Workflow has no durable run identity.") + originals = records("DURABLE_WORKFLOW_UPDATE_RUNS") + original = next(record for record in originals if record["runtime"] == runtime) + if execution.run_id != original["run_id"]: + raise RuntimeError("Update experiment replaced its original run.") + return execution, await client.get_history(execution.workflow_id, execution.run_id) + + +async def start(client): + for runtime in RUNTIMES: + handle = await client.start_workflow(workflow_type=f"polyglot.{runtime}.updates", + workflow_id=workflow_id(runtime), + task_queue=f"polyglot-{runtime}", input=[workflow_id(runtime)]) + deadline = asyncio.get_running_loop().time() + 30 + while True: + execution = await handle.describe() + if execution.status == "waiting" and execution.run_id: + break + if execution.status in ("failed", "completed", "terminated"): + raise RuntimeError(f"Update workflow did not wait: {execution!r}") + if asyncio.get_running_loop().time() >= deadline: + raise TimeoutError(f"{runtime} did not reach its signal wait.") + await asyncio.sleep(.25) + emit(runtime=runtime, workflow_id=execution.workflow_id, run_id=execution.run_id) + + +async def call(client, target, request_id, name): + response = await client.update_workflow(target, name, args=[request("python", request_id)], + wait_for="completed", request_id=request_id) + if response.get("update_status") != "completed": + raise RuntimeError(f"Update did not complete: {response!r}") + emit(caller="python", request_id=request_id, + result=serializer.decode_envelope(response["result"])) + + +async def matrix(client): + results = records("DURABLE_WORKFLOW_UPDATE_RESULTS") + if len(results) != len(DIRECTIONS): + raise RuntimeError("Not all five client/handler directions executed.") + for caller, runtime in DIRECTIONS: + request_id = f"{required('DURABLE_WORKFLOW_UPDATE_ID')}-{caller}-{runtime}" + matches = [result for result in results if result.get("request_id") == request_id] + if len(matches) != 1 or matches[0].get("caller") != caller: + raise RuntimeError("Missing or duplicate SDK client observation.") + execution, history = await observe(client, runtime) + update_id, related, expected = verify_completed(history, request_id, runtime, caller) + if matches[0].get("result") != expected: + raise RuntimeError(f"SDK client result differs from persisted result: {matches[0]!r}") + emit(scenario="client-handler", caller=caller, runtime=runtime, run_id=execution.run_id, + update_id=update_id, result=expected, events=related) + + +async def queued(client): + request_id = f"{required('DURABLE_WORKFLOW_UPDATE_ID')}-queued" + response = await client.update_workflow(workflow_id("rust"), "echo", + args=[request("python", request_id)], + wait_for="accepted", request_id=request_id) + execution, history = await observe(client, "rust") + update_id, related = update_events(history, request_id) + if response.get("update_status") != "accepted" or response.get("update_id") != update_id: + raise RuntimeError(f"Missing durable acceptance while worker is absent: {response!r}") + if any(event["event_type"] in ("UpdateCompleted", "UpdateFailed") for event in related): + raise RuntimeError("Update settled while its worker should be absent.") + emit(scenario="accepted-without-worker", request_id=request_id, update_id=update_id, + run_id=execution.run_id, response=response) + + +async def replacement(client): + original = json.loads(required("DURABLE_WORKFLOW_UPDATE_QUEUED")) + request_id = original["request_id"] + responses = [] + for _ in range(2): + responses.append(await client.update_workflow(workflow_id("rust"), "echo", + args=[request("python", request_id)], wait_for="completed", request_id=request_id)) + execution, history = await observe(client, "rust") + update_id, related, expected = verify_completed(history, request_id, "rust", "python") + if update_id != original["update_id"] or execution.run_id != original["run_id"]: + raise RuntimeError("Replacement changed the original accepted identity.") + if any(response.get("update_id") != update_id or response.get("update_status") != "completed" + or serializer.decode_envelope(response["result"]) != expected for response in responses): + raise RuntimeError("Duplicate request did not retain its original completion.") + emit(scenario="replacement-and-duplicate", runtime="rust", update_id=update_id, + run_id=execution.run_id, result=expected, events=related) + await matrix(client) + + +async def failure(client): + request_id = f"{required('DURABLE_WORKFLOW_UPDATE_ID')}-failure" + response = await client.update_workflow(workflow_id("rust"), "fail", + args=[request("python", request_id)], wait_for="completed", request_id=request_id) + 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 (response.get("update_status") != "failed" or response.get("update_id") != update_id + or 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: {response!r}, {related!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, events=related) + + +async def finish(client): + for runtime in RUNTIMES: + execution, _ = await observe(client, runtime) + handle = client.get_workflow_handle(execution.workflow_id) + await handle.signal("updates-finish") + result = await handle.result(timeout=60, poll_interval=.25) + execution, history = await observe(client, runtime) + if result != {"workflow_runtime": runtime, "request": workflow_id(runtime)}: + raise RuntimeError(f"Unexpected workflow result: {result!r}") + if execution.status != "completed" or sum(event["event_type"] == "WorkflowCompleted" + for event in history["events"]) != 1: + raise RuntimeError("Workflow did not complete once from its original run.") + emit(scenario="workflow-completion", runtime=runtime, run_id=execution.run_id, result=result) + + +async def main(): + parser = argparse.ArgumentParser() + parser.add_argument("phase", choices=["start", "call", "matrix", "queued", "replacement", "failure", "finish"]) + parser.add_argument("arguments", nargs="*") + args = parser.parse_args() + async with Client(required("DURABLE_WORKFLOW_SERVER_URL"), token=required("DURABLE_WORKFLOW_AUTH_TOKEN"), + namespace=required("DURABLE_WORKFLOW_NAMESPACE"), timeout=60) as client: + if args.phase == "call": + if len(args.arguments) != 3: + parser.error("call requires workflow ID, request ID and update name") + await call(client, *args.arguments) + else: + await globals()[args.phase](client) + + +if __name__ == "__main__": + asyncio.run(main()) diff --git a/polyglot/python_worker/tests/test_sdk_updates.py b/polyglot/python_worker/tests/test_sdk_updates.py new file mode 100644 index 00000000..d845cc0b --- /dev/null +++ b/polyglot/python_worker/tests/test_sdk_updates.py @@ -0,0 +1,84 @@ +"""Reject incomplete or contradictory durable update observations.""" + +import copy +import importlib.util +import pathlib +import unittest + +from durable_workflow import serializer + +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) +spec.loader.exec_module(updates) + + +class DurableUpdateObservationTest(unittest.TestCase): + def setUp(self): + self.request_id = "example-python-rust" + self.result = {"handler_runtime": "rust", "request": updates.request("python", self.request_id)} + self.history = {"events": [ + {"event_type": "UpdateAccepted", "payload": {"update_id": "original", + "arguments": serializer.envelope([self.result["request"]])}}, + {"event_type": "UpdateCompleted", "payload": {"update_id": "original", + "result": serializer.envelope(self.result)}}]} + + def verify(self, history): + return updates.verify_completed(history, self.request_id, "rust", "python") + + def test_completed_result_retains_update_identity(self): + self.assertEqual(self.verify(self.history)[0], "original") + + def test_missing_acceptance_is_rejected(self): + self.history["events"].pop(0) + with self.assertRaises(RuntimeError): + self.verify(self.history) + + def test_duplicate_acceptance_is_rejected(self): + self.history["events"].append(copy.deepcopy(self.history["events"][0])) + with self.assertRaises(RuntimeError): + self.verify(self.history) + + def test_missing_completion_is_rejected(self): + self.history["events"].pop() + with self.assertRaises(RuntimeError): + self.verify(self.history) + + def test_duplicate_completion_is_rejected(self): + self.history["events"].append(copy.deepcopy(self.history["events"][-1])) + with self.assertRaises(RuntimeError): + self.verify(self.history) + + def test_completion_for_different_update_is_rejected(self): + self.history["events"][-1]["payload"]["update_id"] = "replacement" + with self.assertRaises(RuntimeError): + self.verify(self.history) + + def test_completion_before_acceptance_is_rejected(self): + self.history["events"].reverse() + with self.assertRaises(RuntimeError): + self.verify(self.history) + + def test_failed_completion_is_rejected(self): + self.history["events"][-1]["payload"]["failure_id"] = "failed" + with self.assertRaises(RuntimeError): + self.verify(self.history) + + def test_wrong_handler_or_payload_is_rejected(self): + for replacement in ("php", {"request": "changed"}, {"count": "42"}): + with self.subTest(replacement=replacement): + history = copy.deepcopy(self.history) + result = copy.deepcopy(self.result) + if isinstance(replacement, str): + result["handler_runtime"] = replacement + elif "count" in replacement: + result["request"]["nested"].update(replacement) + else: + result.update(replacement) + history["events"][-1]["payload"]["result"] = serializer.envelope(result) + with self.assertRaises(RuntimeError): + self.verify(history) + + +if __name__ == "__main__": + unittest.main() diff --git a/polyglot/python_workflow/workflow.py b/polyglot/python_workflow/workflow.py index aed13c77..1e05aac7 100644 --- a/polyglot/python_workflow/workflow.py +++ b/polyglot/python_workflow/workflow.py @@ -85,6 +85,24 @@ def run(self, ctx, request): return {"workflow_runtime": "python", "request": request, "timer_seconds": 30} +@workflow.defn(name="polyglot.python.updates") +class PythonUpdatesWorkflow: + def __init__(self): + self.done = False + + @workflow.update("echo") + def echo(self, request): + return {"handler_runtime": "python", "request": request} + + @workflow.signal("updates-finish") + def finish(self): + self.done = True + + def run(self, ctx, request): + yield ctx.wait_condition(lambda: self.done, key="updates-finish") + return {"workflow_runtime": "python", "request": request} + + @workflow.defn(name="polyglot.python-to-php.greeter") class PythonToPhpGreeterWorkflow: def run(self, ctx, request): # type: ignore[no-untyped-def] @@ -425,6 +443,7 @@ async def main() -> int: workflows=[ PythonGreeterWorkflow, PythonTimerWorkflow, + PythonUpdatesWorkflow, PythonToPhpGreeterWorkflow, PythonToPhpTypeRoundtripWorkflow, PythonToPhpBinaryTypeRoundtripWorkflow, diff --git a/polyglot/rust_worker/src/main.rs b/polyglot/rust_worker/src/main.rs index c17db823..81e3bd27 100644 --- a/polyglot/rust_worker/src/main.rs +++ b/polyglot/rust_worker/src/main.rs @@ -4,6 +4,8 @@ use apache_avro::{from_avro_datum, to_avro_datum, Schema}; use base64::{engine::general_purpose::STANDARD as BASE64, Engine as _}; use durable_workflow::{json, ActivityOptions, AvroValue, Client, Error, Result, Value, Worker}; +mod updates; + const RUST_SAME_WORKFLOW: &str = "polyglot.rust.greeter"; const RUST_TIMER_WORKFLOW: &str = "polyglot.rust.timer"; const RUST_TIMER_DELAY_SECONDS: u64 = 30; @@ -40,6 +42,8 @@ async fn main() -> Result<()> { match mode.as_str() { "workflow" => run_workflow_worker(client).await, "activity" => run_activity_worker(client).await, + "update-client" => updates::call(client).await, + "validator-refusal" => updates::validator_refusal(client).await, other => panic!("unsupported POLYGLOT_RUST_MODE {other:?}; expected workflow or activity"), } } @@ -54,6 +58,8 @@ async fn run_workflow_worker(client: Client) -> Result<()> { .cooperative_cancellation(env_value("POLYGLOT_TIMER_COOPERATIVE", "0") == "1") .poll_timeout(Duration::from_secs(5)); + updates::register(&mut worker); + worker.register_activity("polyglot.rust.echo", |_ctx, args| async move { Ok(runtime_echo(first_argument(&args))) }); diff --git a/polyglot/rust_worker/src/updates.rs b/polyglot/rust_worker/src/updates.rs new file mode 100644 index 00000000..44cd2ead --- /dev/null +++ b/polyglot/rust_worker/src/updates.rs @@ -0,0 +1,41 @@ +use durable_workflow::{json, Client, Error, Result, Worker}; + +pub fn register(worker: &mut Worker) { + worker.register_workflow("polyglot.rust.updates", |ctx, input| async move { + let request = super::first_argument(&input); + ctx.wait_signal("updates-finish").await?; + Ok(json!({"workflow_runtime": "rust", "request": request})) + }); + worker.register_update("polyglot.rust.updates", "echo", |_ctx, input| async move { + Ok(json!({"handler_runtime": "rust", "request": super::first_argument(&input)})) + }); + worker.register_update("polyglot.rust.updates", "fail", |_ctx, _input| async move { + Err(Error::Codec("update-probe-failure".into())) + }); +} + +pub async fn call(client: Client) -> Result<()> { + let args = std::env::args().skip(1).collect::>(); + if args.len() != 3 { + return Err(Error::Codec("expected workflow ID, request ID and update name".into())); + } + let request = json!({"caller": "rust", "request_id": args[1], "value": "hello", "nested": {"enabled": true, "count": 42}}); + let result = client.update_workflow(&args[0], &args[2], json!([request]), Some(&args[1])).await?; + println!("{}", json!({"caller": "rust", "request_id": args[1], "result": result})); + Ok(()) +} + +pub async fn validator_refusal(client: Client) -> Result<()> { + let result = client.register_worker_with_command_contracts( + "rust-unsupported-validator", "polyglot-rust", vec!["polyglot.rust.updates".into()], + vec![], 1, 0, vec!["workflow_updates".into()], + json!({"polyglot.rust.updates": {"updates": ["echo"], "update_validators": ["echo"]}}), + ).await; + match result { + Err(Error::UnsupportedUpdateValidators { workflow_type }) if workflow_type == "polyglot.rust.updates" => { + println!("{}", json!({"scenario": "rust-validator-refusal", "error": "UnsupportedUpdateValidators", "workflow_type": workflow_type})); + Ok(()) + } + other => Err(Error::Codec(format!("expected unsupported-validator refusal, got {other:?}"))), + } +} diff --git a/polyglot/updates/README.md b/polyglot/updates/README.md new file mode 100644 index 00000000..d7062c79 --- /dev/null +++ b/polyglot/updates/README.md @@ -0,0 +1,42 @@ +# Published SDK workflow updates + +Run the five Rust-involving client/update-handler directions from the prepared +Sample App development container: + +```bash +scripts/playground doctor +while IFS= read -r assignment; do export "$assignment"; done \ + < <(scripts/resolve-current-artifacts.sh) +SDK_UPDATES_COMPOSE_PROJECT_NAME=sample-app-sdk-updates scripts/sdk-updates.sh +``` + +The checked-in tuple pins published Server, PHP/Python SDK packages and Rust +crate versions. The command builds the existing workers, starts an isolated +MySQL/Redis stack, and exercises PHP → Rust, Python → Rust, Rust → Rust, +Rust → PHP and Rust → Python updates. These are client-to-handler directions, +not workflow-to-activity directions. The remaining PHP/Python-only update cells +belong to the Server's existing update experiment. + +Each call checks the named handler's result through the real SDK client and +persisted `UpdateAccepted`/`UpdateCompleted` history. It then kills the Rust +worker, accepts an update while that process is absent, starts its replacement, +and requires the same update/run identities and one completion. Repeating the +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. + +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 +side effect, or coverage of every update ordering/race condition. + +The command prints the exact tuple, Server digest, UTC interval, SDK results, +run/update identities and relevant durable events. Report the scenario outcomes +on the owning issue. Its exit trap removes the task stack, volumes and worker +images on success or failure. If externally interrupted, clean up with the same +project name: + +```bash +COMPOSE_PROJECT_NAME=sample-app-sdk-updates \ + docker compose -f polyglot/docker-compose.yml down --volumes --remove-orphans +``` diff --git a/scripts/sdk-updates.sh b/scripts/sdk-updates.sh new file mode 100755 index 00000000..d6bfb19f --- /dev/null +++ b/scripts/sdk-updates.sh @@ -0,0 +1,86 @@ +#!/usr/bin/env bash +set -euo pipefail + +if [[ "${1:-}" == --help ]]; then + printf '%s\n' 'Usage: scripts/sdk-updates.sh' \ + 'Runs five Rust-involving SDK client/update-handler directions, worker replacement, duplicate requests, handler failure and validator refusal.' \ + 'Requires Docker Compose and exact assignments from scripts/resolve-current-artifacts.sh.' \ + 'SDK_UPDATES_COMPOSE_PROJECT_NAME selects an isolated project. All project resources are removed on exit.' + exit 0 +fi +[[ $# == 0 ]] || exit 2 +repo_root="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)" +export COMPOSE_PROJECT_NAME="${SDK_UPDATES_COMPOSE_PROJECT_NAME:-sample-app-sdk-updates-$(date -u +%Y%m%d%H%M%S)}" +[[ "$COMPOSE_PROJECT_NAME" =~ ^[a-z0-9][a-z0-9_-]*$ ]] || exit 2 +export DURABLE_WORKFLOW_UPDATE_ID="$COMPOSE_PROJECT_NAME" +compose=(docker compose --project-directory "$repo_root/polyglot" -f "$repo_root/polyglot/docker-compose.yml") +workers=(php-same-workflow-worker python-workflow-worker rust-workflow-worker) +if [[ -n "$(docker ps -aq --filter "label=com.docker.compose.project=$COMPOSE_PROJECT_NAME")" ]]; then + printf 'Choose a new isolated project. %s already has containers.\n' "$COMPOSE_PROJECT_NAME" >&2 + exit 2 +fi + +cleanup() { + local code=$? + if [[ "$code" != 0 ]]; then "${compose[@]}" logs --no-color --timestamps || true; fi + "${compose[@]}" down --volumes --remove-orphans || return 1 + local suffix image + for suffix in php-sdk-worker python-workflow-worker rust-workflow-worker smoke; do + image="${COMPOSE_PROJECT_NAME}-${suffix}:latest" + if docker image inspect "$image" >/dev/null 2>&1; then docker image rm --no-prune "$image" || return 1; fi + done + return "$code" +} +trap cleanup EXIT + +printf 'SDK updates start: %s\n' "$(date -u +%Y-%m-%dT%H:%M:%SZ)" +for name in DURABLE_SERVER_IMAGE DURABLE_WORKFLOW_PHP_SDK_VERSION DURABLE_WORKFLOW_PYTHON_SDK_VERSION \ + DURABLE_WORKFLOW_RUST_SDK_VERSION DURABLE_WORKFLOW_CLI_VERSION DURABLE_WORKFLOW_WORKFLOW_VERSION DURABLE_WORKFLOW_WATERLINE_VERSION; do + printf '%s=%s\n' "$name" "${!name:?resolve exact published artifacts first}" +done +"${compose[@]}" build smoke "${workers[@]}" +"${compose[@]}" pull --policy always bootstrap server +"${compose[@]}" up -d --wait --wait-timeout 180 --no-build server "${workers[@]}" +docker image inspect "$DURABLE_SERVER_IMAGE" --format '{{json .RepoDigests}}' + +observer() { + "${compose[@]}" run --rm --no-deps --user 1000:1000 \ + -e DURABLE_WORKFLOW_UPDATE_ID -e DURABLE_WORKFLOW_UPDATE_RUNS \ + -e DURABLE_WORKFLOW_UPDATE_RESULTS -e DURABLE_WORKFLOW_UPDATE_QUEUED \ + smoke python /app/scripts/sdk_updates.py "$@" +} + +client() { + local caller=$1 target=$2 request_id=$3 + case "$caller" in + php) "${compose[@]}" exec -T --user 1000:1000 php-same-workflow-worker \ + php /app/update_client.php "${COMPOSE_PROJECT_NAME}-${target}" "$request_id" echo ;; + python) observer call "${COMPOSE_PROJECT_NAME}-${target}" "$request_id" echo ;; + rust) "${compose[@]}" exec -T --user 1000:1000 -e POLYGLOT_RUST_MODE=update-client \ + rust-workflow-worker polyglot-rust-worker "${COMPOSE_PROJECT_NAME}-${target}" "$request_id" echo ;; + esac +} + +DURABLE_WORKFLOW_UPDATE_RUNS="$(observer start)" +export DURABLE_WORKFLOW_UPDATE_RUNS +printf '%s\n' "$DURABLE_WORKFLOW_UPDATE_RUNS" +export DURABLE_WORKFLOW_UPDATE_RESULTS='' +for direction in php:rust python:rust rust:rust rust:php rust:python; do + caller=${direction%:*} + target=${direction#*:} + result="$(client "$caller" "$target" "${COMPOSE_PROJECT_NAME}-${caller}-${target}")" + DURABLE_WORKFLOW_UPDATE_RESULTS+="${result}"$'\n' +done +observer matrix + +"${compose[@]}" kill --signal SIGKILL rust-workflow-worker +DURABLE_WORKFLOW_UPDATE_QUEUED="$(observer queued)" +export DURABLE_WORKFLOW_UPDATE_QUEUED +printf '%s\n' "$DURABLE_WORKFLOW_UPDATE_QUEUED" +"${compose[@]}" up -d --wait --no-build rust-workflow-worker +observer replacement +observer failure +"${compose[@]}" exec -T --user 1000:1000 -e POLYGLOT_RUST_MODE=validator-refusal \ + rust-workflow-worker polyglot-rust-worker +observer finish +printf 'SDK updates pass: %s\n' "$(date -u +%Y-%m-%dT%H:%M:%SZ)" From 5c6403086889996225dc81e6cd91552c9628fafd Mon Sep 17 00:00:00 2001 From: Durable Workflow Date: Wed, 7 Oct 2026 23:24:59 +0000 Subject: [PATCH 2/7] Keep metadata tests isolated and decode persisted Avro updates explicitly --- polyglot/php_worker/worker.php | 2 +- polyglot/python_worker/scripts/sdk_updates.py | 8 ++++---- .../python_worker/tests/test_native_binary_workers.py | 4 ++++ polyglot/python_worker/tests/test_sdk_updates.py | 11 +++++++++-- 4 files changed, 18 insertions(+), 7 deletions(-) diff --git a/polyglot/php_worker/worker.php b/polyglot/php_worker/worker.php index 25d2ac5f..d5531cd0 100644 --- a/polyglot/php_worker/worker.php +++ b/polyglot/php_worker/worker.php @@ -451,7 +451,7 @@ function configureWorkflows(Worker $worker, PayloadCodec $codec): void return ['workflow_runtime' => 'php', 'request' => $request]; }); - $worker->declareSignal('polyglot.php.updates', 'updates-finish', static function (): void {}); + $worker->declareSignal('polyglot.php.updates', 'updates-finish', static fn (): mixed => null); $worker->registerUpdate('polyglot.php.updates', 'echo', static function (QueryContext $context, array $request): array { return ['handler_runtime' => 'php', 'request' => $request]; }); diff --git a/polyglot/python_worker/scripts/sdk_updates.py b/polyglot/python_worker/scripts/sdk_updates.py index f19271ef..4e9abe0f 100644 --- a/polyglot/python_worker/scripts/sdk_updates.py +++ b/polyglot/python_worker/scripts/sdk_updates.py @@ -43,7 +43,7 @@ def update_events(history, request_id): for event in history["events"]: if event["event_type"] != "UpdateAccepted": continue - args = serializer.decode_envelope(event["payload"]["arguments"]) + args = serializer.decode_envelope(event["payload"]["arguments"], codec="avro") if args and isinstance(args[0], dict) and args[0].get("request_id") == request_id: accepted.append(event) if len(accepted) != 1: @@ -62,7 +62,7 @@ def verify_completed(history, request_id, runtime, caller): if types.index("UpdateAccepted") >= types.index("UpdateCompleted"): raise RuntimeError("Update completion precedes acceptance.") expected = {"handler_runtime": runtime, "request": request(caller, request_id)} - result = serializer.decode_envelope(completed[0]["payload"]["result"]) + result = serializer.decode_envelope(completed[0]["payload"]["result"], codec="avro") if result != expected: raise RuntimeError(f"Persisted result changed: {result!r}") return update_id, related, expected @@ -103,7 +103,7 @@ async def call(client, target, request_id, name): if response.get("update_status") != "completed": raise RuntimeError(f"Update did not complete: {response!r}") emit(caller="python", request_id=request_id, - result=serializer.decode_envelope(response["result"])) + result=serializer.decode_envelope(response["result"], codec="avro")) async def matrix(client): @@ -150,7 +150,7 @@ async def replacement(client): if update_id != original["update_id"] or execution.run_id != original["run_id"]: raise RuntimeError("Replacement changed the original accepted identity.") if any(response.get("update_id") != update_id or response.get("update_status") != "completed" - or serializer.decode_envelope(response["result"]) != expected for response in responses): + or serializer.decode_envelope(response["result"], codec="avro") != expected for response in responses): raise RuntimeError("Duplicate request did not retain its original completion.") emit(scenario="replacement-and-duplicate", runtime="rust", update_id=update_id, run_id=execution.run_id, result=expected, events=related) diff --git a/polyglot/python_worker/tests/test_native_binary_workers.py b/polyglot/python_worker/tests/test_native_binary_workers.py index 35664f4f..ba5e89b1 100644 --- a/polyglot/python_worker/tests/test_native_binary_workers.py +++ b/polyglot/python_worker/tests/test_native_binary_workers.py @@ -24,6 +24,10 @@ def signal(_name): # type: ignore[no-untyped-def] def query(_name): # type: ignore[no-untyped-def] return lambda value: value + @staticmethod + def update(_name): # type: ignore[no-untyped-def] + return lambda value: value + durable_workflow = types.ModuleType("durable_workflow") durable_workflow.Client = object diff --git a/polyglot/python_worker/tests/test_sdk_updates.py b/polyglot/python_worker/tests/test_sdk_updates.py index d845cc0b..335b1bf8 100644 --- a/polyglot/python_worker/tests/test_sdk_updates.py +++ b/polyglot/python_worker/tests/test_sdk_updates.py @@ -3,14 +3,21 @@ import copy import importlib.util import pathlib +import sys +import types import unittest +from unittest.mock import patch -from durable_workflow import serializer +# Isolate history/identity checks from codec execution. The real published +# experiment uses the installed SDK's official Avro serializer. +serializer = types.SimpleNamespace(envelope=lambda value: {"decoded": value}, + decode_envelope=lambda value, **_options: value["decoded"]) 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) -spec.loader.exec_module(updates) +with patch.dict(sys.modules, {"durable_workflow": types.SimpleNamespace(Client=object, serializer=serializer)}): + spec.loader.exec_module(updates) class DurableUpdateObservationTest(unittest.TestCase): From 5e0d839776e6383e22868b33029a28801e7f31b7 Mon Sep 17 00:00:00 2001 From: Durable Workflow Date: Wed, 7 Oct 2026 23:30:00 +0000 Subject: [PATCH 3/7] Cover all live update directions and verify complete durable history --- .github/workflows/sdk-updates.yml | 2 +- polyglot/README.md | 1 + polyglot/python_worker/scripts/sdk_updates.py | 30 ++++++++--- .../python_worker/tests/test_sdk_updates.py | 32 +++++++++-- polyglot/rust_worker/src/main.rs | 2 +- polyglot/rust_worker/src/updates.rs | 53 ++++++++++++++----- polyglot/updates/README.md | 10 ++-- scripts/sdk-updates.sh | 4 +- 8 files changed, 101 insertions(+), 33 deletions(-) diff --git a/.github/workflows/sdk-updates.yml b/.github/workflows/sdk-updates.yml index 930d5ffb..90b9a7d1 100644 --- a/.github/workflows/sdk-updates.yml +++ b/.github/workflows/sdk-updates.yml @@ -16,7 +16,7 @@ concurrency: jobs: updates: - name: SDK updates (Rust directions) + name: SDK updates (PHP/Python/Rust) runs-on: ubuntu-latest timeout-minutes: 25 env: diff --git a/polyglot/README.md b/polyglot/README.md index 51d57233..201a7142 100644 --- a/polyglot/README.md +++ b/polyglot/README.md @@ -53,6 +53,7 @@ different isolated project name. | [`sagas/`](sagas/README.md) | Five Rust-involving workflow/compensation directions across PHP, Python, and Rust | | [`schedules/`](schedules/README.md) | PHP- and Python-created automatic schedules dispatched to a Rust workflow worker | | [`timers/`](timers/README.md) | PHP/Python/Rust timer completion, worker SIGKILL and cold replay, Server restart and cooperative cancellation | +| [`updates/`](updates/README.md) | Nine PHP/Python/Rust SDK client/update-handler directions, Rust replacement and duplicate requests | | `python_workflow/` | Python-authored workflow examples | | `laravel/` | Waterline image used to inspect standalone runs | | `docker-compose.yml` | Complete service-mode topology | diff --git a/polyglot/python_worker/scripts/sdk_updates.py b/polyglot/python_worker/scripts/sdk_updates.py index 4e9abe0f..e49d78be 100644 --- a/polyglot/python_worker/scripts/sdk_updates.py +++ b/polyglot/python_worker/scripts/sdk_updates.py @@ -10,8 +10,7 @@ from durable_workflow import Client, serializer RUNTIMES = ("php", "python", "rust") -DIRECTIONS = (("php", "rust"), ("python", "rust"), ("rust", "rust"), - ("rust", "php"), ("rust", "python")) +DIRECTIONS = tuple((caller, runtime) for caller in RUNTIMES for runtime in RUNTIMES) def required(name): @@ -38,6 +37,10 @@ def request(caller, request_id): "nested": {"enabled": True, "count": 42}} +def same_result(actual, expected): + return json.dumps(actual, sort_keys=True) == json.dumps(expected, sort_keys=True) + + def update_events(history, request_id): accepted = [] for event in history["events"]: @@ -63,7 +66,7 @@ def verify_completed(history, request_id, runtime, caller): raise RuntimeError("Update completion precedes acceptance.") expected = {"handler_runtime": runtime, "request": request(caller, request_id)} result = serializer.decode_envelope(completed[0]["payload"]["result"], codec="avro") - if result != expected: + if not same_result(result, expected): raise RuntimeError(f"Persisted result changed: {result!r}") return update_id, related, expected @@ -76,7 +79,19 @@ async def observe(client, runtime): original = next(record for record in originals if record["runtime"] == runtime) if execution.run_id != original["run_id"]: raise RuntimeError("Update experiment replaced its original run.") - return execution, await client.get_history(execution.workflow_id, execution.run_id) + history = {"events": []} + token = None + seen = set() + while True: + page = await client.get_history(execution.workflow_id, execution.run_id, + page_size=100, next_page_token=token) + history["events"].extend(page["events"]) + token = page.get("next_page_token") + if not token: + return execution, history + if token in seen: + raise RuntimeError("History pagination repeated its token.") + seen.add(token) async def start(client): @@ -109,7 +124,7 @@ async def call(client, target, request_id, name): async def matrix(client): results = records("DURABLE_WORKFLOW_UPDATE_RESULTS") if len(results) != len(DIRECTIONS): - raise RuntimeError("Not all five client/handler directions executed.") + raise RuntimeError("Not all nine client/handler directions executed.") for caller, runtime in DIRECTIONS: request_id = f"{required('DURABLE_WORKFLOW_UPDATE_ID')}-{caller}-{runtime}" matches = [result for result in results if result.get("request_id") == request_id] @@ -117,7 +132,7 @@ async def matrix(client): raise RuntimeError("Missing or duplicate SDK client observation.") execution, history = await observe(client, runtime) update_id, related, expected = verify_completed(history, request_id, runtime, caller) - if matches[0].get("result") != expected: + if not same_result(matches[0].get("result"), expected): raise RuntimeError(f"SDK client result differs from persisted result: {matches[0]!r}") emit(scenario="client-handler", caller=caller, runtime=runtime, run_id=execution.run_id, update_id=update_id, result=expected, events=related) @@ -150,7 +165,8 @@ async def replacement(client): if update_id != original["update_id"] or execution.run_id != original["run_id"]: raise RuntimeError("Replacement changed the original accepted identity.") if any(response.get("update_id") != update_id or response.get("update_status") != "completed" - or serializer.decode_envelope(response["result"], codec="avro") != expected for response in responses): + or not same_result(serializer.decode_envelope(response["result"], codec="avro"), expected) + for response in responses): raise RuntimeError("Duplicate request did not retain its original completion.") emit(scenario="replacement-and-duplicate", runtime="rust", update_id=update_id, run_id=execution.run_id, result=expected, events=related) diff --git a/polyglot/python_worker/tests/test_sdk_updates.py b/polyglot/python_worker/tests/test_sdk_updates.py index 335b1bf8..135e1329 100644 --- a/polyglot/python_worker/tests/test_sdk_updates.py +++ b/polyglot/python_worker/tests/test_sdk_updates.py @@ -2,11 +2,13 @@ import copy import importlib.util +import json +import os import pathlib import sys import types import unittest -from unittest.mock import patch +from unittest.mock import AsyncMock, patch # Isolate history/identity checks from codec execution. The real published # experiment uses the installed SDK's official Avro serializer. @@ -72,13 +74,13 @@ def test_failed_completion_is_rejected(self): self.verify(self.history) def test_wrong_handler_or_payload_is_rejected(self): - for replacement in ("php", {"request": "changed"}, {"count": "42"}): + for replacement in ("php", {"request": "changed"}, {"count": "42"}, {"count": 42.0}, {"enabled": 1}): with self.subTest(replacement=replacement): history = copy.deepcopy(self.history) result = copy.deepcopy(self.result) if isinstance(replacement, str): result["handler_runtime"] = replacement - elif "count" in replacement: + elif "count" in replacement or "enabled" in replacement: result["request"]["nested"].update(replacement) else: result.update(replacement) @@ -87,5 +89,29 @@ def test_wrong_handler_or_payload_is_rejected(self): self.verify(history) +class HistoryPaginationTest(unittest.IsolatedAsyncioTestCase): + async def test_observer_reads_later_pages(self): + client = types.SimpleNamespace( + describe_workflow=AsyncMock(return_value=types.SimpleNamespace(workflow_id="example-rust", run_id="original")), + get_history=AsyncMock(side_effect=[ + {"events": [{"event_type": "UpdateAccepted"}], "next_page_token": "page-2"}, + {"events": [{"event_type": "UpdateCompleted"}], "next_page_token": None}, + ]), + ) + with patch.dict(os.environ, {"DURABLE_WORKFLOW_UPDATE_ID": "example", + "DURABLE_WORKFLOW_UPDATE_RUNS": json.dumps({"runtime": "rust", "run_id": "original"})}): + _, history = await updates.observe(client, "rust") + self.assertEqual([event["event_type"] for event in history["events"]], ["UpdateAccepted", "UpdateCompleted"]) + self.assertEqual(client.get_history.await_args.kwargs["next_page_token"], "page-2") + + async def test_replacement_run_is_rejected(self): + client = types.SimpleNamespace(describe_workflow=AsyncMock( + return_value=types.SimpleNamespace(workflow_id="example-rust", run_id="replacement"))) + with patch.dict(os.environ, {"DURABLE_WORKFLOW_UPDATE_ID": "example", + "DURABLE_WORKFLOW_UPDATE_RUNS": json.dumps({"runtime": "rust", "run_id": "original"})}): + with self.assertRaisesRegex(RuntimeError, "original run"): + await updates.observe(client, "rust") + + if __name__ == "__main__": unittest.main() diff --git a/polyglot/rust_worker/src/main.rs b/polyglot/rust_worker/src/main.rs index 81e3bd27..a659f6ca 100644 --- a/polyglot/rust_worker/src/main.rs +++ b/polyglot/rust_worker/src/main.rs @@ -44,7 +44,7 @@ async fn main() -> Result<()> { "activity" => run_activity_worker(client).await, "update-client" => updates::call(client).await, "validator-refusal" => updates::validator_refusal(client).await, - other => panic!("unsupported POLYGLOT_RUST_MODE {other:?}; expected workflow or activity"), + other => panic!("unsupported POLYGLOT_RUST_MODE {other:?}; expected workflow, activity, update-client or validator-refusal"), } } diff --git a/polyglot/rust_worker/src/updates.rs b/polyglot/rust_worker/src/updates.rs index 44cd2ead..4065ddd8 100644 --- a/polyglot/rust_worker/src/updates.rs +++ b/polyglot/rust_worker/src/updates.rs @@ -6,9 +6,13 @@ pub fn register(worker: &mut Worker) { ctx.wait_signal("updates-finish").await?; Ok(json!({"workflow_runtime": "rust", "request": request})) }); - worker.register_update("polyglot.rust.updates", "echo", |_ctx, input| async move { - Ok(json!({"handler_runtime": "rust", "request": super::first_argument(&input)})) - }); + worker.register_update( + "polyglot.rust.updates", + "echo", + |_ctx, input| async move { + Ok(json!({"handler_runtime": "rust", "request": super::first_argument(&input)})) + }, + ); worker.register_update("polyglot.rust.updates", "fail", |_ctx, _input| async move { Err(Error::Codec("update-probe-failure".into())) }); @@ -17,25 +21,46 @@ pub fn register(worker: &mut Worker) { pub async fn call(client: Client) -> Result<()> { let args = std::env::args().skip(1).collect::>(); if args.len() != 3 { - return Err(Error::Codec("expected workflow ID, request ID and update name".into())); + return Err(Error::Codec( + "expected workflow ID, request ID and update name".into(), + )); } let request = json!({"caller": "rust", "request_id": args[1], "value": "hello", "nested": {"enabled": true, "count": 42}}); - let result = client.update_workflow(&args[0], &args[2], json!([request]), Some(&args[1])).await?; - println!("{}", json!({"caller": "rust", "request_id": args[1], "result": result})); + let result = client + .update_workflow(&args[0], &args[2], json!([request]), Some(&args[1])) + .await?; + println!( + "{}", + json!({"caller": "rust", "request_id": args[1], "result": result}) + ); Ok(()) } pub async fn validator_refusal(client: Client) -> Result<()> { - let result = client.register_worker_with_command_contracts( - "rust-unsupported-validator", "polyglot-rust", vec!["polyglot.rust.updates".into()], - vec![], 1, 0, vec!["workflow_updates".into()], - json!({"polyglot.rust.updates": {"updates": ["echo"], "update_validators": ["echo"]}}), - ).await; + let result = client + .register_worker_with_command_contracts( + "rust-unsupported-validator", + "polyglot-rust", + vec!["polyglot.rust.updates".into()], + vec![], + 1, + 0, + vec!["workflow_updates".into()], + json!({"polyglot.rust.updates": {"updates": ["echo"], "update_validators": ["echo"]}}), + ) + .await; match result { - Err(Error::UnsupportedUpdateValidators { workflow_type }) if workflow_type == "polyglot.rust.updates" => { - println!("{}", json!({"scenario": "rust-validator-refusal", "error": "UnsupportedUpdateValidators", "workflow_type": workflow_type})); + Err(Error::UnsupportedUpdateValidators { workflow_type }) + if workflow_type == "polyglot.rust.updates" => + { + println!( + "{}", + json!({"scenario": "rust-validator-refusal", "error": "UnsupportedUpdateValidators", "workflow_type": workflow_type}) + ); Ok(()) } - other => Err(Error::Codec(format!("expected unsupported-validator refusal, got {other:?}"))), + other => Err(Error::Codec(format!( + "expected unsupported-validator refusal, got {other:?}" + ))), } } diff --git a/polyglot/updates/README.md b/polyglot/updates/README.md index d7062c79..bc41a31c 100644 --- a/polyglot/updates/README.md +++ b/polyglot/updates/README.md @@ -1,6 +1,6 @@ # Published SDK workflow updates -Run the five Rust-involving client/update-handler directions from the prepared +Run all nine PHP/Python/Rust client/update-handler directions from the prepared Sample App development container: ```bash @@ -12,10 +12,10 @@ SDK_UPDATES_COMPOSE_PROJECT_NAME=sample-app-sdk-updates scripts/sdk-updates.sh The checked-in tuple pins published Server, PHP/Python SDK packages and Rust crate versions. The command builds the existing workers, starts an isolated -MySQL/Redis stack, and exercises PHP → Rust, Python → Rust, Rust → Rust, -Rust → PHP and Rust → Python updates. These are client-to-handler directions, -not workflow-to-activity directions. The remaining PHP/Python-only update cells -belong to the Server's existing update experiment. +MySQL/Redis stack, and exercises each of the three client languages against +each of the three handler languages. These are client-to-handler directions. +The Server's existing update experiment separately covers its embedded probe, +PHP process boundary, Python SDK surface fixtures, validators and diagnostics. Each call checks the named handler's result through the real SDK client and persisted `UpdateAccepted`/`UpdateCompleted` history. It then kills the Rust diff --git a/scripts/sdk-updates.sh b/scripts/sdk-updates.sh index d6bfb19f..3922c36e 100755 --- a/scripts/sdk-updates.sh +++ b/scripts/sdk-updates.sh @@ -3,7 +3,7 @@ set -euo pipefail if [[ "${1:-}" == --help ]]; then printf '%s\n' 'Usage: scripts/sdk-updates.sh' \ - 'Runs five Rust-involving SDK client/update-handler directions, worker replacement, duplicate requests, handler failure and validator refusal.' \ + 'Runs nine PHP/Python/Rust client/update-handler directions, Rust worker replacement, duplicate requests, handler failure and validator refusal.' \ 'Requires Docker Compose and exact assignments from scripts/resolve-current-artifacts.sh.' \ 'SDK_UPDATES_COMPOSE_PROJECT_NAME selects an isolated project. All project resources are removed on exit.' exit 0 @@ -65,7 +65,7 @@ DURABLE_WORKFLOW_UPDATE_RUNS="$(observer start)" export DURABLE_WORKFLOW_UPDATE_RUNS printf '%s\n' "$DURABLE_WORKFLOW_UPDATE_RUNS" export DURABLE_WORKFLOW_UPDATE_RESULTS='' -for direction in php:rust python:rust rust:rust rust:php rust:python; do +for direction in php:php php:python php:rust python:php python:python python:rust rust:php rust:python rust:rust; do caller=${direction%:*} target=${direction#*:} result="$(client "$caller" "$target" "${COMPOSE_PROJECT_NAME}-${caller}-${target}")" From 3d26a084865f1463ca4b74cd5fb81ada5d5c1bb2 Mon Sep 17 00:00:00 2001 From: Durable Workflow Date: Wed, 7 Oct 2026 23:35:05 +0000 Subject: [PATCH 4/7] Expose SDK client failure output in the published update check --- scripts/sdk-updates.sh | 7 +++++-- 1 file changed, 5 insertions(+), 2 deletions(-) diff --git a/scripts/sdk-updates.sh b/scripts/sdk-updates.sh index 3922c36e..2485a069 100755 --- a/scripts/sdk-updates.sh +++ b/scripts/sdk-updates.sh @@ -54,7 +54,7 @@ client() { local caller=$1 target=$2 request_id=$3 case "$caller" in php) "${compose[@]}" exec -T --user 1000:1000 php-same-workflow-worker \ - php /app/update_client.php "${COMPOSE_PROJECT_NAME}-${target}" "$request_id" echo ;; + php -d display_errors=stderr /app/update_client.php "${COMPOSE_PROJECT_NAME}-${target}" "$request_id" echo ;; python) observer call "${COMPOSE_PROJECT_NAME}-${target}" "$request_id" echo ;; rust) "${compose[@]}" exec -T --user 1000:1000 -e POLYGLOT_RUST_MODE=update-client \ rust-workflow-worker polyglot-rust-worker "${COMPOSE_PROJECT_NAME}-${target}" "$request_id" echo ;; @@ -68,7 +68,10 @@ export DURABLE_WORKFLOW_UPDATE_RESULTS='' for direction in php:php php:python php:rust python:php python:python python:rust rust:php rust:python rust:rust; do caller=${direction%:*} target=${direction#*:} - result="$(client "$caller" "$target" "${COMPOSE_PROJECT_NAME}-${caller}-${target}")" + if ! result="$(client "$caller" "$target" "${COMPOSE_PROJECT_NAME}-${caller}-${target}")"; then + printf 'SDK client %s failed:\n%s\n' "$direction" "$result" >&2 + exit 1 + fi DURABLE_WORKFLOW_UPDATE_RESULTS+="${result}"$'\n' done observer matrix From fab9e2ad6dff1a09ceec3fbdc475636527600fa2 Mon Sep 17 00:00:00 2001 From: Durable Workflow Date: Wed, 7 Oct 2026 23:38:48 +0000 Subject: [PATCH 5/7] Report typed update refusal reason and response details --- polyglot/php_worker/update_client.php | 12 +++++++++++- 1 file changed, 11 insertions(+), 1 deletion(-) diff --git a/polyglot/php_worker/update_client.php b/polyglot/php_worker/update_client.php index 2ff07d73..d7fc099a 100644 --- a/polyglot/php_worker/update_client.php +++ b/polyglot/php_worker/update_client.php @@ -3,6 +3,7 @@ declare(strict_types=1); use DurableWorkflow\Client; +use DurableWorkflow\Exception\ServerException; require __DIR__.'/vendor/autoload.php'; @@ -15,5 +16,14 @@ namespace: (string) getenv('DURABLE_WORKFLOW_NAMESPACE'), ); $request = ['caller' => 'php', 'request_id' => $argv[2], 'value' => 'hello', 'nested' => ['enabled' => true, 'count' => 42]]; -$result = $client->updateWorkflow($argv[1], $argv[3], [$request], requestId: $argv[2]); +try { + $result = $client->updateWorkflow($argv[1], $argv[3], [$request], requestId: $argv[2]); +} catch (ServerException $exception) { + fwrite(STDERR, json_encode([ + 'caller' => 'php', 'workflow_id' => $argv[1], 'request_id' => $argv[2], + 'error' => get_class($exception), 'status' => $exception->status, + 'reason' => $exception->reason, 'details' => $exception->details, + ], JSON_THROW_ON_ERROR).PHP_EOL); + exit(1); +} fwrite(STDOUT, json_encode(['caller' => 'php', 'request_id' => $argv[2], 'result' => $result], JSON_THROW_ON_ERROR).PHP_EOL); From 82aba89b159acd87633368cd09016f91add87752 Mon Sep 17 00:00:00 2001 From: Durable Workflow Date: Wed, 7 Oct 2026 23:55:03 +0000 Subject: [PATCH 6/7] Qualify update handlers with published Rust SDK 3.2.2 --- playground/templates/rust/Cargo.lock | 4 +- polyglot/python_worker/scripts/sdk_updates.py | 23 ++++++----- .../python_worker/tests/test_sdk_updates.py | 38 ++++++++++++++++++- polyglot/qualified-artifact-tuple.json | 2 +- polyglot/rust_worker/Cargo.lock | 4 +- polyglot/rust_worker/Cargo.toml | 2 +- polyglot/rust_worker/src/updates.rs | 10 ++--- polyglot/updates/README.md | 3 ++ 8 files changed, 63 insertions(+), 23 deletions(-) diff --git a/playground/templates/rust/Cargo.lock b/playground/templates/rust/Cargo.lock index 3f89d396..0316bc66 100644 --- a/playground/templates/rust/Cargo.lock +++ b/playground/templates/rust/Cargo.lock @@ -268,9 +268,9 @@ dependencies = [ [[package]] name = "durable-workflow" -version = "3.2.1" +version = "3.2.2" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "9d0482e9e070aefa6844a16132ca35eea8dcebc1524859f879961e57fc0985ae" +checksum = "305d4db324c3bb35177473c8d3aa3a16674e710fdeccde34d43d1da123286a95" dependencies = [ "apache-avro", "base64", diff --git a/polyglot/python_worker/scripts/sdk_updates.py b/polyglot/python_worker/scripts/sdk_updates.py index e49d78be..026b395e 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, serializer +from durable_workflow import Client, InvalidArgument, serializer RUNTIMES = ("php", "python", "rust") DIRECTIONS = tuple((caller, runtime) for caller in RUNTIMES for runtime in RUNTIMES) @@ -118,7 +118,7 @@ async def call(client, target, request_id, name): if response.get("update_status") != "completed": raise RuntimeError(f"Update did not complete: {response!r}") emit(caller="python", request_id=request_id, - result=serializer.decode_envelope(response["result"], codec="avro")) + result=serializer.decode_envelope(response["result_envelope"], codec="avro")) async def matrix(client): @@ -165,7 +165,7 @@ async def replacement(client): if update_id != original["update_id"] or execution.run_id != original["run_id"]: raise RuntimeError("Replacement changed the original accepted identity.") if any(response.get("update_id") != update_id or response.get("update_status") != "completed" - or not same_result(serializer.decode_envelope(response["result"], codec="avro"), expected) + or not same_result(serializer.decode_envelope(response["result_envelope"], codec="avro"), expected) for response in responses): raise RuntimeError("Duplicate request did not retain its original completion.") emit(scenario="replacement-and-duplicate", runtime="rust", update_id=update_id, @@ -175,19 +175,24 @@ async def replacement(client): async def failure(client): request_id = f"{required('DURABLE_WORKFLOW_UPDATE_ID')}-failure" - response = await client.update_workflow(workflow_id("rust"), "fail", - args=[request("python", request_id)], wait_for="completed", request_id=request_id) + 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.") 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 (response.get("update_status") != "failed" or response.get("update_id") != update_id - or len(failed) != 1 or "update-probe-failure" not in json.dumps(failed[0]["payload"]) + 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: {response!r}, {related!r}") + raise RuntimeError(f"Missing durable handler failure: {sdk_error!r}, {related!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, events=related) + emit(scenario="handler-failure", runtime="rust", update_id=update_id, + sdk_error=sdk_error, 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 135e1329..64933c00 100644 --- a/polyglot/python_worker/tests/test_sdk_updates.py +++ b/polyglot/python_worker/tests/test_sdk_updates.py @@ -18,7 +18,12 @@ 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) -with patch.dict(sys.modules, {"durable_workflow": types.SimpleNamespace(Client=object, serializer=serializer)}): +class InvalidArgument(Exception): + pass + + +with patch.dict(sys.modules, {"durable_workflow": types.SimpleNamespace( + Client=object, InvalidArgument=InvalidArgument, serializer=serializer)}): spec.loader.exec_module(updates) @@ -89,6 +94,37 @@ def test_wrong_handler_or_payload_is_rejected(self): self.verify(history) +class ClientResultTest(unittest.IsolatedAsyncioTestCase): + async def test_client_decodes_the_envelope_field(self): + result = {"handler_runtime": "rust", "request": updates.request("python", "request")} + client = types.SimpleNamespace(update_workflow=AsyncMock(return_value={ + "update_status": "completed", "result": result, + "result_envelope": serializer.envelope(result), + })) + with patch.object(updates, "emit") as emit: + await updates.call(client, "workflow", "request", "echo") + self.assertEqual(emit.call_args.kwargs["result"], result) + + 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") + 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"}}, + ]} + 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" + with self.assertRaisesRegex(RuntimeError, "durable handler failure"): + await updates.failure(client) + + class HistoryPaginationTest(unittest.IsolatedAsyncioTestCase): async def test_observer_reads_later_pages(self): client = types.SimpleNamespace( diff --git a/polyglot/qualified-artifact-tuple.json b/polyglot/qualified-artifact-tuple.json index 2c0a6e97..df65c4f5 100644 --- a/polyglot/qualified-artifact-tuple.json +++ b/polyglot/qualified-artifact-tuple.json @@ -6,7 +6,7 @@ "cli": "2.1.3", "sdk-php": "2.2.3", "sdk-python": "2.4.2", - "sdk-rust": "3.2.1", + "sdk-rust": "3.2.2", "server": "2.5.9", "waterline": "2.3.1", "workflow": "2.5.3" diff --git a/polyglot/rust_worker/Cargo.lock b/polyglot/rust_worker/Cargo.lock index eb9fbe93..27cfb51e 100644 --- a/polyglot/rust_worker/Cargo.lock +++ b/polyglot/rust_worker/Cargo.lock @@ -244,9 +244,9 @@ dependencies = [ [[package]] name = "durable-workflow" -version = "3.2.1" +version = "3.2.2" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "9d0482e9e070aefa6844a16132ca35eea8dcebc1524859f879961e57fc0985ae" +checksum = "305d4db324c3bb35177473c8d3aa3a16674e710fdeccde34d43d1da123286a95" dependencies = [ "apache-avro", "base64", diff --git a/polyglot/rust_worker/Cargo.toml b/polyglot/rust_worker/Cargo.toml index 570620d2..bf03a634 100644 --- a/polyglot/rust_worker/Cargo.toml +++ b/polyglot/rust_worker/Cargo.toml @@ -7,7 +7,7 @@ publish = false [dependencies] apache-avro = "=0.21.0" base64 = "=0.22.1" -durable-workflow = "3.2.1" +durable-workflow = "3.2.2" tokio = { version = "1.47.1", features = ["macros", "rt-multi-thread", "signal"] } [dev-dependencies] diff --git a/polyglot/rust_worker/src/updates.rs b/polyglot/rust_worker/src/updates.rs index 4065ddd8..f85b1dd9 100644 --- a/polyglot/rust_worker/src/updates.rs +++ b/polyglot/rust_worker/src/updates.rs @@ -6,13 +6,9 @@ pub fn register(worker: &mut Worker) { ctx.wait_signal("updates-finish").await?; Ok(json!({"workflow_runtime": "rust", "request": request})) }); - worker.register_update( - "polyglot.rust.updates", - "echo", - |_ctx, input| async move { - Ok(json!({"handler_runtime": "rust", "request": super::first_argument(&input)})) - }, - ); + worker.register_update("polyglot.rust.updates", "echo", |_ctx, input| async move { + Ok(json!({"handler_runtime": "rust", "request": super::first_argument(&input)})) + }); worker.register_update("polyglot.rust.updates", "fail", |_ctx, _input| async move { Err(Error::Codec("update-probe-failure".into())) }); diff --git a/polyglot/updates/README.md b/polyglot/updates/README.md index bc41a31c..e48ec872 100644 --- a/polyglot/updates/README.md +++ b/polyglot/updates/README.md @@ -17,6 +17,9 @@ each of the three handler languages. These are client-to-handler directions. The Server's existing update experiment separately covers its embedded probe, PHP process boundary, Python SDK surface fixtures, validators and diagnostics. +Rust update workers require SDK 3.2.2 or later so their registration includes +the argument contracts Server records when starting a run. + Each call checks the named handler's result through the real SDK client and persisted `UpdateAccepted`/`UpdateCompleted` history. It then kills the Rust worker, accepts an update while that process is absent, starts its replacement, From 3f19dd6e480b41a3cba5eac96df657c94162d63d Mon Sep 17 00:00:00 2001 From: Durable Workflow Date: Thu, 8 Oct 2026 00:07:47 +0000 Subject: [PATCH 7/7] Declare Rust completion signal using published SDK 3.3.0 --- playground/templates/rust/Cargo.lock | 4 ++-- polyglot/qualified-artifact-tuple.json | 2 +- polyglot/rust_worker/Cargo.lock | 4 ++-- polyglot/rust_worker/Cargo.toml | 2 +- polyglot/rust_worker/src/updates.rs | 3 +++ polyglot/updates/README.md | 5 +++-- 6 files changed, 12 insertions(+), 8 deletions(-) diff --git a/playground/templates/rust/Cargo.lock b/playground/templates/rust/Cargo.lock index 0316bc66..d6367ffd 100644 --- a/playground/templates/rust/Cargo.lock +++ b/playground/templates/rust/Cargo.lock @@ -268,9 +268,9 @@ dependencies = [ [[package]] name = "durable-workflow" -version = "3.2.2" +version = "3.3.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "305d4db324c3bb35177473c8d3aa3a16674e710fdeccde34d43d1da123286a95" +checksum = "c5f5cc23a542ec5249798a396ac3e5745c0819ec94f682372d453e338d88278c" dependencies = [ "apache-avro", "base64", diff --git a/polyglot/qualified-artifact-tuple.json b/polyglot/qualified-artifact-tuple.json index df65c4f5..722a6207 100644 --- a/polyglot/qualified-artifact-tuple.json +++ b/polyglot/qualified-artifact-tuple.json @@ -6,7 +6,7 @@ "cli": "2.1.3", "sdk-php": "2.2.3", "sdk-python": "2.4.2", - "sdk-rust": "3.2.2", + "sdk-rust": "3.3.0", "server": "2.5.9", "waterline": "2.3.1", "workflow": "2.5.3" diff --git a/polyglot/rust_worker/Cargo.lock b/polyglot/rust_worker/Cargo.lock index 27cfb51e..9bf7bcd5 100644 --- a/polyglot/rust_worker/Cargo.lock +++ b/polyglot/rust_worker/Cargo.lock @@ -244,9 +244,9 @@ dependencies = [ [[package]] name = "durable-workflow" -version = "3.2.2" +version = "3.3.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "305d4db324c3bb35177473c8d3aa3a16674e710fdeccde34d43d1da123286a95" +checksum = "c5f5cc23a542ec5249798a396ac3e5745c0819ec94f682372d453e338d88278c" dependencies = [ "apache-avro", "base64", diff --git a/polyglot/rust_worker/Cargo.toml b/polyglot/rust_worker/Cargo.toml index bf03a634..a99a433c 100644 --- a/polyglot/rust_worker/Cargo.toml +++ b/polyglot/rust_worker/Cargo.toml @@ -7,7 +7,7 @@ publish = false [dependencies] apache-avro = "=0.21.0" base64 = "=0.22.1" -durable-workflow = "3.2.2" +durable-workflow = "3.3.0" tokio = { version = "1.47.1", features = ["macros", "rt-multi-thread", "signal"] } [dev-dependencies] diff --git a/polyglot/rust_worker/src/updates.rs b/polyglot/rust_worker/src/updates.rs index f85b1dd9..9df71ff8 100644 --- a/polyglot/rust_worker/src/updates.rs +++ b/polyglot/rust_worker/src/updates.rs @@ -6,6 +6,9 @@ pub fn register(worker: &mut Worker) { ctx.wait_signal("updates-finish").await?; Ok(json!({"workflow_runtime": "rust", "request": request})) }); + worker + .declare_workflow_signals("polyglot.rust.updates", &["updates-finish"]) + .expect("declare the update workflow's completion signal"); worker.register_update("polyglot.rust.updates", "echo", |_ctx, input| async move { Ok(json!({"handler_runtime": "rust", "request": super::first_argument(&input)})) }); diff --git a/polyglot/updates/README.md b/polyglot/updates/README.md index e48ec872..2e5e24c5 100644 --- a/polyglot/updates/README.md +++ b/polyglot/updates/README.md @@ -17,8 +17,9 @@ each of the three handler languages. These are client-to-handler directions. The Server's existing update experiment separately covers its embedded probe, PHP process boundary, Python SDK surface fixtures, validators and diagnostics. -Rust update workers require SDK 3.2.2 or later so their registration includes -the argument contracts Server records when starting a run. +Rust update workers require SDK 3.3.0 or later. Their registration includes the +handler argument contracts and declared completion signal Server records when +starting a run. Each call checks the named handler's result through the real SDK client and persisted `UpdateAccepted`/`UpdateCompleted` history. It then kills the Rust