-
Notifications
You must be signed in to change notification settings - Fork 94
Expand file tree
/
Copy pathcancellation.go
More file actions
183 lines (144 loc) · 4.81 KB
/
Copy pathcancellation.go
File metadata and controls
183 lines (144 loc) · 4.81 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
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
package main
import (
"context"
"errors"
"log"
"os"
"os/signal"
"time"
"github.com/cschleiden/go-workflows/backend"
"github.com/cschleiden/go-workflows/client"
"github.com/cschleiden/go-workflows/samples"
"github.com/cschleiden/go-workflows/worker"
"github.com/cschleiden/go-workflows/workflow"
"github.com/google/uuid"
)
func main() {
ctx := context.Background()
ctx, cancel := context.WithCancel(ctx)
b := samples.GetBackend("cancellation", true)
// Run worker
go RunWorker(ctx, b)
// Start workflow via client
c := client.New(b)
startWorkflow(ctx, c)
c2 := make(chan os.Signal, 1)
signal.Notify(c2, os.Interrupt)
<-c2
cancel()
}
func startWorkflow(ctx context.Context, c *client.Client) {
wf, err := c.CreateWorkflowInstance(ctx, client.WorkflowInstanceOptions{
InstanceID: uuid.NewString(),
}, Workflow1, "Hello world")
if err != nil {
panic("could not start workflow")
}
time.Sleep(2 * time.Second)
if err := c.CancelWorkflowInstance(ctx, wf); err != nil {
panic("could not cancel workflow")
}
}
func RunWorker(ctx context.Context, mb backend.Backend) {
w := worker.New(mb, nil)
w.RegisterWorkflow(Workflow1)
w.RegisterWorkflow(Workflow2)
w.RegisterActivity(ActivityCancel)
w.RegisterActivity(ActivitySkip)
w.RegisterActivity(ActivitySuccess)
w.RegisterActivity(ActivityCleanup)
if err := w.Start(ctx); err != nil {
panic("could not start worker")
}
}
func Workflow1(ctx workflow.Context, msg string) (string, error) {
logger := workflow.Logger(ctx)
logger.Debug("Entering Workflow1", "msg", msg)
defer logger.Debug("Leaving Workflow1")
defer func() {
if errors.Is(ctx.Err(), workflow.Canceled) {
logger.Debug("Workflow1 was canceled")
logger.Debug("Do cleanup")
ctx := workflow.NewDisconnectedContext(ctx)
if _, err := workflow.ExecuteActivity[any](ctx, workflow.DefaultActivityOptions, ActivityCleanup).Get(ctx); err != nil {
panic("could not execute cleanup activity")
}
logger.Debug("Done with cleanup")
}
}()
logger.Debug("schedule ActivitySuccess")
if r0, err := workflow.ExecuteActivity[int](ctx, workflow.DefaultActivityOptions, ActivitySuccess, 1, 2).Get(ctx); err != nil {
logger.Debug("error getting activity success result", "err", err)
} else {
logger.Debug("ActivitySuccess result:", "r0", r0)
}
logger.Debug("Run SubWorkflow: Workflow2")
f := workflow.CreateSubWorkflowInstance[string](ctx, workflow.SubWorkflowOptions{
InstanceID: uuid.NewString(),
}, Workflow2, "hello sub")
workflow.Select(ctx,
workflow.Await(f, func(ctx workflow.Context, f workflow.Future[string]) {
rw, err := f.Get(ctx)
if err != nil {
logger.Debug("error getting workflow2 result", "err", err)
} else {
logger.Debug("Workflow2 result:", "rw", rw)
}
}),
)
logger.Debug("schedule ActivitySkip")
if r2, err := workflow.ExecuteActivity[int](ctx, workflow.DefaultActivityOptions, ActivitySkip, 1, 2).Get(ctx); err != nil {
logger.Debug("error getting activity skip result", "err", err)
} else {
logger.Debug("ActivitySkip result:", "r2", r2)
}
logger.Debug("Workflow finished")
return "result", nil
}
func Workflow2(ctx workflow.Context, msg string) (ret string, err error) {
logger := workflow.Logger(ctx)
logger.Debug("Entering Workflow2", "msg", msg)
defer logger.Debug("Leaving Workflow2")
defer func() {
if errors.Is(ctx.Err(), workflow.Canceled) {
logger.Debug("Workflow2 was canceled")
logger.Debug("Do cleanup")
ctx := workflow.NewDisconnectedContext(ctx)
if _, err := workflow.ExecuteActivity[any](ctx, workflow.DefaultActivityOptions, ActivityCleanup).Get(ctx); err != nil {
panic("could not execute cleanup activity")
}
logger.Debug("Done with cleanup")
ret = "cleanup result"
}
}()
logger.Debug("schedule ActivityCancel")
if r1, err := workflow.ExecuteActivity[int](ctx, workflow.DefaultActivityOptions, ActivityCancel, 1, 2).Get(ctx); err != nil {
logger.Debug("error getting activity cancel result", "err", err)
} else {
logger.Debug("ActivityCancel result:", "r1", r1)
}
return "some result", nil
}
func ActivitySuccess(ctx context.Context, a, b int) (int, error) {
log.Println("Entering ActivitySuccess")
defer log.Println("Leaving ActivitySuccess")
return a + b, nil
}
func ActivityCancel(ctx context.Context, a, b int) (int, error) {
log.Println("Entering ActivityCancel")
defer log.Println("Leaving ActivityCancel")
// Wait for 10s, this will cause the cancellation event to be fired while waiting here
time.Sleep(10 * time.Second)
return a + b, nil
}
func ActivitySkip(ctx context.Context, a, b int) (int, error) {
log.Println("Entering ActivitySkip")
defer log.Println("Leaving ActivitySkip")
return a + b, nil
}
func ActivityCleanup(ctx context.Context) error {
log.Println("Entering ActivityCleanup")
defer log.Println("Leaving ActivityCleanup")
log.Println("Do some cleanup")
return nil
}