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 }