From ec99509ade947bc04d1f1a6ea729d5c2d0710d86 Mon Sep 17 00:00:00 2001 From: Durable Workflow Date: Mon, 5 Oct 2026 21:54:13 +0000 Subject: [PATCH] Add published Temporal cancellation and cleanup comparison --- .../workflows/cancellation-competitors.yml | 37 ++++ conformance/cancellation/temporal/README.md | 67 +++++++ conformance/cancellation/temporal/compose.yml | 38 ++++ .../cancellation/temporal/requirements.txt | 1 + conformance/cancellation/temporal/scenario.py | 179 ++++++++++++++++++ .../cancellation/temporal/start-runtime.py | 26 +++ .../cancellation/temporal/versions.json | 6 + conformance/cancellation/temporal/worker.py | 128 +++++++++++++ 8 files changed, 482 insertions(+) create mode 100644 conformance/cancellation/temporal/README.md create mode 100644 conformance/cancellation/temporal/compose.yml create mode 100644 conformance/cancellation/temporal/requirements.txt create mode 100644 conformance/cancellation/temporal/scenario.py create mode 100644 conformance/cancellation/temporal/start-runtime.py create mode 100644 conformance/cancellation/temporal/versions.json create mode 100644 conformance/cancellation/temporal/worker.py diff --git a/.github/workflows/cancellation-competitors.yml b/.github/workflows/cancellation-competitors.yml index f9214f8..fa23229 100644 --- a/.github/workflows/cancellation-competitors.yml +++ b/.github/workflows/cancellation-competitors.yml @@ -52,3 +52,40 @@ jobs: path: conformance/cancellation/restate/evidence if-no-files-found: error retention-days: 90 + temporal: + name: Temporal published behavior + runs-on: ubuntu-latest + timeout-minutes: 10 + defaults: + run: + working-directory: conformance/cancellation/temporal + steps: + - uses: actions/checkout@d23441a48e516b6c34aea4fa41551a30e30af803 # v6 + with: + persist-credentials: false + - name: Run documented cancellation cases + run: | + set -euo pipefail + mkdir -p evidence + chmod 0777 evidence + git rev-parse HEAD > evidence/source.txt + sha256sum *.py requirements.txt versions.json compose.yml > evidence/source-files.txt + docker compose -p cancellation-temporal up -d + docker compose -p cancellation-temporal exec -T sdk sh -ec ' + python -m venv /tmp/venv + /tmp/venv/bin/pip install --no-cache-dir --report /evidence/packages.json -r requirements.txt + /tmp/venv/bin/python scenario.py + ' | tee evidence/run.log + - name: Collect observations and remove the fixture + if: always() + run: | + mkdir -p evidence + docker compose -p cancellation-temporal logs runtime > evidence/runtime.log + docker compose -p cancellation-temporal down -v --remove-orphans + - uses: actions/upload-artifact@043fb46d1a93c77aae656e7c1c64a875d1fc6a0a # v7 + if: always() + with: + name: temporal-cancellation-behavior + path: conformance/cancellation/temporal/evidence + if-no-files-found: error + retention-days: 90 diff --git a/conformance/cancellation/temporal/README.md b/conformance/cancellation/temporal/README.md new file mode 100644 index 0000000..dd35c40 --- /dev/null +++ b/conformance/cancellation/temporal/README.md @@ -0,0 +1,67 @@ +# Published Temporal cancellation behavior + +This fixture tests a Python parent, child and activity with published Temporal +artifacts. Exact versions and the verified Linux amd64 CLI archive are in +`versions.json`. The CLI runs its embedded development Server with SQLite. +The Python SDK is pinned in `requirements.txt` and both container images are +pinned by digest in `compose.yml`. + +The fixture uses explicit `WAIT_CANCELLATION_COMPLETED` policies for the child +and activity, `REQUEST_CANCEL` on parent close, shielded cleanup, and an explicit +workflow acknowledgement after an awaited operation returns. This last step +keeps a pending cancellation from becoming a normal successful workflow result +when an activity completes before observing the request. + +Both workflow-task timeouts are two seconds. The workflow cache retains its SDK +default. Heartbeating remote activities configure a two-second heartbeat timeout +and 100 ms worker heartbeat throttles. These are deliberate comparison settings, +not a production sizing recommendation. + +## Run + +Use Docker Compose on Linux amd64. No host language tooling or published ports +are needed. Downloading the checksum-verified CLI requires access to GitHub. + +```sh +mkdir -p evidence +chmod 0777 evidence +docker compose -p cancellation-temporal up -d +docker compose -p cancellation-temporal exec -T sdk sh -ec ' + python -m venv /tmp/venv + /tmp/venv/bin/pip install --no-cache-dir --report /evidence/packages.json -r requirements.txt + /tmp/venv/bin/python scenario.py +' +docker compose -p cancellation-temporal down -v --remove-orphans +``` + +The scenario waits for the Server and default namespace, then runs three +repetitions of six cases: + +| Case | Physical callback observation | Durable outcome | +| --- | --- | --- | +| Async remote, application heartbeats | Stops before its late effect | Parent and child CANCELED, cleanup completes | +| Async remote, no application heartbeats | Runs twelve seconds and performs the late effect | Parent and child CANCELED after explicit acknowledgement, cleanup completes | +| Async local, no application heartbeats | Stops before its late effect | Parent and child CANCELED, cleanup completes | +| Blocking local, no application heartbeats | Thread exits after its twelve-second sleep, without the late effect | Parent and child CANCELED, cleanup completes | +| Worker SIGKILL during committed root cleanup timer | Fresh worker resumes the original timer, one cleanup entry | Parent and child CANCELED, cleanup completes | +| Another cancellation during committed root cleanup timer | One original request and reason, one original timer | Parent and child CANCELED, cleanup completes | + +The last two cases also submit stale activity results while root cleanup is +running and require rejection. All remote cases submit another stale result +after workflow closure and require rejection. `results.jsonl` retains complete +parent/child histories, callback process IDs, physical exit markers, request +timestamps, cleanup outcomes and per-repetition elapsed times. A failed assertion +exits nonzero after preserving the observations. + +These are cancellation semantics checks. Latency is recorded to distinguish +physical callback exit from durable acknowledgement, not to rank throughput or +production configurations. The original local comparison stopped heartbeating +remote callbacks in about one second and async local callbacks in about two +seconds. Recovery and repeated cancellation both completed durable cleanup. + +Temporal documents the heartbeat requirement for remote activity cancellation +and the different local activity behavior in its +[Python cancellation guide](https://docs.temporal.io/develop/python/workflows/cancellation). +The fixture tests that documented distinction with explicit waiting policies. +It does not establish a general winner or describe another Temporal SDK's +callback execution model. diff --git a/conformance/cancellation/temporal/compose.yml b/conformance/cancellation/temporal/compose.yml new file mode 100644 index 0000000..a37ac92 --- /dev/null +++ b/conformance/cancellation/temporal/compose.yml @@ -0,0 +1,38 @@ +services: + runtime: + image: python@sha256:02108f5d322dd89f1c9e552442c25acb0543dfdbc455693a5599624f20d9155d + user: "1000:1000" + init: true + working_dir: /experiment + command: ["python", "start-runtime.py"] + environment: + USER: qualification + volumes: + - ./start-runtime.py:/experiment/start-runtime.py:ro + - ./versions.json:/experiment/versions.json:ro + tmpfs: + - /tmp:rw,exec,nosuid,nodev,uid=1000,gid=1000,mode=1777,size=512m + networks: + default: + aliases: [temporal] + mem_limit: 1536m + memswap_limit: 1536m + cpus: 2 + sdk: + image: python@sha256:02108f5d322dd89f1c9e552442c25acb0543dfdbc455693a5599624f20d9155d + user: "1000:1000" + init: true + working_dir: /experiment + command: ["sleep", "3600"] + environment: + PYTHONUNBUFFERED: "1" + PYTHONDONTWRITEBYTECODE: "1" + volumes: + - ./worker.py:/experiment/worker.py:ro + - ./scenario.py:/experiment/scenario.py:ro + - ./requirements.txt:/experiment/requirements.txt:ro + - ./versions.json:/experiment/versions.json:ro + - ./evidence:/evidence + mem_limit: 768m + memswap_limit: 768m + cpus: 2 diff --git a/conformance/cancellation/temporal/requirements.txt b/conformance/cancellation/temporal/requirements.txt new file mode 100644 index 0000000..64f0c83 --- /dev/null +++ b/conformance/cancellation/temporal/requirements.txt @@ -0,0 +1 @@ +temporalio==1.34.0 diff --git a/conformance/cancellation/temporal/scenario.py b/conformance/cancellation/temporal/scenario.py new file mode 100644 index 0000000..027659c --- /dev/null +++ b/conformance/cancellation/temporal/scenario.py @@ -0,0 +1,179 @@ +import asyncio +import json +import os +import signal +import subprocess +import sys +import time +import uuid +from datetime import timedelta +from pathlib import Path + +from google.protobuf.json_format import MessageToDict +from temporalio.api.enums.v1 import EventType +from temporalio.api.workflowservice.v1 import DescribeNamespaceRequest, GetClusterInfoRequest +from temporalio.client import Client + +from worker import Parent + +ROOT = Path('/evidence') +process = None +log = (ROOT / 'worker.log').open('a') + + +def start_worker(): + global process + process = subprocess.Popen([sys.executable, 'worker.py'], stdout=log, stderr=subprocess.STDOUT) + + +def events(token): + path = ROOT / 'events.jsonl' + return [row for line in path.read_text().splitlines() + if (row := json.loads(line))['token'] == token] if path.exists() else [] + + +async def wait_for(token, stage, seconds=20): + until = time.monotonic() + seconds + while time.monotonic() < until: + found = [row for row in events(token) if row['stage'] == stage] + if found: + return found[-1] + if process.poll() is not None: + raise RuntimeError('Worker exited unexpectedly') + await asyncio.sleep(0.02) + raise RuntimeError('Missing ' + stage) + + +async def case(client, profile, repetition, kill_cleanup=False, duplicate=False): + token = 'comparison-' + uuid.uuid4().hex + arg = dict(token=token, local=profile.startswith('local'), sync=profile == 'local_sync', + heartbeat=profile == 'remote_heartbeat', kill_cleanup=kill_cleanup, duplicate=duplicate) + handle = await client.start_workflow(Parent.run, arg, id=token, + task_queue='cancellation-comparison', task_timeout=timedelta(seconds=2)) + started = await wait_for(token, 'callback.started') + requested = time.monotonic() + await handle.cancel(reason='comparison cancellation') + killed_pid = None + replacement_pid = None + stale_during_cleanup = None + if kill_cleanup or duplicate: + await wait_for(token, 'root.cleanup.entered') + # Wait until the original cleanup timer is in canonical history. + until = time.monotonic() + 5 + while time.monotonic() < until: + history = await handle.fetch_history() + if any(event.event_type == EventType.EVENT_TYPE_TIMER_STARTED for event in history.events): + break + await asyncio.sleep(0.02) + else: + raise RuntimeError('Root cleanup timer did not commit') + try: + await client.get_async_activity_handle(task_token=bytes.fromhex(started['task_token'])).complete('stale-during-cleanup') + stale_during_cleanup = 'accepted' + except Exception as error: + stale_during_cleanup = type(error).__name__ + ': ' + str(error) + if duplicate: + await handle.cancel(reason='duplicate cancellation') + if kill_cleanup: + killed_pid = process.pid + os.kill(killed_pid, signal.SIGKILL) + process.wait(timeout=5) + start_worker() + replacement_pid = process.pid + try: + await asyncio.wait_for(handle.result(), timeout=35) + result = 'returned' + except Exception as error: + result = type(error).__name__ + ': ' + str(error) + stopped = await wait_for(token, 'callback.exited', 15) + child = client.get_workflow_handle(token + '-child') + parent_history = await handle.fetch_history() + child_history = await child.fetch_history() + rows = events(token) + cleanup = [event for event in rows if event['stage'] == 'root.cleanup.completed'] + stale = None + if not arg['local']: + try: + await client.get_async_activity_handle(task_token=bytes.fromhex(started['task_token'])).complete('stale-result') + stale = 'accepted' + except Exception as error: + stale = type(error).__name__ + ': ' + str(error) + row = dict(token=token, profile=profile, repetition=repetition, requested=requested, + callback_exit_elapsed=stopped['monotonic'] - requested, + cleanup_elapsed=cleanup[-1]['monotonic'] - requested if cleanup else None, + late_effect=any(event['stage'] == 'callback.late_effect' for event in rows), + parent_status=(await handle.describe()).status.name, + child_status=(await child.describe()).status.name, + result=result, stale_completion=stale, duplicate=duplicate, + killed_pid=killed_pid, replacement_pid=replacement_pid, + stale_during_cleanup=stale_during_cleanup, events=rows, + parent_history=[MessageToDict(event) for event in parent_history.events], + child_history=[MessageToDict(event) for event in child_history.events]) + with (ROOT / 'results.jsonl').open('a') as output: + output.write(json.dumps(row) + '\n') + print(json.dumps({key: row[key] for key in ('profile', 'repetition', 'parent_status', 'child_status', + 'callback_exit_elapsed', 'cleanup_elapsed', 'late_effect', 'duplicate', 'killed_pid', 'replacement_pid')}), flush=True) + + # Require the durable outcomes and physical observations, independently. + assert row['parent_status'] == row['child_status'] == 'CANCELED', row + assert cleanup and row['cleanup_elapsed'] < 30, row + assert sum(event['stage'] == 'root.cleanup.entered' for event in rows) == 1, row + assert sum(event['stage'] == 'root.cleanup.completed' for event in rows) == 1, row + assert sum(event['stage'] == 'child.cleanup.entered' for event in rows) == 1, row + assert sum(event['stage'] == 'child.cleanup.completed' for event in rows) == 1, row + original = [event for event in row['parent_history'] + if event['eventType'] == 'EVENT_TYPE_WORKFLOW_EXECUTION_CANCEL_REQUESTED'] + assert len(original) == 1, original + assert original[0]['workflowExecutionCancelRequestedEventAttributes']['cause'] == 'comparison cancellation', original + if not arg['local']: + assert stale is not None and stale != 'accepted', row + if profile == 'remote_no_heartbeat': + assert row['late_effect'] and row['callback_exit_elapsed'] > 10, row + elif profile == 'local_sync': + assert not row['late_effect'] and row['callback_exit_elapsed'] > 10, row + else: + assert not row['late_effect'], row + if kill_cleanup: + assert killed_pid != replacement_pid, row + if kill_cleanup or duplicate: + assert stale_during_cleanup is not None and stale_during_cleanup != 'accepted', row + for event_type in ('EVENT_TYPE_TIMER_STARTED', 'EVENT_TYPE_TIMER_FIRED'): + assert sum(event['eventType'] == event_type for event in row['parent_history']) == 1, row + + +async def main(): + until = time.monotonic() + 20 + while True: + try: + client = await Client.connect('temporal:7233') + cluster = await client.workflow_service.get_cluster_info(GetClusterInfoRequest()) + await client.workflow_service.describe_namespace(DescribeNamespaceRequest(namespace='default')) + break + except Exception: + if time.monotonic() >= until: + raise + await asyncio.sleep(0.1) + (ROOT / 'cluster.json').write_text(json.dumps(MessageToDict(cluster), indent=2)) + versions = json.loads(Path('versions.json').read_text()) + assert cluster.server_version == versions['server'], cluster + start_worker() + try: + for profile in ('remote_heartbeat', 'remote_no_heartbeat', 'local_async', 'local_sync'): + for repetition in range(1, 4): + await case(client, profile, repetition) + for repetition in range(1, 4): + await case(client, 'remote_heartbeat', repetition, kill_cleanup=True) + for repetition in range(1, 4): + await case(client, 'remote_heartbeat', repetition, duplicate=True) + finally: + if process is not None and process.poll() is None: + process.terminate() + try: + process.wait(timeout=5) + except subprocess.TimeoutExpired: + process.kill() + process.wait(timeout=5) + log.close() + + +asyncio.run(main()) diff --git a/conformance/cancellation/temporal/start-runtime.py b/conformance/cancellation/temporal/start-runtime.py new file mode 100644 index 0000000..b5e3238 --- /dev/null +++ b/conformance/cancellation/temporal/start-runtime.py @@ -0,0 +1,26 @@ +"""Run the verified published CLI's isolated development server.""" + +import hashlib +import io +import json +import os +import tarfile +import urllib.request +from pathlib import Path + +versions = json.loads(Path('versions.json').read_text()) +with urllib.request.urlopen(versions['archive'], timeout=60) as response: + archive = response.read() +if hashlib.sha256(archive).hexdigest() != versions['sha256']: + raise RuntimeError('Published Temporal CLI archive checksum mismatch') +with tarfile.open(fileobj=io.BytesIO(archive), mode='r:gz') as bundle: + binary = bundle.getmember('temporal') + if not binary.isfile(): + raise RuntimeError('Missing Temporal CLI binary') + source = bundle.extractfile(binary) + if source is None: + raise RuntimeError('Unreadable Temporal CLI binary') + Path('/tmp/temporal').write_bytes(source.read()) +Path('/tmp/temporal').chmod(0o755) +os.execv('/tmp/temporal', ['temporal', 'server', 'start-dev', '--ip', '0.0.0.0', + '--db-filename', '/tmp/temporal.sqlite', '--log-level', 'error']) diff --git a/conformance/cancellation/temporal/versions.json b/conformance/cancellation/temporal/versions.json new file mode 100644 index 0000000..fb517bd --- /dev/null +++ b/conformance/cancellation/temporal/versions.json @@ -0,0 +1,6 @@ +{ + "cli": "1.9.1", + "server": "1.32.0", + "archive": "https://github.com/temporalio/cli/releases/download/v1.9.1/temporal_cli_1.9.1_linux_amd64.tar.gz", + "sha256": "09a0326a51db84d02735e53542b9ebd8c4758daf47482a9ab0abce15844e60d5" +} diff --git a/conformance/cancellation/temporal/worker.py b/conformance/cancellation/temporal/worker.py new file mode 100644 index 0000000..206813b --- /dev/null +++ b/conformance/cancellation/temporal/worker.py @@ -0,0 +1,128 @@ +import asyncio +import json +import os +import time +from concurrent.futures import ThreadPoolExecutor +from datetime import timedelta +from pathlib import Path + +from temporalio import activity, workflow +from temporalio.client import Client +from temporalio.common import RetryPolicy +from temporalio.exceptions import is_cancelled_exception +from temporalio.worker import Worker + +ROOT = Path('/evidence') + + +def observe(token, stage, **extra): + row = dict(token=token, stage=stage, monotonic=time.monotonic(), pid=os.getpid(), **extra) + with (ROOT / 'events.jsonl').open('a') as output: + output.write(json.dumps(row) + '\n') + + +@activity.defn +def mark(arg: dict) -> None: + observe(arg['token'], arg['stage'], reason=arg.get('reason')) + + +def started(arg): + info = activity.info() + observe(arg['token'], 'callback.started', local=info.is_local, + task_token=info.task_token.hex(), attempt=info.attempt) + + +@activity.defn +async def async_work(arg: dict) -> str: + started(arg) + until = time.monotonic() + 12 + try: + while time.monotonic() < until: + if arg['heartbeat']: + activity.heartbeat('working') + await asyncio.sleep(0.02) + observe(arg['token'], 'callback.late_effect') + return 'late-result' + finally: + observe(arg['token'], 'callback.exited') + + +@activity.defn +def sync_work(arg: dict) -> str: + started(arg) + try: + time.sleep(12) + observe(arg['token'], 'callback.late_effect') + return 'late-result' + finally: + observe(arg['token'], 'callback.exited') + + +async def cleanup(arg, stage): + await workflow.execute_local_activity(mark, dict(token=arg['token'], stage=stage + '.cleanup.entered', + reason=workflow.cancellation_reason()), start_to_close_timeout=timedelta(seconds=5), + retry_policy=RetryPolicy(maximum_attempts=1)) + if stage == 'root' and (arg['kill_cleanup'] or arg['duplicate']): + await asyncio.sleep(3) + await workflow.execute_local_activity(mark, dict(token=arg['token'], stage=stage + '.cleanup.completed'), + start_to_close_timeout=timedelta(seconds=5), retry_policy=RetryPolicy(maximum_attempts=1)) + + +@workflow.defn +class Child: + @workflow.run + async def run(self, arg: dict) -> str: + try: + callback = sync_work if arg['sync'] else async_work + options = dict(start_to_close_timeout=timedelta(seconds=30), + cancellation_type=workflow.ActivityCancellationType.WAIT_CANCELLATION_COMPLETED, + retry_policy=RetryPolicy(maximum_attempts=1)) + if arg['local']: + result = await workflow.execute_local_activity(callback, arg, **options) + else: + if arg['heartbeat']: + options['heartbeat_timeout'] = timedelta(seconds=2) + result = await workflow.execute_activity(callback, arg, **options) + if workflow.cancellation_reason() is not None: + raise asyncio.CancelledError() + return result + except BaseException as error: + if not is_cancelled_exception(error): + raise + await asyncio.shield(cleanup(arg, 'child')) + raise + + +@workflow.defn +class Parent: + @workflow.run + async def run(self, arg: dict) -> str: + try: + result = await workflow.execute_child_workflow(Child.run, arg, + id=arg['token'] + '-child', + cancellation_type=workflow.ChildWorkflowCancellationType.WAIT_CANCELLATION_COMPLETED, + parent_close_policy=workflow.ParentClosePolicy.REQUEST_CANCEL, + task_timeout=timedelta(seconds=2)) + if workflow.cancellation_reason() is not None: + raise asyncio.CancelledError() + return result + except BaseException as error: + if not is_cancelled_exception(error): + raise + await asyncio.shield(cleanup(arg, 'root')) + raise + + +async def main(): + client = await Client.connect('temporal:7233') + with ThreadPoolExecutor(max_workers=4) as executor: + async with Worker(client, task_queue='cancellation-comparison', workflows=[Parent, Child], + activities=[mark, async_work, sync_work], activity_executor=executor, + sticky_queue_schedule_to_start_timeout=timedelta(seconds=2), + max_heartbeat_throttle_interval=timedelta(milliseconds=100), + default_heartbeat_throttle_interval=timedelta(milliseconds=100)): + await asyncio.Event().wait() + + +if __name__ == '__main__': + asyncio.run(main())