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
37 changes: 37 additions & 0 deletions .github/workflows/cancellation-competitors.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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
67 changes: 67 additions & 0 deletions conformance/cancellation/temporal/README.md
Original file line number Diff line number Diff line change
@@ -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.
38 changes: 38 additions & 0 deletions conformance/cancellation/temporal/compose.yml
Original file line number Diff line number Diff line change
@@ -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
1 change: 1 addition & 0 deletions conformance/cancellation/temporal/requirements.txt
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
temporalio==1.34.0
179 changes: 179 additions & 0 deletions conformance/cancellation/temporal/scenario.py
Original file line number Diff line number Diff line change
@@ -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())
26 changes: 26 additions & 0 deletions conformance/cancellation/temporal/start-runtime.py
Original file line number Diff line number Diff line change
@@ -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'])
6 changes: 6 additions & 0 deletions conformance/cancellation/temporal/versions.json
Original file line number Diff line number Diff line change
@@ -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"
}
Loading
Loading