Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
29 changes: 29 additions & 0 deletions .github/workflows/sdk-timers.yml
Original file line number Diff line number Diff line change
@@ -0,0 +1,29 @@
name: published SDK timers

on:
push:
branches: [main]
pull_request:
branches: [main]
workflow_dispatch:

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)
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
4 changes: 2 additions & 2 deletions playground/templates/rust/Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 1 addition & 1 deletion playground/templates/rust/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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"] }
2 changes: 1 addition & 1 deletion polyglot/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 |
Expand Down
12 changes: 10 additions & 2 deletions polyglot/docker-compose.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand All @@ -171,6 +172,7 @@ services:
condition: service_healthy

python-workflow-worker:
init: true
build:
context: ./python_workflow
args:
Expand All @@ -181,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
Expand All @@ -207,13 +210,15 @@ services:
condition: service_healthy

php-same-workflow-worker:
init: true
build: *php-sdk-worker-build
image: *php-sdk-worker-image
environment:
DURABLE_WORKFLOW_SERVER_URL: "http://server:8080"
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
Expand Down Expand Up @@ -300,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:
Expand All @@ -312,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
Expand Down Expand Up @@ -386,6 +393,7 @@ services:
retries: 30

smoke:
init: true
build:
context: ./python_worker
args:
Expand Down
18 changes: 16 additions & 2 deletions polyglot/php_worker/worker.php
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -22,6 +23,7 @@

const WORKFLOW_TYPES = [
'polyglot.php.greeter',
'polyglot.php.timer',
'polyglot.PolyglotWorkflow',
'polyglot.php-to-python.greeter',
'polyglot.php-to-python.type-roundtrip',
Expand Down Expand Up @@ -443,6 +445,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';
Expand Down Expand Up @@ -793,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);

Expand All @@ -805,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",
Expand Down
203 changes: 203 additions & 0 deletions polyglot/python_worker/scripts/sdk_timers.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,203 @@
"""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)


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")
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 and execution.status == "waiting":
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}")
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.")
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]
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"]):
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.")
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")])


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))
Loading
Loading