From 71962a2d515736b914886abe88d57d109d80b07c Mon Sep 17 00:00:00 2001 From: Durable Workflow Date: Wed, 7 Oct 2026 22:24:23 +0000 Subject: [PATCH 1/5] Exercise published PHP Python and Rust timers across restarts and cancellation --- .github/workflows/sdk-timers.yml | 25 +++ polyglot/README.md | 2 +- polyglot/docker-compose.yml | 5 +- polyglot/php_worker/worker.php | 7 + polyglot/python_worker/scripts/sdk_timers.py | 173 +++++++++++++++++++ polyglot/python_workflow/workflow.py | 8 + polyglot/qualified-artifact-tuple.json | 2 +- polyglot/timers/README.md | 50 +++++- scripts/resolve-current-artifacts.sh | 18 +- scripts/sdk-timers.sh | 64 +++++++ tests/Unit/PolyglotComposeContractTest.php | 10 +- 11 files changed, 349 insertions(+), 15 deletions(-) create mode 100644 .github/workflows/sdk-timers.yml create mode 100644 polyglot/python_worker/scripts/sdk_timers.py create mode 100755 scripts/sdk-timers.sh diff --git a/.github/workflows/sdk-timers.yml b/.github/workflows/sdk-timers.yml new file mode 100644 index 00000000..b0ea6b1b --- /dev/null +++ b/.github/workflows/sdk-timers.yml @@ -0,0 +1,25 @@ +name: published SDK timers + +on: + push: + branches: [main] + pull_request: + branches: [main] + workflow_dispatch: + +permissions: + contents: read + +jobs: + timers: + name: SDK timers (PHP/Python/Rust) + runs-on: ubuntu-latest + timeout-minutes: 30 + env: + SDK_TIMERS_COMPOSE_PROJECT_NAME: sample-app-sdk-timers-${{ 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 timer experiment + run: scripts/sdk-timers.sh diff --git a/polyglot/README.md b/polyglot/README.md index 3ea97031..51d57233 100644 --- a/polyglot/README.md +++ b/polyglot/README.md @@ -52,7 +52,7 @@ different isolated project name. | [`child-workflows/`](child-workflows/README.md) | Runnable PHP/Python/Rust parent-child matrix using the local playground runtime | | [`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) | Rust durable timer fired while its workflow worker is stopped, then cold replayed | +| [`timers/`](timers/README.md) | PHP/Python/Rust timer completion, worker SIGKILL and cold replay, Server restart and cooperative cancellation | | `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/docker-compose.yml b/polyglot/docker-compose.yml index e2fb8e3c..384a7570 100644 --- a/polyglot/docker-compose.yml +++ b/polyglot/docker-compose.yml @@ -158,9 +158,10 @@ services: server: condition: service_healthy - # Opt-in queue consumer for durable workflow timers in the Rust timer example. + # Opt-in queue consumer for the SDK timer examples. timer-queue: - profiles: ["timer-rust"] + profiles: ["timer-rust", "timers"] + init: true image: *durable-server-image command: ["php", "artisan", "queue:work", "--sleep=1", "--tries=3", "--max-time=3600"] environment: diff --git a/polyglot/php_worker/worker.php b/polyglot/php_worker/worker.php index a99960c8..0811b28d 100644 --- a/polyglot/php_worker/worker.php +++ b/polyglot/php_worker/worker.php @@ -22,6 +22,7 @@ const WORKFLOW_TYPES = [ 'polyglot.php.greeter', + 'polyglot.php.timer', 'polyglot.PolyglotWorkflow', 'polyglot.php-to-python.greeter', 'polyglot.php-to-python.type-roundtrip', @@ -443,6 +444,12 @@ function signalQueryWorkflow(): Closure function configureWorkflows(Worker $worker, PayloadCodec $codec): void { + $worker->registerWorkflow('polyglot.php.timer', static function (WorkflowContext $context, string $request): array { + $context->sleep(30); + + return ['workflow_runtime' => 'php', 'request' => $request, 'timer_seconds' => 30]; + }); + $workflowQueue = getenv('POLYGLOT_WORKFLOW_TASK_QUEUE') ?: 'polyglot-workflow'; $pythonQueue = getenv('POLYGLOT_PHP2PY_TASK_QUEUE') ?: 'polyglot-php-to-python'; $rustQueue = getenv('POLYGLOT_TO_RUST_TASK_QUEUE') ?: 'polyglot-to-rust'; diff --git a/polyglot/python_worker/scripts/sdk_timers.py b/polyglot/python_worker/scripts/sdk_timers.py new file mode 100644 index 00000000..b20940e8 --- /dev/null +++ b/polyglot/python_worker/scripts/sdk_timers.py @@ -0,0 +1,173 @@ +"""Exercise published PHP/Python/Rust workflow timers through the public API.""" + +from __future__ import annotations + +import argparse +import asyncio +import json +import os +from datetime import datetime, timezone + +from durable_workflow import Client + + +RUNTIMES = ("php", "python", "rust") +TIMER_SECONDS = 30 + + +def required(name: str) -> str: + value = os.environ.get(name, "").strip() + if not value: + raise RuntimeError(f"Set {name} before running the SDK timer experiment.") + return value + + +def events(history: dict, event_type: str) -> list[dict]: + return [event for event in history["events"] if event["event_type"] == event_type] + + +def timestamp(value: str) -> datetime: + return datetime.fromisoformat(value.replace("Z", "+00:00")) + + +def emit(**record) -> None: + print(json.dumps(record, sort_keys=True), flush=True) + + +async def observe(client: Client, runtime: str, scenario: str): + workflow_id = f"{required('DURABLE_WORKFLOW_TIMER_ID')}-{scenario}-{runtime}" + handle = client.get_workflow_handle(workflow_id, workflow_type=f"polyglot.{runtime}.timer") + execution = await handle.describe() + if not execution.run_id: + raise RuntimeError(f"No run ID for {workflow_id}.") + history = await client.get_history(workflow_id, execution.run_id) + return handle, execution, history + + +async def pending(client: Client, runtime: str, scenario: str) -> None: + workflow_id = f"{required('DURABLE_WORKFLOW_TIMER_ID')}-{scenario}-{runtime}" + await client.start_workflow( + workflow_type=f"polyglot.{runtime}.timer", + workflow_id=workflow_id, + task_queue=f"polyglot-{runtime}", + input=[workflow_id], + ) + deadline = asyncio.get_running_loop().time() + 20 + while True: + handle, execution, history = await observe(client, runtime, scenario) + scheduled = events(history, "TimerScheduled") + if scheduled: + break + if execution.status in ("failed", "cancelled", "terminated", "completed"): + raise RuntimeError(f"{runtime} closed before scheduling its timer: {history!r}") + if asyncio.get_running_loop().time() >= deadline: + raise TimeoutError(f"{runtime} did not schedule its durable timer.") + await asyncio.sleep(0.25) + if len(scheduled) != 1 or events(history, "TimerFired") or execution.status != "waiting": + raise RuntimeError(f"Expected one pending timer and public waiting status: {history!r}") + timer = scheduled[0]["payload"] + if timer.get("delay_seconds") != TIMER_SECONDS or not timer.get("timer_id") or not timer.get("fire_at"): + raise RuntimeError(f"Invalid timer: {timer!r}") + if timestamp(timer["fire_at"]) <= datetime.now(timezone.utc): + raise RuntimeError("The pending timer's deadline has already passed.") + emit(phase="pending", runtime=runtime, scenario=scenario, workflow_id=workflow_id, + run_id=execution.run_id, status=execution.status, timer=timer) + if scenario == "cancellation": + first = await handle.request_cancellation(reason="timer conformance", cleanup_timeout_seconds=20) + duplicate = await handle.request_cancellation(reason="duplicate request", cleanup_timeout_seconds=60) + original, repeated = first["cancellation_request"], duplicate["cancellation_request"] + if (first["duplicate"] or not duplicate["duplicate"] + or original["request_id"] != repeated["request_id"] + or original["cleanup_deadline_at"] != repeated["cleanup_deadline_at"]): + raise RuntimeError(f"Cancellation identity or original deadline changed: {first!r}, {duplicate!r}") + emit(phase="cancellation_requested", runtime=runtime, workflow_id=workflow_id, + request_id=original["request_id"], cleanup_deadline_at=original["cleanup_deadline_at"]) + + +async def fired(client: Client, runtime: str, scenario: str) -> None: + deadline = asyncio.get_running_loop().time() + 60 + while True: + _, execution, history = await observe(client, runtime, scenario) + if events(history, "TimerFired"): + if events(history, "WorkflowCompleted"): + raise RuntimeError("Workflow completed while its SDK worker should be killed.") + emit(phase="fired_without_worker", runtime=runtime, scenario=scenario, + workflow_id=execution.workflow_id, run_id=execution.run_id) + return + if asyncio.get_running_loop().time() >= deadline: + raise TimeoutError(f"{runtime} timer did not fire while worker was absent.") + await asyncio.sleep(0.25) + + +async def verify(client: Client, runtime: str, scenario: str) -> None: + handle, execution, history = await observe(client, runtime, scenario) + if scenario == "cancellation": + deadline = asyncio.get_running_loop().time() + 60 + while True: + scheduled = events(history, "TimerScheduled") + due = timestamp(scheduled[0]["payload"]["fire_at"]) + if execution.status == "cancelled" and datetime.now(timezone.utc) > due: + break + if execution.status in ("failed", "terminated", "completed"): + raise RuntimeError(f"Unexpected cancellation outcome: {history!r}") + if asyncio.get_running_loop().time() >= deadline: + raise TimeoutError(f"{runtime} cancellation did not converge.") + await asyncio.sleep(0.25) + handle, execution, history = await observe(client, runtime, scenario) + for event_type in ("TimerScheduled", "TimerCancelled", "CooperativeCancellationRequested", + "CooperativeCancellationDelivered", "WorkflowCancelled"): + if len(events(history, event_type)) != 1: + raise RuntimeError(f"Expected one {event_type}: {history!r}") + if events(history, "TimerFired") or events(history, "WorkflowCompleted"): + raise RuntimeError(f"Cancelled timer fired or workflow completed: {history!r}") + cancelled = events(history, "TimerCancelled")[0]["payload"] + if cancelled["timer_id"] != scheduled[0]["payload"]["timer_id"]: + raise RuntimeError("Cancelled timer identity differs from its scheduled identity.") + result = None + else: + result = await handle.result(timeout=120, poll_interval=0.25) + handle, execution, history = await observe(client, runtime, scenario) + expected = {"workflow_runtime": runtime, "request": execution.workflow_id, "timer_seconds": TIMER_SECONDS} + if not isinstance(result, dict) or any(result.get(key) != value for key, value in expected.items()): + raise RuntimeError(f"Unexpected {runtime} result: {result!r}") + if execution.status != "completed": + raise RuntimeError(f"Workflow did not complete: {execution!r}") + for event_type in ("TimerScheduled", "TimerFired", "WorkflowCompleted"): + if len(events(history, event_type)) != 1: + raise RuntimeError(f"Expected one {event_type}: {history!r}") + scheduled = events(history, "TimerScheduled")[0] + fired_event = events(history, "TimerFired")[0] + if scheduled["payload"]["timer_id"] != fired_event["payload"]["timer_id"]: + raise RuntimeError("Fired timer identity differs from its scheduled identity.") + if timestamp(fired_event["payload"]["fired_at"]) < timestamp(scheduled["payload"]["fire_at"]): + raise RuntimeError("Timer fired before its original deadline.") + types = [event["event_type"] for event in history["events"]] + if not types.index("TimerScheduled") < types.index("TimerFired") < types.index("WorkflowCompleted"): + raise RuntimeError(f"Timer history is out of order: {types!r}") + if scenario == "worker-restart": + stopped = timestamp(required("DURABLE_WORKFLOW_WORKER_STOPPED_AT")) + restarted = timestamp(required("DURABLE_WORKFLOW_WORKER_RESTART_AT")) + if not stopped <= timestamp(fired_event["payload"]["fired_at"]) < restarted: + raise RuntimeError("Timer did not fire while the SDK worker was absent.") + emit(phase="verified", runtime=runtime, scenario=scenario, workflow_id=execution.workflow_id, + run_id=execution.run_id, status=execution.status, result=result, + history_events=[event["event_type"] for event in history["events"]], + timer_events=[event for event in history["events"] if event["event_type"].startswith("Timer")]) + + +async def main(phase: str, scenario: str) -> None: + async with Client( + required("DURABLE_WORKFLOW_SERVER_URL"), + namespace=required("DURABLE_WORKFLOW_NAMESPACE"), + control_token=required("DURABLE_WORKFLOW_AUTH_TOKEN"), + ) as client: + action = {"start": pending, "fired": fired, "verify": verify}[phase] + await asyncio.gather(*(action(client, runtime, scenario) for runtime in RUNTIMES)) + + +if __name__ == "__main__": + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("phase", choices=("start", "fired", "verify")) + parser.add_argument("scenario", choices=("completion", "worker-restart", "server-restart", "cancellation")) + args = parser.parse_args() + asyncio.run(main(args.phase, args.scenario)) diff --git a/polyglot/python_workflow/workflow.py b/polyglot/python_workflow/workflow.py index edafc763..9237c4a7 100644 --- a/polyglot/python_workflow/workflow.py +++ b/polyglot/python_workflow/workflow.py @@ -78,6 +78,13 @@ def run(self, ctx, request): # type: ignore[no-untyped-def] } +@workflow.defn(name="polyglot.python.timer") +class PythonTimerWorkflow: + def run(self, ctx, request): + yield ctx.sleep(30) + return {"workflow_runtime": "python", "request": request, "timer_seconds": 30} + + @workflow.defn(name="polyglot.python-to-php.greeter") class PythonToPhpGreeterWorkflow: def run(self, ctx, request): # type: ignore[no-untyped-def] @@ -413,6 +420,7 @@ async def main() -> int: task_queue=TASK_QUEUE, workflows=[ PythonGreeterWorkflow, + PythonTimerWorkflow, PythonToPhpGreeterWorkflow, PythonToPhpTypeRoundtripWorkflow, PythonToPhpBinaryTypeRoundtripWorkflow, diff --git a/polyglot/qualified-artifact-tuple.json b/polyglot/qualified-artifact-tuple.json index ec34b39a..892bb6fd 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": "2.1.5", + "sdk-rust": "3.2.0", "server": "2.5.9", "waterline": "2.3.1", "workflow": "2.5.3" diff --git a/polyglot/timers/README.md b/polyglot/timers/README.md index 968ff0eb..9e221e11 100644 --- a/polyglot/timers/README.md +++ b/polyglot/timers/README.md @@ -1,4 +1,52 @@ -# Rust durable timer across a worker restart +# Published SDK durable timers + +Run PHP, Python and Rust timer workflows against the published Server with: + +```bash +scripts/playground doctor +while IFS= read -r assignment; do export "$assignment"; done \ + < <(scripts/resolve-current-artifacts.sh) +SDK_TIMERS_COMPOSE_PROJECT_NAME=sample-app-sdk-timers scripts/sdk-timers.sh +``` + +Use the prepared Sample App development container with Docker Compose. The +checked-in artifact tuple selects exact published packages and images. The +experiment builds workers from those packages, starts a disposable MySQL/Redis +Server stack, and checks four scenarios for each workflow language: + +- Normal completion after the original 30-second deadline, with exactly one + scheduled timer, matching fire and completed result. +- Worker `SIGKILL` while waiting, timer fire while all SDK workers are absent, + then cold replay in replacement processes without rescheduling or duplicate + completion. +- Server and timer-queue restart across the deadline, preserving timer identity + and completing once from the original history. +- Cooperative cancellation while waiting, duplicate requests preserving the + original request identity and cleanup deadline, one delivered cancellation + and cancelled timer, and no fire or completion after the timer's original + due time has passed. + +Every start also checks the public API's `waiting` status. The Python SDK is the +observer and control client. PHP, Python and Rust author and execute their own +timer workflows. Timers have no remote activity or cross-language timer-worker +direction to multiply into a workflow/activity matrix. + +The runner prints the exact tuple, Server digest, timestamps, per-language +results and persisted timer events. Record the command, Sample App commit, +UTC interval and twelve scenario outcomes in the owning GitHub issue. Its exit +trap removes the isolated Compose stack and volumes on success or failure. +If interrupted externally, repeat the removal with the same project name: + +```bash +COMPOSE_PROJECT_NAME=sample-app-sdk-timers COMPOSE_PROFILES=timers \ + docker compose -f polyglot/docker-compose.yml down --volumes --remove-orphans +``` + +This is focused SDK timer coverage. Concurrent distinct deadlines, nested +cancellation scopes and timer-bearing application upgrades have separate +qualification requirements. + +## Single Rust restart example This focused service-mode experiment runs a timer authored with the published Rust SDK against a published Server image. The Python SDK client observes one diff --git a/scripts/resolve-current-artifacts.sh b/scripts/resolve-current-artifacts.sh index 10810f68..d0976475 100755 --- a/scripts/resolve-current-artifacts.sh +++ b/scripts/resolve-current-artifacts.sh @@ -26,6 +26,7 @@ const path = process.argv[2]; const expectedSchema = 'durable-workflow.sample-app.polyglot-qualified-artifact-tuple'; const keys = ['server', 'cli', 'sdk-php', 'sdk-python', 'sdk-rust', 'workflow', 'waterline']; const stableV2 = /^2\.\d+\.\d+(?:\+[0-9A-Za-z.-]+)?$/; +const stableRust = /^[23]\.\d+\.\d+(?:\+[0-9A-Za-z.-]+)?$/; let tuple; try { @@ -50,8 +51,9 @@ if (unknown.length > 0) { for (const key of keys) { const version = artifacts[key]; - if (typeof version !== 'string' || !stableV2.test(version)) { - throw new Error(`${path} artifact ${key} must be a stable 2.x version`); + const supported = key === 'sdk-rust' ? stableRust : stableV2; + if (typeof version !== 'string' || !supported.test(version)) { + throw new Error(`${path} artifact ${key} must be a stable ${key === 'sdk-rust' ? '2.x or 3.x' : '2.x'} version`); } process.stdout.write(`${key}=${version}\n`); } @@ -66,10 +68,16 @@ done < <(parse_tuple) stable_version() { local name="$1" local value="$2" + local supported_major=2 + local supported_label=2.x + if [[ "$name" == SAMPLE_APP_RUST_SDK_VERSION ]]; then + supported_major='[23]' + supported_label='2.x or 3.x' + fi - if [[ ! "$value" =~ ^2\.[0-9]+\.[0-9]+(\+[0-9A-Za-z.-]+)?$ ]]; then - printf 'resolve-current-artifacts: %s must be a stable 2.x version; received %s\n' \ - "$name" "$value" >&2 + if [[ ! "$value" =~ ^${supported_major}\.[0-9]+\.[0-9]+(\+[0-9A-Za-z.-]+)?$ ]]; then + printf 'resolve-current-artifacts: %s must be a stable %s version; received %s\n' \ + "$name" "$supported_label" "$value" >&2 exit 1 fi diff --git a/scripts/sdk-timers.sh b/scripts/sdk-timers.sh new file mode 100755 index 00000000..9555e0cf --- /dev/null +++ b/scripts/sdk-timers.sh @@ -0,0 +1,64 @@ +#!/usr/bin/env bash +set -euo pipefail + +if [[ "${1:-}" == --help ]]; then + printf '%s\n' 'Usage: scripts/sdk-timers.sh' \ + 'Runs completion, worker SIGKILL/cold replay, Server restart and cooperative cancellation for PHP/Python/Rust timers.' \ + 'Requires Docker Compose and exact artifact assignments from scripts/resolve-current-artifacts.sh.' \ + 'SDK_TIMERS_COMPOSE_PROJECT_NAME selects an isolated project. All project resources are removed on exit.' + exit 0 +fi +[[ $# == 0 ]] || { printf '%s\n' 'Use --help for usage.' >&2; exit 2; } + +repo_root="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)" +export COMPOSE_PROJECT_NAME="${SDK_TIMERS_COMPOSE_PROJECT_NAME:-sample-app-sdk-timers-$(date -u +%Y%m%d%H%M%S)}" +[[ "$COMPOSE_PROJECT_NAME" =~ ^[a-z0-9][a-z0-9_-]*$ ]] || exit 2 +export COMPOSE_PROFILES=timers +export DURABLE_WORKFLOW_TIMER_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) + +cleanup() { + local code=$? + if [[ "$code" != 0 ]]; then "${compose[@]}" logs --no-color --timestamps; fi + "${compose[@]}" down --volumes --remove-orphans || return 1 + return "$code" +} +trap cleanup EXIT + +printf 'SDK timers 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 timer-queue +"${compose[@]}" up -d --wait --wait-timeout 180 --no-build server timer-queue "${workers[@]}" +docker image inspect "$DURABLE_SERVER_IMAGE" --format '{{json .RepoDigests}}' + +client() { + "${compose[@]}" run --rm --no-deps --user 1000:1000 \ + -e DURABLE_WORKFLOW_TIMER_ID -e DURABLE_WORKFLOW_WORKER_STOPPED_AT -e DURABLE_WORKFLOW_WORKER_RESTART_AT \ + smoke python /app/scripts/sdk_timers.py "$@" +} + +client start completion +client verify completion + +client start worker-restart +"${compose[@]}" kill --signal SIGKILL "${workers[@]}" +export DURABLE_WORKFLOW_WORKER_STOPPED_AT="$(date -u +%Y-%m-%dT%H:%M:%S.%NZ)" +client fired worker-restart +export DURABLE_WORKFLOW_WORKER_RESTART_AT="$(date -u +%Y-%m-%dT%H:%M:%S.%NZ)" +"${compose[@]}" up -d --wait --no-build "${workers[@]}" +client verify worker-restart + +client start server-restart +"${compose[@]}" stop server timer-queue +sleep 32 +"${compose[@]}" up -d --wait --wait-timeout 180 --no-build server timer-queue +client verify server-restart + +client start cancellation +client verify cancellation +printf 'SDK timers pass: %s\n' "$(date -u +%Y-%m-%dT%H:%M:%SZ)" diff --git a/tests/Unit/PolyglotComposeContractTest.php b/tests/Unit/PolyglotComposeContractTest.php index 74ae7389..aac4a9ad 100644 --- a/tests/Unit/PolyglotComposeContractTest.php +++ b/tests/Unit/PolyglotComposeContractTest.php @@ -25,8 +25,8 @@ public function test_stable_artifact_tuple_drives_every_public_runtime(): void array_keys($tuple['artifacts']), ); - foreach ($tuple['artifacts'] as $version) { - $this->assertMatchesRegularExpression('/^2\.\d+\.\d+$/', $version); + foreach ($tuple['artifacts'] as $artifact => $version) { + $this->assertMatchesRegularExpression($artifact === 'sdk-rust' ? '/^[23]\.\d+\.\d+$/' : '/^2\.\d+\.\d+$/', $version); } $assignments = $this->resolveArtifacts(); @@ -46,11 +46,11 @@ public function test_stable_artifact_tuple_drives_every_public_runtime(): void public function test_artifact_resolver_accepts_stable_overrides_and_rejects_prereleases(): void { $assignments = $this->resolveArtifacts([ - 'SAMPLE_APP_RUST_SDK_VERSION' => '2.3.4', + 'SAMPLE_APP_RUST_SDK_VERSION' => '3.2.0', 'SAMPLE_APP_PHP_SDK_PIN' => 'durable-workflow/sdk:2.4.5', ]); - $this->assertSame('2.3.4', $assignments['DURABLE_WORKFLOW_RUST_SDK_VERSION']); + $this->assertSame('3.2.0', $assignments['DURABLE_WORKFLOW_RUST_SDK_VERSION']); $this->assertSame('2.4.5', $assignments['DURABLE_WORKFLOW_PHP_SDK_VERSION']); $this->assertSame('durable-workflow/sdk:2.4.5', $assignments['DURABLE_WORKFLOW_PHP_SDK_PIN']); @@ -61,7 +61,7 @@ public function test_artifact_resolver_accepts_stable_overrides_and_rejects_prer $process->run(); $this->assertFalse($process->isSuccessful()); - $this->assertStringContainsString('must be a stable 2.x version', $process->getErrorOutput()); + $this->assertStringContainsString('must be a stable 2.x or 3.x version', $process->getErrorOutput()); } public function test_polyglot_compose_uses_isolated_runtime_services_and_resolved_artifacts(): void From ccc51e5736970284113448087c0436184e5dfdc7 Mon Sep 17 00:00:00 2001 From: Durable Workflow Date: Wed, 7 Oct 2026 22:49:43 +0000 Subject: [PATCH 2/5] Use published Rust outage recovery and preserve original timer boundaries --- playground/templates/rust/Cargo.lock | 4 +- playground/templates/rust/Cargo.toml | 2 +- polyglot/docker-compose.yml | 7 + polyglot/php_worker/worker.php | 11 +- polyglot/python_worker/scripts/sdk_timers.py | 30 +++++ .../python_worker/tests/test_sdk_timers.py | 121 ++++++++++++++++++ polyglot/python_workflow/workflow.py | 4 + polyglot/qualified-artifact-tuple.json | 2 +- polyglot/rust_worker/Cargo.lock | 4 +- polyglot/rust_worker/Cargo.toml | 2 +- polyglot/rust_worker/src/main.rs | 2 + polyglot/timers/README.md | 10 +- scripts/sdk-timers.sh | 34 ++++- 13 files changed, 217 insertions(+), 16 deletions(-) create mode 100644 polyglot/python_worker/tests/test_sdk_timers.py diff --git a/playground/templates/rust/Cargo.lock b/playground/templates/rust/Cargo.lock index 0fcbde17..3f89d396 100644 --- a/playground/templates/rust/Cargo.lock +++ b/playground/templates/rust/Cargo.lock @@ -268,9 +268,9 @@ dependencies = [ [[package]] name = "durable-workflow" -version = "2.1.5" +version = "3.2.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "5b47b06c299983d2d0de407d17e340a5d871be636ca0105f8058f7bfaaf60a91" +checksum = "9d0482e9e070aefa6844a16132ca35eea8dcebc1524859f879961e57fc0985ae" dependencies = [ "apache-avro", "base64", diff --git a/playground/templates/rust/Cargo.toml b/playground/templates/rust/Cargo.toml index ae7f36ee..753b63ed 100644 --- a/playground/templates/rust/Cargo.toml +++ b/playground/templates/rust/Cargo.toml @@ -6,6 +6,6 @@ publish = false rust-version = "1.86" [dependencies] -durable-workflow = "2.0" +durable-workflow = "3.2" serde_json = "=1.0.150" tokio = { version = "=1.47.1", features = ["macros", "rt-multi-thread", "signal", "time"] } diff --git a/polyglot/docker-compose.yml b/polyglot/docker-compose.yml index 384a7570..cba92f90 100644 --- a/polyglot/docker-compose.yml +++ b/polyglot/docker-compose.yml @@ -172,6 +172,7 @@ services: condition: service_healthy python-workflow-worker: + init: true build: context: ./python_workflow args: @@ -182,6 +183,7 @@ services: DURABLE_WORKFLOW_AUTH_TOKEN: "test-token" DURABLE_WORKFLOW_NAMESPACE: default DURABLE_WORKFLOW_POLL_TIMEOUT_SECONDS: "5" + POLYGLOT_TIMER_COOPERATIVE: "${POLYGLOT_TIMER_COOPERATIVE:-0}" POLYGLOT_PY_TASK_QUEUE: polyglot-python POLYGLOT_PY2PHP_TASK_QUEUE: polyglot-python-to-php POLYGLOT_TO_RUST_TASK_QUEUE: polyglot-to-rust @@ -208,6 +210,7 @@ services: condition: service_healthy php-same-workflow-worker: + init: true build: *php-sdk-worker-build image: *php-sdk-worker-image environment: @@ -215,6 +218,7 @@ services: DURABLE_WORKFLOW_AUTH_TOKEN: "test-token" DURABLE_WORKFLOW_NAMESPACE: default POLYGLOT_PHP_TASK_QUEUE: polyglot-php + POLYGLOT_TIMER_COOPERATIVE: "${POLYGLOT_TIMER_COOPERATIVE:-0}" depends_on: server: condition: service_healthy @@ -301,6 +305,7 @@ services: command: ["php", "/app/worker.php", "--mode=activity", "--task-queue=polyglot-python-to-php", "--poll-timeout=5"] rust-workflow-worker: + init: true build: context: ./rust_worker args: @@ -313,6 +318,7 @@ services: DURABLE_WORKFLOW_RUST_SDK_VERSION: *durable-rust-sdk-version APACHE_AVRO_RUST_VERSION: *durable-rust-avro-version POLYGLOT_RUST_MODE: workflow + POLYGLOT_TIMER_COOPERATIVE: "${POLYGLOT_TIMER_COOPERATIVE:-0}" POLYGLOT_RUST_TASK_QUEUE: polyglot-rust POLYGLOT_PHP2PY_TASK_QUEUE: polyglot-php-to-python POLYGLOT_PY2PHP_TASK_QUEUE: polyglot-python-to-php @@ -387,6 +393,7 @@ services: retries: 30 smoke: + init: true build: context: ./python_worker args: diff --git a/polyglot/php_worker/worker.php b/polyglot/php_worker/worker.php index 0811b28d..a5a16c81 100644 --- a/polyglot/php_worker/worker.php +++ b/polyglot/php_worker/worker.php @@ -9,6 +9,7 @@ use DurableWorkflow\Codec\PayloadCodec; use DurableWorkflow\Exception\ActivityFailed; use DurableWorkflow\Exception\CodecException; +use DurableWorkflow\Version; use DurableWorkflow\Worker; use DurableWorkflow\Worker\PollResponse; use DurableWorkflow\Worker\QueryContext; @@ -800,7 +801,13 @@ function runStandaloneWorker(): void throw new RuntimeException('Expected --mode=workflow, --mode=activity, --mode=query, or --mode=replay-fixtures.'); } - $client = new Client($serverUrl, token: $token, namespace: $namespace); + $cooperative = getenv('POLYGLOT_TIMER_COOPERATIVE') === '1'; + $client = new Client( + $serverUrl, + token: $token, + namespace: $namespace, + workerProtocolVersion: $cooperative ? '1.20' : Version::WORKER_PROTOCOL, + ); if ($mode === 'activity') { runActivityWorker($client, $workerId, $taskQueue, $pollTimeout); @@ -812,7 +819,7 @@ function runStandaloneWorker(): void return; } - $worker = new Worker($client, $taskQueue, $workerId); + $worker = new Worker($client, $taskQueue, $workerId, enableCooperativeCancellation: $cooperative); configureWorkflows($worker, $client->payloadCodec()); fwrite(STDOUT, sprintf( "polyglot php worker registered: id=%s queue=%s types=[%s]\n", diff --git a/polyglot/python_worker/scripts/sdk_timers.py b/polyglot/python_worker/scripts/sdk_timers.py index b20940e8..2eab8d13 100644 --- a/polyglot/python_worker/scripts/sdk_timers.py +++ b/polyglot/python_worker/scripts/sdk_timers.py @@ -34,6 +34,22 @@ def emit(**record) -> None: print(json.dumps(record, sort_keys=True), flush=True) +def original_pending(runtime: str, scenario: str) -> dict: + records = [json.loads(line) for line in required("DURABLE_WORKFLOW_TIMER_PENDING").splitlines()] + matches = [record for record in records if record.get("phase") == "pending" + and record.get("runtime") == runtime and record.get("scenario") == scenario] + if len(matches) != 1: + raise RuntimeError("Expected exactly one original pending observation for this cell.") + return matches[0] + + +def check_original(runtime: str, scenario: str, execution, scheduled: dict) -> None: + original = original_pending(runtime, scenario) + if (execution.workflow_id != original["workflow_id"] or execution.run_id != original["run_id"] + or scheduled["payload"] != original["timer"]): + raise RuntimeError("Run identity or original scheduled timer changed after the interruption.") + + async def observe(client: Client, runtime: str, scenario: str): workflow_id = f"{required('DURABLE_WORKFLOW_TIMER_ID')}-{scenario}-{runtime}" handle = client.get_workflow_handle(workflow_id, workflow_type=f"polyglot.{runtime}.timer") @@ -120,6 +136,7 @@ async def verify(client: Client, runtime: str, scenario: str) -> None: raise RuntimeError(f"Expected one {event_type}: {history!r}") if events(history, "TimerFired") or events(history, "WorkflowCompleted"): raise RuntimeError(f"Cancelled timer fired or workflow completed: {history!r}") + check_original(runtime, scenario, execution, scheduled[0]) cancelled = events(history, "TimerCancelled")[0]["payload"] if cancelled["timer_id"] != scheduled[0]["payload"]["timer_id"]: raise RuntimeError("Cancelled timer identity differs from its scheduled identity.") @@ -137,6 +154,7 @@ async def verify(client: Client, runtime: str, scenario: str) -> None: raise RuntimeError(f"Expected one {event_type}: {history!r}") scheduled = events(history, "TimerScheduled")[0] fired_event = events(history, "TimerFired")[0] + check_original(runtime, scenario, execution, scheduled) if scheduled["payload"]["timer_id"] != fired_event["payload"]["timer_id"]: raise RuntimeError("Fired timer identity differs from its scheduled identity.") if timestamp(fired_event["payload"]["fired_at"]) < timestamp(scheduled["payload"]["fire_at"]): @@ -149,8 +167,20 @@ async def verify(client: Client, runtime: str, scenario: str) -> None: restarted = timestamp(required("DURABLE_WORKFLOW_WORKER_RESTART_AT")) if not stopped <= timestamp(fired_event["payload"]["fired_at"]) < restarted: raise RuntimeError("Timer did not fire while the SDK worker was absent.") + if scenario == "server-restart": + stopped = timestamp(required("DURABLE_WORKFLOW_SERVER_STOPPED_AT")) + restarted = timestamp(required("DURABLE_WORKFLOW_SERVER_RESTART_AT")) + due = timestamp(scheduled["payload"]["fire_at"]) + if not stopped < due < restarted: + raise RuntimeError("Server downtime did not cross the original timer deadline.") + if timestamp(fired_event["payload"]["fired_at"]) < restarted: + raise RuntimeError("Timer fired before Server restart began.") emit(phase="verified", runtime=runtime, scenario=scenario, workflow_id=execution.workflow_id, run_id=execution.run_id, status=execution.status, result=result, + interruption={key: required(f"DURABLE_WORKFLOW_{key.upper()}") + for key in ({"worker-restart": ("worker_stopped_at", "worker_restart_at"), + "server-restart": ("server_stopped_at", "server_restart_at")} + .get(scenario, ()))}, history_events=[event["event_type"] for event in history["events"]], timer_events=[event for event in history["events"] if event["event_type"].startswith("Timer")]) diff --git a/polyglot/python_worker/tests/test_sdk_timers.py b/polyglot/python_worker/tests/test_sdk_timers.py new file mode 100644 index 00000000..04ef944a --- /dev/null +++ b/polyglot/python_worker/tests/test_sdk_timers.py @@ -0,0 +1,121 @@ +from __future__ import annotations + +import copy +import importlib.util +import sys +import types +import unittest +from pathlib import Path +from unittest.mock import AsyncMock, patch + + +if "durable_workflow" not in sys.modules: + stub = types.ModuleType("durable_workflow") + stub.Client = object + sys.modules["durable_workflow"] = stub +path = Path(__file__).parents[1] / "scripts" / "sdk_timers.py" +spec = importlib.util.spec_from_file_location("sdk_timers", path) +timers = importlib.util.module_from_spec(spec) +spec.loader.exec_module(timers) + + +def event(kind, **payload): + return {"event_type": kind, "payload": payload} + + +class TimerEvidenceTest(unittest.IsolatedAsyncioTestCase): + def setUp(self): + self.execution = types.SimpleNamespace(workflow_id="timer-rust", run_id="run-1", status="completed") + self.handle = types.SimpleNamespace(result=AsyncMock(return_value={ + "workflow_runtime": "rust", "request": "timer-rust", "timer_seconds": 30, + })) + self.history = {"events": [ + event("WorkflowStarted"), + event("TimerScheduled", timer_id="timer-1", fire_at="2026-10-05T00:00:30Z"), + event("TimerFired", timer_id="timer-1", fired_at="2026-10-05T00:00:30Z"), + event("WorkflowCompleted"), + ]} + self.original = {"workflow_id": "timer-rust", "run_id": "run-1", + "timer": copy.deepcopy(self.history["events"][1]["payload"])} + + async def verify(self, history, scenario="completion"): + with patch.object(timers, "observe", AsyncMock(return_value=(self.handle, self.execution, history))), \ + patch.object(timers, "original_pending", return_value=self.original), \ + patch.object(timers, "emit"): + await timers.verify(None, "rust", scenario) + + async def test_matching_timer_and_one_result_pass(self): + await self.verify(self.history) + + async def test_duplicate_or_missing_durable_events_cannot_pass(self): + for kind in ("TimerScheduled", "TimerFired", "WorkflowCompleted"): + for duplicate in (False, True): + with self.subTest(kind=kind, duplicate=duplicate): + history = copy.deepcopy(self.history) + selected = next(item for item in history["events"] if item["event_type"] == kind) + if duplicate: + history["events"].append(selected) + else: + history["events"].remove(selected) + with self.assertRaises(RuntimeError): + await self.verify(history) + + async def test_early_fire_cannot_pass(self): + self.history["events"][2]["payload"]["fired_at"] = "2026-10-05T00:00:29Z" + with self.assertRaisesRegex(RuntimeError, "before its original deadline"): + await self.verify(self.history) + + async def test_wrong_timer_identity_cannot_pass(self): + self.history["events"][2]["payload"]["timer_id"] = "another-timer" + with self.assertRaisesRegex(RuntimeError, "identity differs"): + await self.verify(self.history) + + async def test_rewritten_original_deadline_cannot_pass(self): + self.history["events"][1]["payload"]["fire_at"] = "2026-10-05T00:00:29Z" + with self.assertRaisesRegex(RuntimeError, "original scheduled timer changed"): + await self.verify(self.history) + + async def test_replacement_run_cannot_pass(self): + self.execution.run_id = "run-2" + with self.assertRaisesRegex(RuntimeError, "Run identity"): + await self.verify(self.history) + + async def test_fire_outside_worker_absence_cannot_pass(self): + with patch.dict("os.environ", { + "DURABLE_WORKFLOW_WORKER_STOPPED_AT": "2026-10-05T00:00:10Z", + "DURABLE_WORKFLOW_WORKER_RESTART_AT": "2026-10-05T00:00:20Z", + }), self.assertRaisesRegex(RuntimeError, "while the SDK worker was absent"): + await self.verify(self.history, "worker-restart") + + async def test_cancelled_timer_is_checked_after_due_time(self): + self.execution.status = "cancelled" + self.history["events"] = [ + self.history["events"][1], + event("CooperativeCancellationRequested"), + event("CooperativeCancellationDelivered"), + event("TimerCancelled", timer_id="timer-1"), + event("WorkflowCancelled"), + ] + await self.verify(self.history, "cancellation") + self.history["events"].append(event("TimerFired", timer_id="timer-1")) + with self.assertRaisesRegex(RuntimeError, "Cancelled timer fired"): + await self.verify(self.history, "cancellation") + + async def test_server_restart_must_cross_original_deadline(self): + self.history["events"][2]["payload"]["fired_at"] = "2026-10-05T00:00:33Z" + with patch.dict("os.environ", { + "DURABLE_WORKFLOW_SERVER_STOPPED_AT": "2026-10-05T00:00:10Z", + "DURABLE_WORKFLOW_SERVER_RESTART_AT": "2026-10-05T00:00:32Z", + }): + await self.verify(self.history, "server-restart") + with patch.dict("os.environ", { + "DURABLE_WORKFLOW_SERVER_STOPPED_AT": "2026-10-05T00:00:31Z", + }), self.assertRaisesRegex(RuntimeError, "did not cross"): + await self.verify(self.history, "server-restart") + self.history["events"][2]["payload"]["fired_at"] = "2026-10-05T00:00:31Z" + with self.assertRaisesRegex(RuntimeError, "before Server restart"): + await self.verify(self.history, "server-restart") + + +if __name__ == "__main__": + unittest.main() diff --git a/polyglot/python_workflow/workflow.py b/polyglot/python_workflow/workflow.py index 9237c4a7..aed13c77 100644 --- a/polyglot/python_workflow/workflow.py +++ b/polyglot/python_workflow/workflow.py @@ -407,6 +407,9 @@ async def main() -> int: "POLYGLOT_PY_WORKER_ID", f"py-workflow-worker-{socket.gethostname()}", ) + cooperative = os.environ.get("POLYGLOT_TIMER_COOPERATIVE") == "1" + if cooperative: + os.environ["DURABLE_WORKFLOW_WORKER_PROTOCOL_VERSION"] = "1.20" async with Client( server_url, @@ -418,6 +421,7 @@ async def main() -> int: worker = Worker( client, task_queue=TASK_QUEUE, + capabilities=["cooperative_cancellation"] if cooperative else [], workflows=[ PythonGreeterWorkflow, PythonTimerWorkflow, diff --git a/polyglot/qualified-artifact-tuple.json b/polyglot/qualified-artifact-tuple.json index 892bb6fd..2c0a6e97 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.0", + "sdk-rust": "3.2.1", "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 0520e219..eb9fbe93 100644 --- a/polyglot/rust_worker/Cargo.lock +++ b/polyglot/rust_worker/Cargo.lock @@ -244,9 +244,9 @@ dependencies = [ [[package]] name = "durable-workflow" -version = "2.1.5" +version = "3.2.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "5b47b06c299983d2d0de407d17e340a5d871be636ca0105f8058f7bfaaf60a91" +checksum = "9d0482e9e070aefa6844a16132ca35eea8dcebc1524859f879961e57fc0985ae" dependencies = [ "apache-avro", "base64", diff --git a/polyglot/rust_worker/Cargo.toml b/polyglot/rust_worker/Cargo.toml index 175e29dc..de541e71 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 = "2.0" +durable-workflow = "3.2" tokio = { version = "1.47.1", features = ["macros", "rt-multi-thread", "signal"] } [dev-dependencies] diff --git a/polyglot/rust_worker/src/main.rs b/polyglot/rust_worker/src/main.rs index 5ea41a2f..c17db823 100644 --- a/polyglot/rust_worker/src/main.rs +++ b/polyglot/rust_worker/src/main.rs @@ -50,6 +50,8 @@ async fn run_workflow_worker(client: Client) -> Result<()> { let php_queue = env_value("POLYGLOT_PY2PHP_TASK_QUEUE", "polyglot-python-to-php"); let mut worker = Worker::new(client, rust_queue.clone()) .worker_id("rust-workflow-worker") + .recover_transient_outages(true) + .cooperative_cancellation(env_value("POLYGLOT_TIMER_COOPERATIVE", "0") == "1") .poll_timeout(Duration::from_secs(5)); worker.register_activity("polyglot.rust.echo", |_ctx, args| async move { diff --git a/polyglot/timers/README.md b/polyglot/timers/README.md index 9e221e11..a901b6ab 100644 --- a/polyglot/timers/README.md +++ b/polyglot/timers/README.md @@ -31,10 +31,18 @@ observer and control client. PHP, Python and Rust author and execute their own timer workflows. Timers have no remote activity or cross-language timer-worker direction to multiply into a workflow/activity matrix. +The timer runner enables each worker's cooperative cancellation capability and +protocol 1.20 with `POLYGLOT_TIMER_COOPERATIVE=1`. Ordinary polyglot commands keep +their existing capability selection. The long-running Rust workflow worker uses +`recover_transient_outages(true)` to remain alive during the Server +restart with capped retry backoff. Its startup and permanent errors still require +operator attention or a process supervisor. + The runner prints the exact tuple, Server digest, timestamps, per-language results and persisted timer events. Record the command, Sample App commit, UTC interval and twelve scenario outcomes in the owning GitHub issue. Its exit -trap removes the isolated Compose stack and volumes on success or failure. +trap removes the isolated Compose stack, volumes and task worker images on +success or failure. Shared published base images are preserved. If interrupted externally, repeat the removal with the same project name: ```bash diff --git a/scripts/sdk-timers.sh b/scripts/sdk-timers.sh index 9555e0cf..9a09085c 100755 --- a/scripts/sdk-timers.sh +++ b/scripts/sdk-timers.sh @@ -14,14 +14,26 @@ repo_root="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)" export COMPOSE_PROJECT_NAME="${SDK_TIMERS_COMPOSE_PROJECT_NAME:-sample-app-sdk-timers-$(date -u +%Y%m%d%H%M%S)}" [[ "$COMPOSE_PROJECT_NAME" =~ ^[a-z0-9][a-z0-9_-]*$ ]] || exit 2 export COMPOSE_PROFILES=timers +export POLYGLOT_TIMER_COOPERATIVE=1 export DURABLE_WORKFLOW_TIMER_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; fi + 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 @@ -39,13 +51,21 @@ docker image inspect "$DURABLE_SERVER_IMAGE" --format '{{json .RepoDigests}}' client() { "${compose[@]}" run --rm --no-deps --user 1000:1000 \ -e DURABLE_WORKFLOW_TIMER_ID -e DURABLE_WORKFLOW_WORKER_STOPPED_AT -e DURABLE_WORKFLOW_WORKER_RESTART_AT \ + -e DURABLE_WORKFLOW_SERVER_STOPPED_AT -e DURABLE_WORKFLOW_SERVER_RESTART_AT \ + -e DURABLE_WORKFLOW_TIMER_PENDING \ smoke python /app/scripts/sdk_timers.py "$@" } -client start completion +start() { + DURABLE_WORKFLOW_TIMER_PENDING="$(client start "$1")" + export DURABLE_WORKFLOW_TIMER_PENDING + printf '%s\n' "$DURABLE_WORKFLOW_TIMER_PENDING" +} + +start completion client verify completion -client start worker-restart +start worker-restart "${compose[@]}" kill --signal SIGKILL "${workers[@]}" export DURABLE_WORKFLOW_WORKER_STOPPED_AT="$(date -u +%Y-%m-%dT%H:%M:%S.%NZ)" client fired worker-restart @@ -53,12 +73,14 @@ export DURABLE_WORKFLOW_WORKER_RESTART_AT="$(date -u +%Y-%m-%dT%H:%M:%S.%NZ)" "${compose[@]}" up -d --wait --no-build "${workers[@]}" client verify worker-restart -client start server-restart -"${compose[@]}" stop server timer-queue +start server-restart +"${compose[@]}" stop --timeout 2 timer-queue server +export DURABLE_WORKFLOW_SERVER_STOPPED_AT="$(date -u +%Y-%m-%dT%H:%M:%S.%NZ)" sleep 32 +export DURABLE_WORKFLOW_SERVER_RESTART_AT="$(date -u +%Y-%m-%dT%H:%M:%S.%NZ)" "${compose[@]}" up -d --wait --wait-timeout 180 --no-build server timer-queue client verify server-restart -client start cancellation +start cancellation client verify cancellation printf 'SDK timers pass: %s\n' "$(date -u +%Y-%m-%dT%H:%M:%SZ)" From 2c50bd9003ecdcd0241c9ce5c5c9831a58934294 Mon Sep 17 00:00:00 2001 From: Durable Workflow Date: Wed, 7 Oct 2026 22:54:33 +0000 Subject: [PATCH 3/5] Wait for public waiting state after timer history becomes visible --- polyglot/python_worker/scripts/sdk_timers.py | 2 +- polyglot/python_worker/tests/test_sdk_timers.py | 16 ++++++++++++++++ 2 files changed, 17 insertions(+), 1 deletion(-) diff --git a/polyglot/python_worker/scripts/sdk_timers.py b/polyglot/python_worker/scripts/sdk_timers.py index 2eab8d13..40707105 100644 --- a/polyglot/python_worker/scripts/sdk_timers.py +++ b/polyglot/python_worker/scripts/sdk_timers.py @@ -72,7 +72,7 @@ async def pending(client: Client, runtime: str, scenario: str) -> None: while True: handle, execution, history = await observe(client, runtime, scenario) scheduled = events(history, "TimerScheduled") - if scheduled: + if scheduled and execution.status == "waiting": break if execution.status in ("failed", "cancelled", "terminated", "completed"): raise RuntimeError(f"{runtime} closed before scheduling its timer: {history!r}") diff --git a/polyglot/python_worker/tests/test_sdk_timers.py b/polyglot/python_worker/tests/test_sdk_timers.py index 04ef944a..c34946ac 100644 --- a/polyglot/python_worker/tests/test_sdk_timers.py +++ b/polyglot/python_worker/tests/test_sdk_timers.py @@ -5,6 +5,7 @@ import sys import types import unittest +from datetime import datetime, timedelta, timezone from pathlib import Path from unittest.mock import AsyncMock, patch @@ -47,6 +48,21 @@ async def verify(self, history, scenario="completion"): async def test_matching_timer_and_one_result_pass(self): await self.verify(self.history) + async def test_pending_waits_for_status_after_timer_history_arrives(self): + client = types.SimpleNamespace(start_workflow=AsyncMock()) + history = {"events": [event( + "TimerScheduled", timer_id="timer-1", delay_seconds=30, + fire_at=(datetime.now(timezone.utc) + timedelta(seconds=30)).isoformat(), + )]} + snapshots = [types.SimpleNamespace(workflow_id="timer-completion-rust", run_id="run-1", status=status) + for status in ("running", "waiting")] + observer = AsyncMock(side_effect=[(self.handle, execution, history) for execution in snapshots]) + with patch.dict("os.environ", {"DURABLE_WORKFLOW_TIMER_ID": "timer"}), \ + patch.object(timers, "observe", observer), patch.object(timers, "emit"), \ + patch.object(timers.asyncio, "sleep", AsyncMock()): + await timers.pending(client, "rust", "completion") + self.assertEqual(observer.await_count, 2) + async def test_duplicate_or_missing_durable_events_cannot_pass(self): for kind in ("TimerScheduled", "TimerFired", "WorkflowCompleted"): for duplicate in (False, True): From ca47f18d7b18635412bdf2daa5e8d282c7a0d86b Mon Sep 17 00:00:00 2001 From: Durable Workflow Date: Wed, 7 Oct 2026 22:56:30 +0000 Subject: [PATCH 4/5] Declare the Rust recovery API minimum version --- polyglot/rust_worker/Cargo.toml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/polyglot/rust_worker/Cargo.toml b/polyglot/rust_worker/Cargo.toml index de541e71..570620d2 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" +durable-workflow = "3.2.1" tokio = { version = "1.47.1", features = ["macros", "rt-multi-thread", "signal"] } [dev-dependencies] From 2544cd9c2ccb9828ae6eaf0e073182428eac1910 Mon Sep 17 00:00:00 2001 From: Durable Workflow Date: Wed, 7 Oct 2026 22:57:06 +0000 Subject: [PATCH 5/5] Cancel superseded pull request timer checks --- .github/workflows/sdk-timers.yml | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/.github/workflows/sdk-timers.yml b/.github/workflows/sdk-timers.yml index b0ea6b1b..94f4806b 100644 --- a/.github/workflows/sdk-timers.yml +++ b/.github/workflows/sdk-timers.yml @@ -10,6 +10,10 @@ on: permissions: contents: read +concurrency: + group: sdk-timers-${{ github.event_name }}-${{ github.event_name == 'pull_request' && github.ref || github.sha }} + cancel-in-progress: ${{ github.event_name == 'pull_request' }} + jobs: timers: name: SDK timers (PHP/Python/Rust)