-
Notifications
You must be signed in to change notification settings - Fork 1
Expand file tree
/
Copy pathgreeter.py
More file actions
74 lines (58 loc) · 2.25 KB
/
Copy pathgreeter.py
File metadata and controls
74 lines (58 loc) · 2.25 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
"""End-to-end demo: start a Python-authored workflow with one activity.
Usage (from inside the server's docker network):
SERVER_URL=http://server:8080 \
WORKFLOW_TOKEN=dev-token-123 \
python -m examples.greeter
"""
from __future__ import annotations
import asyncio
import logging
import os
import sys
import uuid
from durable_workflow import Client, Worker, activity, workflow
@activity.defn(name="greet")
def greet(name: str) -> str:
return f"hello, {name}"
@workflow.defn(name="greeter")
class GreeterWorkflow:
def run(self, ctx, *_ignored):
greeting = yield ctx.schedule_activity("greet", ["world"])
return {"greeting": greeting, "length": len(greeting)}
async def main() -> int:
logging.basicConfig(level=logging.INFO, format="%(asctime)s %(name)s %(levelname)s %(message)s")
log = logging.getLogger("greeter")
url = os.environ.get("SERVER_URL", "http://server:8080")
token = os.environ.get("WORKFLOW_TOKEN", "dev-token-123")
task_queue = "python-workers"
wf_id = f"greet-{uuid.uuid4().hex[:6]}"
async with Client(url, token=token, namespace="default") as client:
handle = await client.start_workflow(
workflow_type="greeter",
task_queue=task_queue,
workflow_id=wf_id,
input=[],
)
log.info("started %s run=%s", handle.workflow_id, handle.run_id)
worker = Worker(
client,
task_queue=task_queue,
workflows=[GreeterWorkflow],
activities=[greet],
)
log.info("running worker until workflow completes...")
await worker.run_until(workflow_id=wf_id, timeout=30.0)
log.info("worker stopped, checking result...")
async with Client(url, token=token, namespace="default") as client2:
handle2 = client2.get_workflow_handle(handle.workflow_id, run_id=handle.run_id)
try:
result = await handle2.result(timeout=10.0, poll_interval=1.0)
log.info("RESULT: %r", result)
print(f"RESULT: {result!r}", flush=True)
return 0
except Exception as e:
log.error("FAILED: %s", e)
print(f"FAILED: {e}", flush=True)
return 1
if __name__ == "__main__":
sys.exit(asyncio.run(main()))