go.temporal.io/server/tests/stickytq_test.go
414 LOC · 0 covered · 414 uncovered · 0 ranges · 0 concepts · 0 introducers · 0 tests
1
package tests
2
3
import (
4
"errors"
5
"testing"
6
"time"
7
8
"github.com/google/uuid"
9
commandpb "go.temporal.io/api/command/v1"
10
commonpb "go.temporal.io/api/common/v1"
11
enumspb "go.temporal.io/api/enums/v1"
12
taskqueuepb "go.temporal.io/api/taskqueue/v1"
13
"go.temporal.io/api/workflowservice/v1"
14
"go.temporal.io/server/common/log/tag"
15
"go.temporal.io/server/common/payloads"
16
"go.temporal.io/server/common/testing/parallelsuite"
17
"go.temporal.io/server/tests/testcore"
18
"google.golang.org/protobuf/types/known/durationpb"
19
)
20
21
type StickyTqTestSuite struct {
22
parallelsuite.Suite[*StickyTqTestSuite]
23
}
24
25
func TestStickyTqTestSuite(t *testing.T) {
26
parallelsuite.Run(t, &StickyTqTestSuite{})
27
}
28
29
func (s *StickyTqTestSuite) TestStickyTimeoutNonTransientWorkflowTask() {
30
env := testcore.NewEnv(s.T())
31
id := "functional-sticky-timeout-non-transient-workflow-task"
32
wt := "functional-sticky-timeout-non-transient-command-type"
33
tl := "functional-sticky-timeout-non-transient-workflow-taskqueue"
34
stl := "functional-sticky-timeout-non-transient-workflow-taskqueue-sticky"
35
identity := "worker1"
36
37
stickyTaskQueue := &taskqueuepb.TaskQueue{Name: stl, Kind: enumspb.TASK_QUEUE_KIND_STICKY, NormalName: tl}
38
stickyScheduleToStartTimeout := 2 * time.Second
39
40
// Start workflow execution
41
request := &workflowservice.StartWorkflowExecutionRequest{
42
RequestId: uuid.NewString(),
43
Namespace: env.Namespace().String(),
44
WorkflowId: id,
45
WorkflowType: &commonpb.WorkflowType{Name: wt},
46
TaskQueue: &taskqueuepb.TaskQueue{Name: tl, Kind: enumspb.TASK_QUEUE_KIND_NORMAL},
47
Input: nil,
48
WorkflowRunTimeout: durationpb.New(100 * time.Second),
49
WorkflowTaskTimeout: durationpb.New(1 * time.Second),
50
Identity: identity,
51
}
52
53
we, err0 := env.FrontendClient().StartWorkflowExecution(s.Context(), request)
54
s.NoError(err0)
55
56
env.Logger.Info("StartWorkflowExecution", tag.WorkflowRunID(we.RunId))
57
workflowExecution := &commonpb.WorkflowExecution{
58
WorkflowId: id,
59
RunId: we.RunId,
60
}
61
62
// workflow logic
63
localActivityDone := false
64
failureCount := 5
65
wtHandler := func(task *workflowservice.PollWorkflowTaskQueueResponse) ([]*commandpb.Command, error) {
66
if !localActivityDone {
67
localActivityDone = true
68
69
return []*commandpb.Command{{
70
CommandType: enumspb.COMMAND_TYPE_RECORD_MARKER,
71
Attributes: &commandpb.Command_RecordMarkerCommandAttributes{RecordMarkerCommandAttributes: &commandpb.RecordMarkerCommandAttributes{
72
MarkerName: "local activity marker",
73
Details: map[string]*commonpb.Payloads{
74
"data": payloads.EncodeString("local activity marker"),
75
"result": payloads.EncodeString("local activity result"),
76
}}},
77
}}, nil
78
}
79
80
if failureCount > 0 {
81
// send a signal on third failure to be buffered, forcing a non-transient workflow task when buffer is flushed
82
/*
83
if failureCount == 3 {
84
err := s.FrontendClient().SignalWorkflowExecution(NewContext(), &workflowservice.SignalWorkflowExecutionRequest{
85
Namespace: s.Namespace(),
86
WorkflowExecution: workflowExecution,
87
SignalName: "signalB",
88
Input: codec.EncodeString("signal input"),
89
Identity: identity,
90
RequestId: uuid.NewString(),
91
})
92
s.NoError(err)
93
}
94
*/
95
failureCount--
96
return nil, errors.New("non deterministic error")
97
}
98
99
return []*commandpb.Command{{
100
CommandType: enumspb.COMMAND_TYPE_COMPLETE_WORKFLOW_EXECUTION,
101
Attributes: &commandpb.Command_CompleteWorkflowExecutionCommandAttributes{CompleteWorkflowExecutionCommandAttributes: &commandpb.CompleteWorkflowExecutionCommandAttributes{
102
Result: payloads.EncodeString("Done"),
103
}},
104
}}, nil
105
}
106
107
poller := &testcore.TaskPoller{
108
Client: env.FrontendClient(),
109
Namespace: env.Namespace().String(),
110
TaskQueue: &taskqueuepb.TaskQueue{Name: tl, Kind: enumspb.TASK_QUEUE_KIND_NORMAL},
111
Identity: identity,
112
WorkflowTaskHandler: wtHandler,
113
Logger: env.Logger,
114
T: s.T(),
115
StickyTaskQueue: stickyTaskQueue,
116
StickyScheduleToStartTimeout: stickyScheduleToStartTimeout,
117
}
118
119
_, err := poller.PollAndProcessWorkflowTask(testcore.WithRespondSticky)
120
env.Logger.Info("PollAndProcessWorkflowTask", tag.Error(err))
121
s.NoError(err)
122
123
_, err = env.FrontendClient().SignalWorkflowExecution(s.Context(), &workflowservice.SignalWorkflowExecutionRequest{
124
Namespace: env.Namespace().String(),
125
WorkflowExecution: workflowExecution,
126
SignalName: "signalA",
127
Input: payloads.EncodeString("signal input"),
128
Identity: identity,
129
RequestId: uuid.NewString(),
130
})
131
s.NoError(err)
132
133
// Wait for workflow task timeout
134
stickyTimeout := false
135
WaitForStickyTimeoutLoop:
136
for range 10 {
137
events := env.GetHistory(env.Namespace().String(), workflowExecution)
138
for _, event := range events {
139
if event.GetEventType() == enumspb.EVENT_TYPE_WORKFLOW_TASK_TIMED_OUT {
140
s.EqualHistoryEvents(`
141
1 WorkflowExecutionStarted
142
2 WorkflowTaskScheduled
143
3 WorkflowTaskStarted
144
4 WorkflowTaskCompleted
145
5 MarkerRecorded
146
6 WorkflowExecutionSignaled
147
7 WorkflowTaskScheduled
148
8 WorkflowTaskTimedOut {"TimeoutType":2} // enumspb.TIMEOUT_TYPE_SCHEDULE_TO_START
149
9 WorkflowTaskScheduled`, events)
150
stickyTimeout = true
151
break WaitForStickyTimeoutLoop
152
}
153
}
154
time.Sleep(time.Second) //nolint:forbidigo
155
}
156
s.True(stickyTimeout, "Workflow task not timed out")
157
158
for i := 1; i <= 3; i++ {
159
_, err = poller.PollAndProcessWorkflowTask(testcore.WithDumpHistory, testcore.WithRespondSticky, testcore.WithExpectedAttemptCount(i))
160
env.Logger.Info("PollAndProcessWorkflowTask", tag.Error(err))
161
s.NoError(err)
162
}
163
164
_, err = env.FrontendClient().SignalWorkflowExecution(s.Context(), &workflowservice.SignalWorkflowExecutionRequest{
165
Namespace: env.Namespace().String(),
166
WorkflowExecution: workflowExecution,
167
SignalName: "signalB",
168
Input: payloads.EncodeString("signal input"),
169
Identity: identity,
170
RequestId: uuid.NewString(),
171
})
172
s.NoError(err)
173
174
for i := 1; i <= 2; i++ {
175
_, err = poller.PollAndProcessWorkflowTask(testcore.WithDumpHistory, testcore.WithRespondSticky, testcore.WithExpectedAttemptCount(i))
176
env.Logger.Info("PollAndProcessWorkflowTask", tag.Error(err))
177
s.NoError(err)
178
}
179
180
events := env.GetHistory(env.Namespace().String(), workflowExecution)
181
s.EqualHistoryEvents(`
182
1 WorkflowExecutionStarted
183
2 WorkflowTaskScheduled
184
3 WorkflowTaskStarted
185
4 WorkflowTaskCompleted
186
5 MarkerRecorded
187
6 WorkflowExecutionSignaled
188
7 WorkflowTaskScheduled
189
8 WorkflowTaskTimedOut
190
9 WorkflowTaskScheduled
191
10 WorkflowTaskStarted
192
11 WorkflowTaskFailed
193
12 WorkflowExecutionSignaled
194
13 WorkflowTaskScheduled
195
14 WorkflowTaskStarted
196
15 WorkflowTaskFailed
197
16 WorkflowTaskScheduled`, events)
198
199
// Complete workflow execution
200
_, err = poller.PollAndProcessWorkflowTask(testcore.WithDumpHistory, testcore.WithRespondSticky, testcore.WithExpectedAttemptCount(3))
201
s.NoError(err)
202
203
events = env.GetHistory(env.Namespace().String(), workflowExecution)
204
s.EqualHistoryEvents(`
205
1 WorkflowExecutionStarted
206
2 WorkflowTaskScheduled
207
3 WorkflowTaskStarted
208
4 WorkflowTaskCompleted
209
5 MarkerRecorded
210
6 WorkflowExecutionSignaled
211
7 WorkflowTaskScheduled
212
8 WorkflowTaskTimedOut
213
9 WorkflowTaskScheduled
214
10 WorkflowTaskStarted
215
11 WorkflowTaskFailed // Two WFTs have failed
216
12 WorkflowExecutionSignaled
217
13 WorkflowTaskScheduled
218
14 WorkflowTaskStarted
219
15 WorkflowTaskFailed // Two WFTs have failed
220
16 WorkflowTaskScheduled
221
17 WorkflowTaskStarted
222
18 WorkflowTaskCompleted
223
19 WorkflowExecutionCompleted // Workflow has completed`, events)
224
}
225
226
func (s *StickyTqTestSuite) TestStickyTaskqueueResetThenTimeout() {
227
env := testcore.NewEnv(s.T())
228
id := "functional-reset-sticky-fire-schedule-to-start-timeout"
229
wt := "functional-reset-sticky-fire-schedule-to-start-timeout-type"
230
tl := "functional-reset-sticky-fire-schedule-to-start-timeout-taskqueue"
231
stl := "functional-reset-sticky-fire-schedule-to-start-timeout-taskqueue-sticky"
232
identity := "worker1"
233
234
stickyTaskQueue := &taskqueuepb.TaskQueue{Name: stl, Kind: enumspb.TASK_QUEUE_KIND_STICKY, NormalName: tl}
235
stickyScheduleToStartTimeout := 2 * time.Second
236
237
// Start workflow execution
238
request := &workflowservice.StartWorkflowExecutionRequest{
239
RequestId: uuid.NewString(),
240
Namespace: env.Namespace().String(),
241
WorkflowId: id,
242
WorkflowType: &commonpb.WorkflowType{Name: wt},
243
TaskQueue: &taskqueuepb.TaskQueue{Name: tl, Kind: enumspb.TASK_QUEUE_KIND_NORMAL},
244
Input: nil,
245
WorkflowRunTimeout: durationpb.New(100 * time.Second),
246
WorkflowTaskTimeout: durationpb.New(1 * time.Second),
247
Identity: identity,
248
}
249
250
we, err0 := env.FrontendClient().StartWorkflowExecution(s.Context(), request)
251
s.NoError(err0)
252
253
env.Logger.Info("StartWorkflowExecution", tag.WorkflowRunID(we.RunId))
254
workflowExecution := &commonpb.WorkflowExecution{
255
WorkflowId: id,
256
RunId: we.RunId,
257
}
258
259
// workflow logic
260
localActivityDone := false
261
failureCount := 5
262
wtHandler := func(task *workflowservice.PollWorkflowTaskQueueResponse) ([]*commandpb.Command, error) {
263
if !localActivityDone {
264
localActivityDone = true
265
266
return []*commandpb.Command{{
267
CommandType: enumspb.COMMAND_TYPE_RECORD_MARKER,
268
Attributes: &commandpb.Command_RecordMarkerCommandAttributes{RecordMarkerCommandAttributes: &commandpb.RecordMarkerCommandAttributes{
269
MarkerName: "local activity marker",
270
Details: map[string]*commonpb.Payloads{
271
"data": payloads.EncodeString("local activity marker"),
272
"result": payloads.EncodeString("local activity result"),
273
}}},
274
}}, nil
275
}
276
277
if failureCount > 0 {
278
failureCount--
279
return nil, errors.New("non deterministic error")
280
}
281
282
return []*commandpb.Command{{
283
CommandType: enumspb.COMMAND_TYPE_COMPLETE_WORKFLOW_EXECUTION,
284
Attributes: &commandpb.Command_CompleteWorkflowExecutionCommandAttributes{CompleteWorkflowExecutionCommandAttributes: &commandpb.CompleteWorkflowExecutionCommandAttributes{
285
Result: payloads.EncodeString("Done"),
286
}},
287
}}, nil
288
}
289
290
poller := &testcore.TaskPoller{
291
Client: env.FrontendClient(),
292
Namespace: env.Namespace().String(),
293
TaskQueue: &taskqueuepb.TaskQueue{Name: tl, Kind: enumspb.TASK_QUEUE_KIND_NORMAL},
294
Identity: identity,
295
WorkflowTaskHandler: wtHandler,
296
Logger: env.Logger,
297
T: s.T(),
298
StickyTaskQueue: stickyTaskQueue,
299
StickyScheduleToStartTimeout: stickyScheduleToStartTimeout,
300
}
301
302
_, err := poller.PollAndProcessWorkflowTask(testcore.WithRespondSticky)
303
env.Logger.Info("PollAndProcessWorkflowTask", tag.Error(err))
304
s.NoError(err)
305
306
_, err = env.FrontendClient().SignalWorkflowExecution(s.Context(), &workflowservice.SignalWorkflowExecutionRequest{
307
Namespace: env.Namespace().String(),
308
WorkflowExecution: workflowExecution,
309
SignalName: "signalA",
310
Input: payloads.EncodeString("signal input"),
311
Identity: identity,
312
RequestId: uuid.NewString(),
313
})
314
s.NoError(err)
315
316
// Reset sticky taskqueue before sticky workflow task starts
317
_, err = env.FrontendClient().ResetStickyTaskQueue(s.Context(), &workflowservice.ResetStickyTaskQueueRequest{
318
Namespace: env.Namespace().String(),
319
Execution: workflowExecution,
320
})
321
s.NoError(err)
322
323
// Wait for workflow task timeout
324
stickyTimeout := false
325
WaitForStickyTimeoutLoop:
326
for range 10 {
327
events := env.GetHistory(env.Namespace().String(), workflowExecution)
328
for _, event := range events {
329
if event.GetEventType() == enumspb.EVENT_TYPE_WORKFLOW_TASK_TIMED_OUT {
330
s.EqualHistoryEvents(`
331
1 WorkflowExecutionStarted
332
2 WorkflowTaskScheduled
333
3 WorkflowTaskStarted
334
4 WorkflowTaskCompleted
335
5 MarkerRecorded
336
6 WorkflowExecutionSignaled
337
7 WorkflowTaskScheduled
338
8 WorkflowTaskTimedOut {"TimeoutType":2} // enumspb.TIMEOUT_TYPE_SCHEDULE_TO_START
339
9 WorkflowTaskScheduled`, events)
340
stickyTimeout = true
341
break WaitForStickyTimeoutLoop
342
}
343
}
344
time.Sleep(time.Second) //nolint:forbidigo
345
}
346
s.True(stickyTimeout, "Workflow task not timed out")
347
348
for i := 1; i <= 3; i++ {
349
_, err = poller.PollAndProcessWorkflowTask(testcore.WithDumpHistory, testcore.WithRespondSticky, testcore.WithExpectedAttemptCount(i))
350
env.Logger.Info("PollAndProcessWorkflowTask", tag.Error(err))
351
s.NoError(err)
352
}
353
354
_, err = env.FrontendClient().SignalWorkflowExecution(s.Context(), &workflowservice.SignalWorkflowExecutionRequest{
355
Namespace: env.Namespace().String(),
356
WorkflowExecution: workflowExecution,
357
SignalName: "signalB",
358
Input: payloads.EncodeString("signal input"),
359
Identity: identity,
360
RequestId: uuid.NewString(),
361
})
362
s.NoError(err)
363
364
for i := 1; i <= 2; i++ {
365
_, err = poller.PollAndProcessWorkflowTask(testcore.WithDumpHistory, testcore.WithRespondSticky, testcore.WithExpectedAttemptCount(i))
366
env.Logger.Info("PollAndProcessWorkflowTask", tag.Error(err))
367
s.NoError(err)
368
}
369
370
events := env.GetHistory(env.Namespace().String(), workflowExecution)
371
s.EqualHistoryEvents(`
372
1 WorkflowExecutionStarted
373
2 WorkflowTaskScheduled
374
3 WorkflowTaskStarted
375
4 WorkflowTaskCompleted
376
5 MarkerRecorded
377
6 WorkflowExecutionSignaled
378
7 WorkflowTaskScheduled
379
8 WorkflowTaskTimedOut
380
9 WorkflowTaskScheduled
381
10 WorkflowTaskStarted
382
11 WorkflowTaskFailed
383
12 WorkflowExecutionSignaled
384
13 WorkflowTaskScheduled
385
14 WorkflowTaskStarted
386
15 WorkflowTaskFailed
387
16 WorkflowTaskScheduled`, events)
388
389
// Complete workflow execution
390
_, err = poller.PollAndProcessWorkflowTask(testcore.WithDumpHistory, testcore.WithRespondSticky, testcore.WithExpectedAttemptCount(3))
391
s.NoError(err)
392
393
events = env.GetHistory(env.Namespace().String(), workflowExecution)
394
s.EqualHistoryEvents(`
395
1 WorkflowExecutionStarted
396
2 WorkflowTaskScheduled
397
3 WorkflowTaskStarted
398
4 WorkflowTaskCompleted
399
5 MarkerRecorded
400
6 WorkflowExecutionSignaled
401
7 WorkflowTaskScheduled
402
8 WorkflowTaskTimedOut
403
9 WorkflowTaskScheduled
404
10 WorkflowTaskStarted
405
11 WorkflowTaskFailed // Two WFTs have failed
406
12 WorkflowExecutionSignaled
407
13 WorkflowTaskScheduled
408
14 WorkflowTaskStarted
409
15 WorkflowTaskFailed // Two WFTs have failed
410
16 WorkflowTaskScheduled
411
17 WorkflowTaskStarted
412
18 WorkflowTaskCompleted
413
19 WorkflowExecutionCompleted // Workflow has completed`, events)
414
}