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
54 changes: 54 additions & 0 deletions .github/workflows/cancellation-competitors.yml
Original file line number Diff line number Diff line change
@@ -0,0 +1,54 @@
name: Published cancellation behavior

on:
pull_request:
paths:
- conformance/cancellation/**
- .github/workflows/cancellation-competitors.yml
workflow_dispatch:

permissions:
contents: read

concurrency:
group: cancellation-competitors-${{ github.ref }}
cancel-in-progress: true

jobs:
restate:
name: Restate published behavior
runs-on: ubuntu-latest
timeout-minutes: 10
defaults:
run:
working-directory: conformance/cancellation/restate
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
docker compose -p cancellation-restate up -d --wait
docker compose -p cancellation-restate 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
/tmp/venv/bin/python scenario.py duplicate
' | tee evidence/run.log
- name: Collect observations and remove the fixture
if: always()
run: |
mkdir -p evidence
docker compose -p cancellation-restate logs runtime > evidence/runtime.log
docker compose -p cancellation-restate down -v --remove-orphans
- uses: actions/upload-artifact@043fb46d1a93c77aae656e7c1c64a875d1fc6a0a # v7
if: always()
with:
name: restate-cancellation-behavior
path: conformance/cancellation/restate/evidence
if-no-files-found: error
retention-days: 90
1 change: 1 addition & 0 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@
__pycache__/
*.py[cod]
node_modules/
/conformance/cancellation/*/evidence/

/authoritative-verification.json
/candidate.json
Expand Down
39 changes: 39 additions & 0 deletions conformance/cancellation/restate/README.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,39 @@
# Published Restate cancellation experiment

This fixture supports [cancellation issue 136](https://github.com/durable-workflow/.github/issues/136). It exercises Restate Server 1.7.13 and Python SDK 1.0.5 through the documented cooperative cancellation API. It measures behavior, not runtime throughput.

The parent and child are durable Restate Workflows. The child calls a leaf service. Each handler catches `TerminalError`, performs durable compensation and rethrows, following [Restate's cancellation guidance](https://docs.restate.dev/services/invocation/managing-invocations).

## Cases

Three repetitions each:

- Cancel a durable wait and inspect all three invocation outcomes.
- Cancel an asynchronous `ctx.run_typed` callback, with no application heartbeat. Observe its actual exit and a possible late side effect.
- Cancel a blocking `ctx.run_typed` callback. Wait for actual callback exit even if the durable invocation already completed.
- SIGKILL the SDK endpoint after root cleanup starts. Start a fresh process and inspect journal replay and cleanup completion.
- Repeat cancellation while root cleanup awaits a durable timer. Inspect whether cleanup completes.

The callback duration is 12 seconds. Recovery and duplicate cases insert a three-second durable wait between two cleanup actions. Application markers and durable `sys_invocation` records are retained separately. A cancellation result cannot stand in for callback exit or completed compensation.

## Run

Use Docker Compose and a host checkout owned by UID/GID 1000. No host ports are published. The runtime is isolated and its database is disposable.

```sh
cd conformance/cancellation/restate
mkdir -p evidence
docker compose -p cancellation-restate up -d --wait
docker compose -p cancellation-restate 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
/tmp/venv/bin/python scenario.py duplicate
'
docker compose -p cancellation-restate logs runtime > evidence/runtime.log
docker compose -p cancellation-restate down -v --remove-orphans
```

Keep measured failures and successes. Setup and import failures do not establish cancellation behavior. Read each case's application observations and durable outcomes before drawing a conclusion. The fixture records observed behavior without requiring another product to satisfy Durable Workflow's contract.

The pinned images and installed-package report bind each run to published artifacts. This comparison does not qualify Durable Workflow's unpublished cancellation candidate or replace its required published mixed-language cascade.
47 changes: 47 additions & 0 deletions conformance/cancellation/restate/compose.yml
Original file line number Diff line number Diff line change
@@ -0,0 +1,47 @@
services:
runtime:
image: docker.restate.dev/restatedev/restate@sha256:65abc8016d408b8b424f5e492a5da37c24573207e217a3d43242dd8b3fd8dc61
user: "1000:1000"
init: true
command: ["--no-logo"]
environment:
RESTATE_BASE_DIR: /restate-data
RESTATE_ADVERTISED_HOST: restate
RESTATE_NODE_NAME: cancellation-comparison
tmpfs:
- /restate-data:uid=1000,gid=1000,mode=0755
networks:
default:
aliases: [restate]
mem_limit: 1024m
memswap_limit: 1024m
cpus: 2
healthcheck:
test: ["CMD", "curl", "--fail", "--silent", "http://localhost:9070/health"]
interval: 2s
timeout: 5s
retries: 30
sdk:
image: python@sha256:02108f5d322dd89f1c9e552442c25acb0543dfdbc455693a5599624f20d9155d
user: "1000:1000"
init: true
working_dir: /experiment
command: ["sleep", "3600"]
environment:
PYTHONUNBUFFERED: "1"
volumes:
- ./service.py:/experiment/service.py:ro
- ./scenario.py:/experiment/scenario.py:ro
- ./requirements.txt:/experiment/requirements.txt:ro
- ./evidence:/evidence
tmpfs:
- /experiment:uid=1000,gid=1000,mode=0755
networks:
default:
aliases: [sdk]
mem_limit: 512m
memswap_limit: 512m
cpus: 1
depends_on:
runtime:
condition: service_healthy
3 changes: 3 additions & 0 deletions conformance/cancellation/restate/requirements.txt
Original file line number Diff line number Diff line change
@@ -0,0 +1,3 @@
restate-sdk==1.0.5
Hypercorn==0.18.0
httpx==0.28.1
157 changes: 157 additions & 0 deletions conformance/cancellation/restate/scenario.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,157 @@
import json
import os
import signal
import subprocess
import sys
import time
import uuid
from pathlib import Path

import httpx

ROOT = Path('/evidence')
client = httpx.Client(timeout=5)
process = None
service_log = (ROOT / 'service.log').open('a')


def start():
global process
process = subprocess.Popen([sys.executable, 'service.py'], stdout=service_log,
stderr=subprocess.STDOUT)
until = time.monotonic() + 15
while time.monotonic() < until:
try:
response = client.get('http://sdk:9080/health')
if response.status_code in (200, 404):
return
except httpx.HTTPError:
pass
time.sleep(0.1)
raise RuntimeError('SDK endpoint did not become ready')


def events(token):
path = ROOT / 'events.jsonl'
if not path.exists():
return []
return [row for line in path.read_text().splitlines()
if (row := json.loads(line))['token'] == token]


def wait_for(token, stage, timeout):
until = time.monotonic() + timeout
while time.monotonic() < until:
found = [row for row in events(token) if row['stage'] == stage]
if found:
return found[-1]
time.sleep(0.05)
raise RuntimeError('Missing ' + stage + ' for ' + token)


def query(sql):
response = client.post('http://restate:9070/query', json={'query': sql},
headers={'Accept': 'application/json'})
response.raise_for_status()
return response.json()


def case(mode, repetition, kill_cleanup=False, duration=12, duplicate_request=False):
token = 'comparison-' + uuid.uuid4().hex
arg = dict(token=token, mode=mode, duration=duration, kill_cleanup=kill_cleanup,
duplicate_during_cleanup=duplicate_request)
response = client.post(f'http://restate:8080/CancellationRoot/{token}/run/send', json=arg)
response.raise_for_status()
submitted = response.json()
invocation = submitted['invocationId']
ready = 'leaf.wait.started' if mode == 'durable_wait' else 'leaf.callback.started'
wait_for(token, ready, 15)
requested = time.monotonic()
accepted = client.patch(f'http://restate:9070/invocations/{invocation}/cancel')
accepted.raise_for_status()
duplicate = None
if duplicate_request:
wait_for(token, 'root.cleanup.entered', 20)
duplicate = client.patch(f'http://restate:9070/invocations/{invocation}/cancel')
duplicate.raise_for_status()
killed_pid = None
replacement_pid = None
if kill_cleanup:
wait_for(token, 'root.cleanup.entered', 20)
killed_pid = process.pid
os.kill(killed_pid, signal.SIGKILL)
process.wait(timeout=5)
start()
replacement_pid = process.pid
until = time.monotonic() + max(35, duration + 10)
while time.monotonic() < until:
observed = events(token)
invocation_ids = sorted({event['invocation_id'] for event in observed if 'invocation_id' in event})
call_graph = query("SELECT * FROM sys_invocation WHERE id IN (" +
','.join("'" + value + "'" for value in invocation_ids) + ")")
rows = call_graph['rows']
if len(rows) == 3 and all(item['status'] == 'completed' for item in rows):
break
time.sleep(0.1)
observed = events(token)
completed = [event for event in observed if event['stage'] == 'root.cleanup.completed']
callback_exit = None
settled = time.monotonic() - requested
if mode in ('async_run', 'sync_run'):
callback_exit = wait_for(token, 'leaf.callback.exited', duration + 10)
observed = events(token)
row = dict(mode=mode, repetition=repetition, token=token, request_monotonic=requested,
submitted=submitted, cancel_status=accepted.status_code,
cancel_body=accepted.text, duplicate_status=duplicate.status_code if duplicate is not None else None,
duplicate_body=duplicate.text if duplicate is not None else None,
cleanup_elapsed=completed[-1]['monotonic'] - requested if completed else None,
cascade_settled_elapsed=settled,
cascade_completed=len(rows) == 3 and all(item['status'] == 'completed' for item in rows),
callback_stop_observed=callback_exit is not None,
callback_exit_elapsed=callback_exit['monotonic'] - requested if callback_exit else None,
late_effect_observed=any(event['stage'] == 'leaf.callback.late_effect' for event in observed),
duplicate_during_cleanup=duplicate_request,
killed_pid=killed_pid, replacement_pid=replacement_pid, events=observed)
# Query the durable invocation projection, not just application log entries.
time.sleep(0.2)
row['durable_invocations'] = query("SELECT * FROM sys_invocation WHERE id = '" + invocation + "'")
invocation_ids = sorted({event['invocation_id'] for event in row['events'] if 'invocation_id' in event})
row['call_graph'] = query("SELECT * FROM sys_invocation WHERE id IN (" +
','.join("'" + value + "'" for value in invocation_ids) + ")")
with (ROOT / 'results.jsonl').open('a') as output:
output.write(json.dumps(row) + '\n')
print(json.dumps({key: row[key] for key in (
'mode', 'repetition', 'cascade_completed', 'cleanup_elapsed',
'callback_exit_elapsed', 'late_effect_observed', 'duplicate_during_cleanup',
'killed_pid', 'replacement_pid'
)}), flush=True)


try:
start()
response = client.post('http://restate:9070/deployments', json={'uri': 'http://sdk:9080'})
response.raise_for_status()
(ROOT / 'deployment.json').write_text(json.dumps(response.json(), indent=2))
if len(sys.argv) > 1 and sys.argv[1] == 'duplicate':
for repetition in range(1, 4):
case('durable_wait', repetition, duplicate_request=True, duration=12)
elif len(sys.argv) > 1 and sys.argv[1] == 'callback':
for mode in ('async_run', 'sync_run'):
for repetition in range(1, 4):
case(mode, repetition)
else:
for mode in ('durable_wait', 'async_run', 'sync_run'):
for repetition in range(1, 4):
case(mode, repetition)
for repetition in range(1, 4):
case('durable_wait', repetition, kill_cleanup=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)
service_log.close()
client.close()
Loading
Loading