Skip to content
File

Blob: src/workerd/server/tests/python/workflow-entrypoint/workflow.py

python70 lines
1# Copyright (c) 2025 Cloudflare, Inc.
2# Licensed under the Apache 2.0 license found in the LICENSE file or at:
3# https://opensource.org/licenses/Apache-2.0
4 
5from workers import WorkflowEntrypoint
6 
7 
8class WorkflowEntrypointExample(WorkflowEntrypoint):
9 async def run(self, event, step):
10 async def await_step(fn):
11 try:
12 return await fn()
13 except TypeError as e:
14 print(f"Successfully caught {type(e).__name__}: {e}")
15 
16 @step.do("my_failing")
17 async def my_failing():
18 print("Executing my_failing")
19 raise TypeError("Intentional error in my_failing")
20 
21 @step.do("normal_step")
22 async def normal_step(ctx):
23 assert ctx["attempt"] == "1"
24 return "done"
25 
26 await normal_step()
27 
28 @step.do()
29 async def step_1():
30 print("Executing step 1")
31 return {"foo": "foo"}
32 
33 @step.do()
34 async def step_2():
35 print("Executing step 2")
36 return {"bar": "bar"}
37 
38 # DAG example with error handling
39 @step.do("step_3", concurrent=True)
40 async def step_3():
41 # this should never run because one of the dependencies will fail
42 pass
43 
44 await await_step(step_3)
45 
46 # `step_1` and `step_2` run serially
47 @step.do("step_4", concurrent=True)
48 async def step_4(step_1, ctx, step_2):
49 assert ctx["attempt"] == "1"
50 print("Executing step 4 (depends on step 1 and step 2)")
51 assert step_1["foo"] == "foo"
52 assert step_2["bar"] == "bar"
53 
54 await await_step(step_4)
55 
56 # tests step memoization - steps 1 and 2 are already resolved
57 @step.do("step_5", concurrent=False)
58 async def step_5(ctx, step_1=(), step_2=()):
59 print("Executing step 5 (depends on step 1 and step 2)")
60 assert step_1["foo"] == "foo"
61 assert step_2["bar"] == "bar"
62 
63 return step_1["foo"] + step_2["bar"]
64 
65 return await step_5()
66 
67 
68async def test(ctrl, env, ctx):
69 pass