go.temporal.io/server/tests/activity_test.go

1558 LOC · 0 covered · 1558 uncovered · 0 ranges · 0 concepts · 0 introducers · 0 tests

1 package tests
2
3 import (
4 "bytes"
5 "context"
6 "encoding/binary"
7 "errors"
8 "math/rand"
9 "strings"
10 "sync"
11 "sync/atomic"
12 "testing"
13 "time"
14
15 "github.com/google/uuid"
16 "github.com/stretchr/testify/assert"
17 "github.com/stretchr/testify/require"
18 commandpb "go.temporal.io/api/command/v1"
19 commonpb "go.temporal.io/api/common/v1"
20 enumspb "go.temporal.io/api/enums/v1"
21 historypb "go.temporal.io/api/history/v1"
22 "go.temporal.io/api/serviceerror"
23 taskqueuepb "go.temporal.io/api/taskqueue/v1"
24 "go.temporal.io/api/workflowservice/v1"
25 "go.temporal.io/sdk/activity"
26 sdkclient "go.temporal.io/sdk/client"
27 "go.temporal.io/sdk/temporal"
28 "go.temporal.io/sdk/worker"
29 "go.temporal.io/sdk/workflow"
30 "go.temporal.io/server/common/convert"
31 "go.temporal.io/server/common/log/tag"
32 "go.temporal.io/server/common/payload"
33 "go.temporal.io/server/common/payloads"
34 "go.temporal.io/server/common/testing/parallelsuite"
35 "go.temporal.io/server/service/history/consts"
36 "go.temporal.io/server/tests/testcore"
37 "google.golang.org/protobuf/types/known/durationpb"
38 )
39
40 type ActivityTestSuite struct {
41 parallelsuite.Suite[*ActivityTestSuite]
42 }
43
44 type ActivityClientTestSuite struct {
45 parallelsuite.Suite[*ActivityClientTestSuite]
46 }
47
48 func TestActivityTestSuite(t *testing.T) {
49 parallelsuite.Run(t, &ActivityTestSuite{})
50 }
51
52 func TestActivityClientTestSuite(t *testing.T) {
53 parallelsuite.Run(t, &ActivityClientTestSuite{})
54 }
55
56 func (s *ActivityClientTestSuite) TestActivityScheduleToClose_FiredDuringBackoff() {
57 env := testcore.NewEnv(s.T())
58 // We have activity that always fails.
59 // We have backoff timers and schedule_to_close activity timeout happens during that backoff timer.
60 // activity will be scheduled twice. After second failure (that should happen at ~4.2 sec) next retry will not
61 // be scheduled because "schedule_to_close" will happen before retry happens
62 initialRetryInterval := time.Second * 2
63 scheduleToCloseTimeout := 3 * time.Second
64 startToCloseTimeout := 1 * time.Second
65
66 activityRetryPolicy := &temporal.RetryPolicy{
67 InitialInterval: initialRetryInterval,
68 BackoffCoefficient: 1,
69 MaximumInterval: time.Second * 10,
70 }
71
72 var activityCompleted atomic.Int32
73 activityFunction := func() (string, error) {
74 activityErr := errors.New("bad-luck-please-retry") //nolint:err113
75 activityCompleted.Add(1)
76 return "", activityErr
77 }
78
79 workflowFn := func(ctx workflow.Context) (string, error) {
80 var ret string
81 err := workflow.ExecuteActivity(workflow.WithActivityOptions(ctx, workflow.ActivityOptions{
82 DisableEagerExecution: true,
83 StartToCloseTimeout: startToCloseTimeout,
84 ScheduleToCloseTimeout: scheduleToCloseTimeout,
85 RetryPolicy: activityRetryPolicy,
86 }), activityFunction).Get(ctx, &ret)
87 return "done!", err
88 }
89
90 env.SdkWorker().RegisterWorkflow(workflowFn)
91 env.SdkWorker().RegisterActivity(activityFunction)
92
93 wfId := "functional-test-gethistoryreverse"
94 workflowOptions := sdkclient.StartWorkflowOptions{
95 ID: wfId,
96 TaskQueue: env.WorkerTaskQueue(),
97 }
98 workflowRun, err := env.SdkClient().ExecuteWorkflow(s.Context(), workflowOptions, workflowFn)
99 s.NoError(err)
100
101 var out string
102 err = workflowRun.Get(s.Context(), &out)
103
104 s.Error(err)
105 var wfExecutionError *temporal.WorkflowExecutionError
106 s.ErrorAs(err, &wfExecutionError)
107 var activityError *temporal.ActivityError
108 s.ErrorAs(wfExecutionError, &activityError)
109 s.Equal(enumspb.RETRY_STATE_TIMEOUT, activityError.RetryState())
110
111 s.Equal(int32(2), activityCompleted.Load())
112
113 }
114
115 func (s *ActivityClientTestSuite) TestActivityScheduleToClose_FiredDuringActivityRun() {
116 env := testcore.NewEnv(s.T())
117 // We have activity that always fails.
118 // We have backoff timers and schedule_to_close activity timeout happens while activity is running.
119 // activity will be scheduled twice.
120 // "schedule_to_close" timer should fire while activity is running for a second time
121 scheduleToCloseTimeout := 7 * time.Second
122 startToCloseTimeout := 3 * time.Second
123
124 activityRetryPolicy := &temporal.RetryPolicy{
125 InitialInterval: time.Second * 1,
126 BackoffCoefficient: 1,
127 }
128 var activityFinishedAt time.Time
129 var workflowFinishedAt time.Time
130
131 var activityCompleted atomic.Int32
132 var wg sync.WaitGroup
133
134 // we need schedule_to_close timeout to fire when 3rd retry is running
135 // with 2 sec run time and 1 sec between retries 3rd retry will start around 2+1+2+1=6 sec
136 // with 2 sec run time it will finish at 8 sec
137 // schedule to close is set to 7 sec. This way schedule to close timeout should fire.
138 activityFunction := func() (string, error) {
139 defer wg.Done()
140 activityErr := errors.New("bad-luck-please-retry") //nolint:err113
141 time.Sleep(2 * time.Second) //nolint:forbidigo
142 activityCompleted.Add(1)
143 activityFinishedAt = time.Now().UTC()
144 return "", activityErr
145 }
146
147 wg.Add(3) // activity should be executed 3 times
148 workflowFn := func(ctx workflow.Context) (string, error) {
149 var ret string
150 err := workflow.ExecuteActivity(workflow.WithActivityOptions(ctx, workflow.ActivityOptions{
151 DisableEagerExecution: true,
152 StartToCloseTimeout: startToCloseTimeout,
153 ScheduleToCloseTimeout: scheduleToCloseTimeout,
154 RetryPolicy: activityRetryPolicy,
155 }), activityFunction).Get(ctx, &ret)
156 return "done!", err
157 }
158
159 env.SdkWorker().RegisterWorkflow(workflowFn)
160 env.SdkWorker().RegisterActivity(activityFunction)
161
162 workflowOptions := sdkclient.StartWorkflowOptions{
163 ID: s.T().Name(),
164 TaskQueue: env.WorkerTaskQueue(),
165 }
166 workflowRun, err := env.SdkClient().ExecuteWorkflow(s.Context(), workflowOptions, workflowFn)
167 s.NoError(err)
168
169 var out string
170 err = workflowRun.Get(s.Context(), &out)
171 var activityError *temporal.ActivityError
172 s.ErrorAs(err, &activityError)
173 s.Equal(enumspb.RETRY_STATE_TIMEOUT, activityError.RetryState())
174 var timeoutError *temporal.TimeoutError
175 s.ErrorAs(activityError, &timeoutError)
176 s.Equal(enumspb.TIMEOUT_TYPE_SCHEDULE_TO_CLOSE, timeoutError.TimeoutType())
177 // schedule to close timeout should fire while last activity is still running.
178 s.Equal(int32(2), activityCompleted.Load())
179
180 // we expect activities to be executed 3 times. Wait for the last activity to finish
181 workflowFinishedAt = time.Now().UTC()
182 wg.Wait() // Wait for activity to finish
183 s.Equal(int32(3), activityCompleted.Load())
184 s.True(activityFinishedAt.After(workflowFinishedAt))
185 }
186
187 func (s *ActivityClientTestSuite) Test_ActivityTimeouts() {
188 env := testcore.NewEnv(s.T())
189 activityFn := func(ctx context.Context) error {
190 info := activity.GetInfo(ctx)
191 if strings.HasPrefix(info.ActivityID, "Heartbeat") {
192 go func() {
193 // NOTE: due to client side heartbeat batching, heartbeat may be sent
194 // later than expected.
195 // e.g. if activity heartbeat timeout is 2s,
196 // and we call RecordHeartbeat() at 0s, 0.5s, 1s, 1.5s
197 // the client by default will send two heartbeats at 0s and 2*0.8=1.6s
198 // Now if when running the test, this heartbeat goroutine becomes slow,
199 // and call RecordHeartbeat() after 1.6s, then that heartbeat will be sent
200 // to server at 3.2s (the next batch).
201 // Since the entire activity will finish at 5s, there won't be
202 // any heartbeat timeout error.
203 // so here, we reduce the duration between two heartbeats, so that they are
204 // more likely be sent in the heartbeat batch at 1.6s
205 // (basically increasing the room for delay in heartbeat goroutine from 0.1s to 1s)
206 for i := range 3 {
207 activity.RecordHeartbeat(ctx, i)
208 time.Sleep(200 * time.Millisecond) //nolint:forbidigo
209 }
210 }()
211 }
212
213 time.Sleep(5 * time.Second) //nolint:forbidigo
214 return nil
215 }
216
217 var err1, err2, err3, err4, err5, err6, err7, err8 error
218 workflowFn := func(ctx workflow.Context) error {
219 noRetryPolicy := &temporal.RetryPolicy{
220 MaximumAttempts: 1, // disable retry
221 }
222 ctx1 := workflow.WithActivityOptions(ctx, workflow.ActivityOptions{
223 ActivityID: "ScheduleToStart",
224 ScheduleToStartTimeout: 2 * time.Second,
225 StartToCloseTimeout: 2 * time.Second,
226 TaskQueue: "NoWorkerTaskQueue",
227 RetryPolicy: noRetryPolicy,
228 })
229 f1 := workflow.ExecuteActivity(ctx1, activityFn)
230
231 ctx2 := workflow.WithActivityOptions(ctx, workflow.ActivityOptions{
232 ActivityID: "StartToClose",
233 ScheduleToStartTimeout: 2 * time.Second,
234 StartToCloseTimeout: 2 * time.Second,
235 RetryPolicy: noRetryPolicy,
236 })
237 f2 := workflow.ExecuteActivity(ctx2, activityFn)
238
239 ctx3 := workflow.WithActivityOptions(ctx, workflow.ActivityOptions{
240 ActivityID: "ScheduleToClose",
241 ScheduleToCloseTimeout: 2 * time.Second,
242 StartToCloseTimeout: 3 * time.Second,
243 RetryPolicy: noRetryPolicy,
244 })
245 f3 := workflow.ExecuteActivity(ctx3, activityFn)
246
247 ctx4 := workflow.WithActivityOptions(ctx, workflow.ActivityOptions{
248 ActivityID: "ScheduleToCloseNotSet",
249 StartToCloseTimeout: 2 * time.Second,
250 RetryPolicy: noRetryPolicy,
251 })
252 f4 := workflow.ExecuteActivity(ctx4, activityFn)
253
254 ctx5 := workflow.WithActivityOptions(ctx, workflow.ActivityOptions{
255 ActivityID: "StartToCloseNotSet",
256 ScheduleToCloseTimeout: 2 * time.Second,
257 RetryPolicy: noRetryPolicy,
258 })
259 f5 := workflow.ExecuteActivity(ctx5, activityFn)
260
261 ctx6 := workflow.WithActivityOptions(ctx, workflow.ActivityOptions{
262 ActivityID: "Heartbeat",
263 StartToCloseTimeout: 10 * time.Second,
264 HeartbeatTimeout: 1 * time.Second,
265 RetryPolicy: noRetryPolicy,
266 })
267 f6 := workflow.ExecuteActivity(ctx6, activityFn)
268
269 ctx7 := workflow.WithActivityOptions(ctx, workflow.ActivityOptions{
270 ActivityID: "HeartbeatWithScheduleToClose",
271 ScheduleToCloseTimeout: 2 * time.Second,
272 HeartbeatTimeout: 1 * time.Second,
273 })
274 f7 := workflow.ExecuteActivity(ctx7, activityFn)
275
276 ctx8 := workflow.WithActivityOptions(ctx, workflow.ActivityOptions{
277 ActivityID: "StartToCloseWithScheduleToClose",
278 ScheduleToCloseTimeout: 3 * time.Second,
279 StartToCloseTimeout: 2 * time.Second,
280 RetryPolicy: &temporal.RetryPolicy{
281 InitialInterval: 2 * time.Second,
282 },
283 })
284 f8 := workflow.ExecuteActivity(ctx8, activityFn)
285
286 err1 = f1.Get(ctx1, nil)
287 err2 = f2.Get(ctx2, nil)
288 err3 = f3.Get(ctx3, nil)
289 err4 = f4.Get(ctx4, nil)
290 err5 = f5.Get(ctx5, nil)
291 err6 = f6.Get(ctx6, nil)
292 err7 = f7.Get(ctx7, nil)
293 err8 = f8.Get(ctx8, nil)
294 return nil
295 }
296
297 env.SdkWorker().RegisterActivity(activityFn)
298 env.SdkWorker().RegisterWorkflow(workflowFn)
299
300 workflowOptions := sdkclient.StartWorkflowOptions{
301 ID: "functional-test-activity-timeouts",
302 TaskQueue: env.WorkerTaskQueue(),
303 WorkflowRunTimeout: 20 * time.Second,
304 }
305 workflowRun, err := env.SdkClient().ExecuteWorkflow(s.Context(), workflowOptions, workflowFn)
306 s.NoError(err)
307
308 s.NotNil(workflowRun)
309 s.NotEmpty(workflowRun.GetRunID())
310 err = workflowRun.Get(s.Context(), nil)
311 s.NoError(err)
312
313 // verify activity timeout type
314 var activityErr *temporal.ActivityError
315 s.ErrorAs(err1, &activityErr)
316 s.Equal("ScheduleToStart", activityErr.ActivityID())
317 timeoutErr, ok := activityErr.Unwrap().(*temporal.TimeoutError)
318 s.True(ok)
319 s.Equal(enumspb.TIMEOUT_TYPE_SCHEDULE_TO_START, timeoutErr.TimeoutType())
320
321 s.ErrorAs(err2, &activityErr)
322 s.Equal("StartToClose", activityErr.ActivityID())
323 timeoutErr, ok = activityErr.Unwrap().(*temporal.TimeoutError)
324 s.True(ok)
325 s.Equal(enumspb.TIMEOUT_TYPE_START_TO_CLOSE, timeoutErr.TimeoutType())
326
327 s.ErrorAs(err3, &activityErr)
328 s.Equal("ScheduleToClose", activityErr.ActivityID())
329 timeoutErr, ok = activityErr.Unwrap().(*temporal.TimeoutError)
330 s.True(ok)
331 s.Equal(enumspb.TIMEOUT_TYPE_SCHEDULE_TO_CLOSE, timeoutErr.TimeoutType())
332
333 s.ErrorAs(err4, &activityErr)
334 s.Equal("ScheduleToCloseNotSet", activityErr.ActivityID())
335 timeoutErr, ok = activityErr.Unwrap().(*temporal.TimeoutError)
336 s.True(ok)
337 s.Equal(enumspb.TIMEOUT_TYPE_START_TO_CLOSE, timeoutErr.TimeoutType())
338
339 s.ErrorAs(err5, &activityErr)
340 s.Equal("StartToCloseNotSet", activityErr.ActivityID())
341 timeoutErr, ok = activityErr.Unwrap().(*temporal.TimeoutError)
342 s.True(ok)
343 s.Equal(enumspb.TIMEOUT_TYPE_SCHEDULE_TO_CLOSE, timeoutErr.TimeoutType())
344
345 s.ErrorAs(err6, &activityErr)
346 s.Equal("Heartbeat", activityErr.ActivityID())
347 timeoutErr, ok = activityErr.Unwrap().(*temporal.TimeoutError)
348 s.True(ok)
349 s.Equal(enumspb.TIMEOUT_TYPE_HEARTBEAT, timeoutErr.TimeoutType())
350 s.True(timeoutErr.HasLastHeartbeatDetails())
351 var v int
352 s.NoError(timeoutErr.LastHeartbeatDetails(&v))
353 s.Equal(2, v)
354
355 s.ErrorAs(err7, &activityErr)
356 s.Equal("HeartbeatWithScheduleToClose", activityErr.ActivityID())
357 timeoutErr, ok = activityErr.Unwrap().(*temporal.TimeoutError)
358 s.True(ok)
359 s.Equal(enumspb.TIMEOUT_TYPE_SCHEDULE_TO_CLOSE, timeoutErr.TimeoutType())
360 s.Equal("Not enough time to schedule next retry before activity ScheduleToClose timeout, giving up retrying (type: ScheduleToClose)", timeoutErr.Error())
361 s.True(timeoutErr.HasLastHeartbeatDetails())
362 v = 0
363 s.NoError(timeoutErr.LastHeartbeatDetails(&v))
364 s.Equal(2, v)
365
366 s.ErrorAs(err8, &activityErr)
367 s.Equal("StartToCloseWithScheduleToClose", activityErr.ActivityID())
368 timeoutErr, ok = activityErr.Unwrap().(*temporal.TimeoutError)
369 s.True(ok)
370 s.Equal(enumspb.TIMEOUT_TYPE_SCHEDULE_TO_CLOSE, timeoutErr.TimeoutType())
371 s.Equal("Not enough time to schedule next retry before activity ScheduleToClose timeout, giving up retrying (type: ScheduleToClose)", timeoutErr.Error())
372 }
373
374 func (s *ActivityTestSuite) TestActivityHeartBeatWorkflow_Success() {
375 env := testcore.NewEnv(s.T())
376 id := "functional-heartbeat-test"
377 wt := "functional-heartbeat-test-type"
378 tl := "functional-heartbeat-test-taskqueue"
379 identity := "worker1"
380 activityName := "activity_timer"
381
382 workflowType := &commonpb.WorkflowType{Name: wt}
383
384 taskQueue := &taskqueuepb.TaskQueue{Name: tl, Kind: enumspb.TASK_QUEUE_KIND_NORMAL}
385
386 header := &commonpb.Header{
387 Fields: map[string]*commonpb.Payload{"tracing": payload.EncodeString("sample data")},
388 }
389
390 request := &workflowservice.StartWorkflowExecutionRequest{
391 RequestId: uuid.NewString(),
392 Namespace: env.Namespace().String(),
393 WorkflowId: id,
394 WorkflowType: workflowType,
395 TaskQueue: taskQueue,
396 Input: nil,
397 Header: header,
398 WorkflowRunTimeout: durationpb.New(100 * time.Second),
399 WorkflowTaskTimeout: durationpb.New(1 * time.Second),
400 Identity: identity,
401 }
402
403 we, err0 := env.FrontendClient().StartWorkflowExecution(s.Context(), request)
404 s.NoError(err0)
405
406 env.Logger.Info("StartWorkflowExecution", tag.WorkflowRunID(we.RunId))
407
408 workflowComplete := false
409 activityCount := int32(1)
410 activityCounter := int32(0)
411
412 wtHandler := func(task *workflowservice.PollWorkflowTaskQueueResponse) ([]*commandpb.Command, error) {
413 if activityCounter < activityCount {
414 activityCounter++
415
416 return []*commandpb.Command{{
417 CommandType: enumspb.COMMAND_TYPE_SCHEDULE_ACTIVITY_TASK,
418 Attributes: &commandpb.Command_ScheduleActivityTaskCommandAttributes{ScheduleActivityTaskCommandAttributes: &commandpb.ScheduleActivityTaskCommandAttributes{
419 ActivityId: convert.Int32ToString(activityCounter),
420 ActivityType: &commonpb.ActivityType{Name: activityName},
421 TaskQueue: &taskqueuepb.TaskQueue{Name: tl, Kind: enumspb.TASK_QUEUE_KIND_NORMAL},
422 Input: payloads.EncodeString("activity-input"),
423 Header: header,
424 ScheduleToCloseTimeout: durationpb.New(15 * time.Second),
425 ScheduleToStartTimeout: durationpb.New(1 * time.Second),
426 StartToCloseTimeout: durationpb.New(15 * time.Second),
427 HeartbeatTimeout: durationpb.New(1 * time.Second),
428 }},
429 }}, nil
430 }
431
432 env.Logger.Info("Completing Workflow")
433
434 workflowComplete = true
435 return []*commandpb.Command{{
436 CommandType: enumspb.COMMAND_TYPE_COMPLETE_WORKFLOW_EXECUTION,
437 Attributes: &commandpb.Command_CompleteWorkflowExecutionCommandAttributes{CompleteWorkflowExecutionCommandAttributes: &commandpb.CompleteWorkflowExecutionCommandAttributes{
438 Result: payloads.EncodeString("Done"),
439 }},
440 }}, nil
441 }
442
443 activityExecutedCount := 0
444 atHandler := func(task *workflowservice.PollActivityTaskQueueResponse) (*commonpb.Payloads, bool, error) {
445 s.Equal(id, task.WorkflowExecution.GetWorkflowId())
446 s.Equal(activityName, task.ActivityType.GetName())
447 for i := range 10 {
448 env.Logger.Info("Heartbeating for activity", tag.ActivityID(task.ActivityId), tag.Counter(i))
449 _, err := env.FrontendClient().RecordActivityTaskHeartbeat(s.Context(), &workflowservice.RecordActivityTaskHeartbeatRequest{
450 Namespace: env.Namespace().String(),
451 TaskToken: task.TaskToken,
452 Details: payloads.EncodeString("details"),
453 })
454 s.NoError(err)
455 time.Sleep(10 * time.Millisecond) //nolint:forbidigo
456 }
457 activityExecutedCount++
458 return payloads.EncodeString("Activity Result"), false, nil
459 }
460
461 poller := &testcore.TaskPoller{
462 Client: env.FrontendClient(),
463 Namespace: env.Namespace().String(),
464 TaskQueue: taskQueue,
465 Identity: identity,
466 WorkflowTaskHandler: wtHandler,
467 ActivityTaskHandler: atHandler,
468 Logger: env.Logger,
469 T: s.T(),
470 }
471
472 _, err := poller.PollAndProcessWorkflowTask()
473 s.True(err == nil || errors.Is(err, testcore.ErrNoTasks))
474
475 err = poller.PollAndProcessActivityTask(false)
476 s.True(err == nil || errors.Is(err, testcore.ErrNoTasks))
477
478 env.Logger.Info("Waiting for workflow to complete", tag.WorkflowRunID(we.RunId))
479
480 s.False(workflowComplete)
481 _, err = poller.PollAndProcessWorkflowTask(testcore.WithDumpHistory)
482 s.NoError(err)
483 s.True(workflowComplete)
484 s.Equal(1, activityExecutedCount)
485
486 // go over history and verify that the activity task scheduled event has header on it
487 events := env.GetHistory(env.Namespace().String(), &commonpb.WorkflowExecution{
488 WorkflowId: id,
489 RunId: we.GetRunId(),
490 })
491
492 s.EqualHistoryEvents(`
493 1 WorkflowExecutionStarted
494 2 WorkflowTaskScheduled
495 3 WorkflowTaskStarted
496 4 WorkflowTaskCompleted
497 5 ActivityTaskScheduled {"Header":{"Fields":{"tracing":{"Data":"\"sample data\""}}} }
498 6 ActivityTaskStarted
499 7 ActivityTaskCompleted
500 8 WorkflowTaskScheduled
501 9 WorkflowTaskStarted
502 10 WorkflowTaskCompleted
503 11 WorkflowExecutionCompleted`, events)
504 }
505
506 func (s *ActivityTestSuite) TestActivityRetry() {
507 env := testcore.NewEnv(s.T())
508 tv := env.Tv()
509
510 activityName := "activity_retry"
511 timeoutActivityName := "timeout_activity"
512 request := &workflowservice.StartWorkflowExecutionRequest{
513 RequestId: uuid.NewString(),
514 Namespace: env.Namespace().String(),
515 WorkflowId: tv.WorkflowID(),
516 WorkflowType: tv.WorkflowType(),
517 TaskQueue: tv.TaskQueue(),
518 Input: nil,
519 WorkflowRunTimeout: durationpb.New(100 * time.Second),
520 WorkflowTaskTimeout: durationpb.New(1 * time.Second),
521 Identity: tv.WorkerIdentity(),
522 }
523
524 we, err0 := env.FrontendClient().StartWorkflowExecution(s.Context(), request)
525 s.NoError(err0)
526
527 env.Logger.Info("StartWorkflowExecution", tag.WorkflowRunID(we.RunId))
528
529 workflowComplete := false
530 activitiesScheduled := false
531 var activityAScheduled, activityAFailed, activityBScheduled, activityBTimeout *historypb.HistoryEvent
532
533 wtHandler := func(task *workflowservice.PollWorkflowTaskQueueResponse) ([]*commandpb.Command, error) {
534 if !activitiesScheduled {
535 activitiesScheduled = true
536
537 return []*commandpb.Command{
538 {
539 CommandType: enumspb.COMMAND_TYPE_SCHEDULE_ACTIVITY_TASK,
540 Attributes: &commandpb.Command_ScheduleActivityTaskCommandAttributes{ScheduleActivityTaskCommandAttributes: &commandpb.ScheduleActivityTaskCommandAttributes{
541 ActivityId: "A",
542 ActivityType: &commonpb.ActivityType{Name: activityName},
543 TaskQueue: tv.TaskQueue(),
544 Input: payloads.EncodeString("1"),
545 ScheduleToCloseTimeout: durationpb.New(4 * time.Second),
546 ScheduleToStartTimeout: durationpb.New(4 * time.Second),
547 StartToCloseTimeout: durationpb.New(4 * time.Second),
548 HeartbeatTimeout: durationpb.New(1 * time.Second),
549 RetryPolicy: &commonpb.RetryPolicy{
550 InitialInterval: durationpb.New(1 * time.Second),
551 MaximumAttempts: 3,
552 MaximumInterval: durationpb.New(1 * time.Second),
553 NonRetryableErrorTypes: []string{"bad-bug"},
554 BackoffCoefficient: 1,
555 },
556 }}},
557 {
558 CommandType: enumspb.COMMAND_TYPE_SCHEDULE_ACTIVITY_TASK,
559 Attributes: &commandpb.Command_ScheduleActivityTaskCommandAttributes{ScheduleActivityTaskCommandAttributes: &commandpb.ScheduleActivityTaskCommandAttributes{
560 ActivityId: "B",
561 ActivityType: &commonpb.ActivityType{Name: timeoutActivityName},
562 TaskQueue: &taskqueuepb.TaskQueue{Name: "no_worker_taskqueue", Kind: enumspb.TASK_QUEUE_KIND_NORMAL},
563 Input: payloads.EncodeString("2"),
564 ScheduleToCloseTimeout: durationpb.New(5 * time.Second),
565 ScheduleToStartTimeout: durationpb.New(5 * time.Second),
566 StartToCloseTimeout: durationpb.New(5 * time.Second),
567 HeartbeatTimeout: durationpb.New(0 * time.Second),
568 }}},
569 }, nil
570 } else if task.PreviousStartedEventId > 0 {
571 for _, event := range task.History.Events[task.PreviousStartedEventId:] {
572 switch event.GetEventType() { // nolint:exhaustive
573 case enumspb.EVENT_TYPE_ACTIVITY_TASK_SCHEDULED:
574 switch event.GetActivityTaskScheduledEventAttributes().GetActivityId() {
575 case "A":
576 activityAScheduled = event
577 case "B":
578 activityBScheduled = event
579 }
580
581 case enumspb.EVENT_TYPE_ACTIVITY_TASK_FAILED:
582 if event.GetActivityTaskFailedEventAttributes().GetScheduledEventId() == activityAScheduled.GetEventId() {
583 activityAFailed = event
584 }
585
586 case enumspb.EVENT_TYPE_ACTIVITY_TASK_TIMED_OUT:
587 if event.GetActivityTaskTimedOutEventAttributes().GetScheduledEventId() == activityBScheduled.GetEventId() {
588 activityBTimeout = event
589 }
590 }
591 }
592 }
593
594 if activityAFailed != nil && activityBTimeout != nil {
595 env.Logger.Info("Completing Workflow")
596 workflowComplete = true
597 return []*commandpb.Command{{
598 CommandType: enumspb.COMMAND_TYPE_COMPLETE_WORKFLOW_EXECUTION,
599 Attributes: &commandpb.Command_CompleteWorkflowExecutionCommandAttributes{CompleteWorkflowExecutionCommandAttributes: &commandpb.CompleteWorkflowExecutionCommandAttributes{
600 Result: payloads.EncodeString("Done"),
601 }},
602 }}, nil
603 }
604
605 return []*commandpb.Command{}, nil
606 }
607
608 activityExecutedCount := 0
609 atHandler := func(task *workflowservice.PollActivityTaskQueueResponse) (*commonpb.Payloads, bool, error) {
610 s.Equal(tv.WorkflowID(), task.WorkflowExecution.GetWorkflowId())
611 s.Equal(activityName, task.ActivityType.GetName())
612 var err error
613 if activityExecutedCount == 0 {
614 err = errors.New("bad-luck-please-retry") //nolint:err113
615 } else if activityExecutedCount == 1 {
616 err = temporal.NewNonRetryableApplicationError("bad-bug", "", nil)
617 }
618 activityExecutedCount++
619 return nil, false, err
620 }
621
622 poller := &testcore.TaskPoller{
623 Client: env.FrontendClient(),
624 Namespace: env.Namespace().String(),
625 TaskQueue: tv.TaskQueue(),
626 Identity: tv.WorkerIdentity(),
627 WorkflowTaskHandler: wtHandler,
628 ActivityTaskHandler: atHandler,
629 Logger: env.Logger,
630 T: s.T(),
631 }
632
633 describeWorkflowExecution := func() (*workflowservice.DescribeWorkflowExecutionResponse, error) {
634 return env.FrontendClient().DescribeWorkflowExecution(s.Context(), &workflowservice.DescribeWorkflowExecutionRequest{
635 Namespace: env.Namespace().String(),
636 Execution: &commonpb.WorkflowExecution{
637 WorkflowId: tv.WorkflowID(),
638 RunId: we.RunId,
639 },
640 })
641 }
642
643 _, err := poller.PollAndProcessWorkflowTask()
644 s.NoError(err)
645
646 _, err = env.TaskPoller().PollAndHandleActivityTask(tv,
647 func(task *workflowservice.PollActivityTaskQueueResponse) (*workflowservice.RespondActivityTaskCompletedRequest, error) {
648 s.Equal(tv.WorkflowID(), task.WorkflowExecution.GetWorkflowId())
649 s.Equal(activityName, task.ActivityType.GetName())
650 return nil, errors.New("bad-luck-please-retry")
651 })
652 s.NoError(err)
653
654 descResp, err := describeWorkflowExecution()
655 s.NoError(err)
656 for _, pendingActivity := range descResp.GetPendingActivities() {
657 if pendingActivity.GetActivityId() == "A" {
658 s.NotNil(pendingActivity.GetLastFailure().GetApplicationFailureInfo())
659 expectedErrString := "bad-luck-please-retry"
660 s.Equal(expectedErrString, pendingActivity.GetLastFailure().GetMessage())
661 s.False(pendingActivity.GetLastFailure().GetApplicationFailureInfo().GetNonRetryable())
662 s.Equal(tv.WorkerIdentity(), pendingActivity.GetLastWorkerIdentity())
663 }
664 }
665
666 _, err = env.TaskPoller().PollAndHandleActivityTask(tv,
667 func(task *workflowservice.PollActivityTaskQueueResponse) (*workflowservice.RespondActivityTaskCompletedRequest, error) {
668 s.Equal(tv.WorkflowID(), task.WorkflowExecution.GetWorkflowId())
669 s.Equal(activityName, task.ActivityType.GetName())
670 return nil, temporal.NewNonRetryableApplicationError("bad-bug", "", nil)
671 })
672 s.NoError(err)
673
674 descResp, err = describeWorkflowExecution()
675 s.NoError(err)
676 s.Len(descResp.GetPendingActivities(), 1)
677 s.Equal("B", descResp.GetPendingActivities()[0].GetActivityId())
678
679 env.Logger.Info("Waiting for workflow to complete", tag.WorkflowRunID(we.RunId))
680 for i := range 3 {
681 s.False(workflowComplete)
682
683 env.Logger.Info("Processing workflow task:", tag.Counter(i))
684 _, err := poller.PollAndProcessWorkflowTask(testcore.WithRetries(1))
685 if err != nil {
686 s.PrintHistoryEvents(env.GetHistory(env.Namespace().String(), &commonpb.WorkflowExecution{
687 WorkflowId: tv.WorkflowID(),
688 RunId: we.GetRunId(),
689 }))
690 }
691 s.NoError(err, "Poll for workflow task failed")
692
693 if workflowComplete {
694 break
695 }
696 }
697
698 s.True(workflowComplete)
699 }
700
701 func (s *ActivityTestSuite) TestActivityRetry_Infinite() {
702 env := testcore.NewEnv(s.T())
703 id := "functional-activity-retry-test"
704 wt := "functional-activity-retry-type"
705 tl := "functional-activity-retry-taskqueue"
706 identity := "worker1"
707 activityName := "activity_retry"
708
709 workflowType := &commonpb.WorkflowType{Name: wt}
710
711 taskQueue := &taskqueuepb.TaskQueue{Name: tl, Kind: enumspb.TASK_QUEUE_KIND_NORMAL}
712
713 request := &workflowservice.StartWorkflowExecutionRequest{
714 RequestId: uuid.NewString(),
715 Namespace: env.Namespace().String(),
716 WorkflowId: id,
717 WorkflowType: workflowType,
718 TaskQueue: taskQueue,
719 Input: nil,
720 WorkflowRunTimeout: durationpb.New(100 * time.Second),
721 WorkflowTaskTimeout: durationpb.New(1 * time.Second),
722 Identity: identity,
723 }
724
725 we, err0 := env.FrontendClient().StartWorkflowExecution(s.Context(), request)
726 s.NoError(err0)
727
728 env.Logger.Info("StartWorkflowExecution", tag.WorkflowRunID(we.RunId))
729
730 workflowComplete := false
731 activitiesScheduled := false
732
733 wtHandler := func(task *workflowservice.PollWorkflowTaskQueueResponse) ([]*commandpb.Command, error) {
734 if !activitiesScheduled {
735 activitiesScheduled = true
736
737 return []*commandpb.Command{{
738 CommandType: enumspb.COMMAND_TYPE_SCHEDULE_ACTIVITY_TASK,
739 Attributes: &commandpb.Command_ScheduleActivityTaskCommandAttributes{
740 ScheduleActivityTaskCommandAttributes: &commandpb.ScheduleActivityTaskCommandAttributes{
741 ActivityId: "A",
742 ActivityType: &commonpb.ActivityType{Name: activityName},
743 TaskQueue: &taskqueuepb.TaskQueue{Name: tl, Kind: enumspb.TASK_QUEUE_KIND_NORMAL},
744 Input: payloads.EncodeString("1"),
745 StartToCloseTimeout: durationpb.New(100 * time.Second),
746 RetryPolicy: &commonpb.RetryPolicy{
747 MaximumAttempts: 0,
748 BackoffCoefficient: 1,
749 },
750 }},
751 }}, nil
752 }
753
754 workflowComplete = true
755 return []*commandpb.Command{{
756 CommandType: enumspb.COMMAND_TYPE_COMPLETE_WORKFLOW_EXECUTION,
757 Attributes: &commandpb.Command_CompleteWorkflowExecutionCommandAttributes{
758 CompleteWorkflowExecutionCommandAttributes: &commandpb.CompleteWorkflowExecutionCommandAttributes{
759 Result: payloads.EncodeString("Done"),
760 },
761 },
762 }}, nil
763 }
764
765 activityExecutedCount := 0
766 activityExecutedLimit := 4
767 atHandler := func(task *workflowservice.PollActivityTaskQueueResponse) (*commonpb.Payloads, bool, error) {
768 s.Equal(id, task.WorkflowExecution.GetWorkflowId())
769 s.Equal(activityName, task.ActivityType.GetName())
770
771 var err error
772 if activityExecutedCount < activityExecutedLimit {
773 err = errors.New("retry-error") //nolint:err113
774 } else if activityExecutedCount == activityExecutedLimit {
775 err = nil
776 }
777 activityExecutedCount++
778 return nil, false, err
779 }
780
781 poller := &testcore.TaskPoller{
782 Client: env.FrontendClient(),
783 Namespace: env.Namespace().String(),
784 TaskQueue: taskQueue,
785 Identity: identity,
786 WorkflowTaskHandler: wtHandler,
787 ActivityTaskHandler: atHandler,
788 Logger: env.Logger,
789 T: s.T(),
790 }
791
792 _, err := poller.PollAndProcessWorkflowTask()
793 s.NoError(err)
794
795 for i := 0; i <= activityExecutedLimit; i++ {
796 err = poller.PollAndProcessActivityTask(false)
797 s.NoError(err)
798 }
799
800 _, err = poller.PollAndProcessWorkflowTask(testcore.WithRetries(1))
801 s.NoError(err)
802 s.True(workflowComplete)
803 }
804
805 func (s *ActivityTestSuite) TestActivityHeartBeatWorkflow_Timeout() {
806 env := testcore.NewEnv(s.T())
807 id := "functional-heartbeat-timeout-test"
808 wt := "functional-heartbeat-timeout-test-type"
809 tl := "functional-heartbeat-timeout-test-taskqueue"
810 identity := "worker1"
811 activityName := "activity_timer"
812
813 workflowType := &commonpb.WorkflowType{Name: wt}
814
815 taskQueue := &taskqueuepb.TaskQueue{Name: tl, Kind: enumspb.TASK_QUEUE_KIND_NORMAL}
816
817 request := &workflowservice.StartWorkflowExecutionRequest{
818 RequestId: uuid.NewString(),
819 Namespace: env.Namespace().String(),
820 WorkflowId: id,
821 WorkflowType: workflowType,
822 TaskQueue: taskQueue,
823 Input: nil,
824 WorkflowRunTimeout: durationpb.New(100 * time.Second),
825 WorkflowTaskTimeout: durationpb.New(1 * time.Second),
826 Identity: identity,
827 }
828
829 we, err0 := env.FrontendClient().StartWorkflowExecution(s.Context(), request)
830 s.NoError(err0)
831
832 env.Logger.Info("StartWorkflowExecution", tag.WorkflowRunID(we.RunId))
833
834 workflowComplete := false
835 activityCount := int32(1)
836 activityCounter := int32(0)
837
838 wtHandler := func(task *workflowservice.PollWorkflowTaskQueueResponse) ([]*commandpb.Command, error) {
839
840 env.Logger.Info("Calling WorkflowTask Handler", tag.Counter(int(activityCounter)), tag.Number(int64(activityCount)))
841
842 if activityCounter < activityCount {
843 activityCounter++
844 buf := new(bytes.Buffer)
845 s.NoError(binary.Write(buf, binary.LittleEndian, activityCounter))
846
847 return []*commandpb.Command{{
848 CommandType: enumspb.COMMAND_TYPE_SCHEDULE_ACTIVITY_TASK,
849 Attributes: &commandpb.Command_ScheduleActivityTaskCommandAttributes{ScheduleActivityTaskCommandAttributes: &commandpb.ScheduleActivityTaskCommandAttributes{
850 ActivityId: convert.Int32ToString(activityCounter),
851 ActivityType: &commonpb.ActivityType{Name: activityName},
852 TaskQueue: &taskqueuepb.TaskQueue{Name: tl, Kind: enumspb.TASK_QUEUE_KIND_NORMAL},
853 Input: payloads.EncodeBytes(buf.Bytes()),
854 ScheduleToCloseTimeout: durationpb.New(15 * time.Second),
855 ScheduleToStartTimeout: durationpb.New(1 * time.Second),
856 StartToCloseTimeout: durationpb.New(15 * time.Second),
857 HeartbeatTimeout: durationpb.New(1 * time.Second),
858 }},
859 }}, nil
860 }
861
862 workflowComplete = true
863 return []*commandpb.Command{{
864 CommandType: enumspb.COMMAND_TYPE_COMPLETE_WORKFLOW_EXECUTION,
865 Attributes: &commandpb.Command_CompleteWorkflowExecutionCommandAttributes{CompleteWorkflowExecutionCommandAttributes: &commandpb.CompleteWorkflowExecutionCommandAttributes{
866 Result: payloads.EncodeString("Done"),
867 }},
868 }}, nil
869 }
870
871 activityExecutedCount := 0
872 atHandler := func(task *workflowservice.PollActivityTaskQueueResponse) (*commonpb.Payloads, bool, error) {
873 s.Equal(id, task.WorkflowExecution.GetWorkflowId())
874 s.Equal(activityName, task.ActivityType.GetName())
875 // Timing out more than HB time.
876 time.Sleep(2 * time.Second) //nolint:forbidigo
877 activityExecutedCount++
878 return payloads.EncodeString("Activity Result"), false, nil
879 }
880
881 poller := &testcore.TaskPoller{
882 Client: env.FrontendClient(),
883 Namespace: env.Namespace().String(),
884 TaskQueue: taskQueue,
885 Identity: identity,
886 WorkflowTaskHandler: wtHandler,
887 ActivityTaskHandler: atHandler,
888 Logger: env.Logger,
889 T: s.T(),
890 }
891
892 _, err := poller.PollAndProcessWorkflowTask()
893 s.True(err == nil || errors.Is(err, testcore.ErrNoTasks))
894
895 err = poller.PollAndProcessActivityTask(false)
896 // Not s.ErrorIs() because error goes through RPC.
897 s.IsType(consts.ErrActivityTaskNotFound, err)
898 s.Equal(consts.ErrActivityTaskNotFound.Error(), err.Error())
899
900 env.Logger.Info("Waiting for workflow to complete", tag.WorkflowRunID(we.RunId))
901
902 s.False(workflowComplete)
903 _, err = poller.PollAndProcessWorkflowTask(testcore.WithDumpHistory)
904 s.NoError(err)
905 s.True(workflowComplete)
906 }
907
908 func (s *ActivityTestSuite) TestTryActivityCancellationFromWorkflow() {
909 env := testcore.NewEnv(s.T())
910
911 id := "functional-activity-cancellation-test"
912 wt := "functional-activity-cancellation-test-type"
913 tl := "functional-activity-cancellation-test-taskqueue"
914 identity := "worker1"
915 activityName := "activity_timer"
916
917 workflowType := &commonpb.WorkflowType{Name: wt}
918
919 taskQueue := &taskqueuepb.TaskQueue{Name: tl, Kind: enumspb.TASK_QUEUE_KIND_NORMAL}
920
921 request := &workflowservice.StartWorkflowExecutionRequest{
922 RequestId: uuid.NewString(),
923 Namespace: env.Namespace().String(),
924 WorkflowId: id,
925 WorkflowType: workflowType,
926 TaskQueue: taskQueue,
927 Input: nil,
928 WorkflowRunTimeout: durationpb.New(100 * time.Second),
929 WorkflowTaskTimeout: durationpb.New(1 * time.Second),
930 Identity: identity,
931 }
932
933 we, err0 := env.FrontendClient().StartWorkflowExecution(s.Context(), request)
934 s.NoError(err0)
935
936 env.Logger.Info("StartWorkflowExecution: response", tag.WorkflowRunID(we.GetRunId()))
937
938 activityCounter := int32(0)
939 scheduleActivity := true
940 requestCancellation := false
941 activityScheduledID := int64(0)
942
943 wtHandler := func(task *workflowservice.PollWorkflowTaskQueueResponse) ([]*commandpb.Command, error) {
944 if scheduleActivity {
945 activityCounter++
946 buf := new(bytes.Buffer)
947 s.NoError(binary.Write(buf, binary.LittleEndian, activityCounter))
948
949 activityScheduledID = task.StartedEventId + 2
950 return []*commandpb.Command{{
951 CommandType: enumspb.COMMAND_TYPE_SCHEDULE_ACTIVITY_TASK,
952 Attributes: &commandpb.Command_ScheduleActivityTaskCommandAttributes{ScheduleActivityTaskCommandAttributes: &commandpb.ScheduleActivityTaskCommandAttributes{
953 ActivityId: convert.Int32ToString(activityCounter),
954 ActivityType: &commonpb.ActivityType{Name: activityName},
955 TaskQueue: &taskqueuepb.TaskQueue{Name: tl, Kind: enumspb.TASK_QUEUE_KIND_NORMAL},
956 Input: payloads.EncodeBytes(buf.Bytes()),
957 ScheduleToCloseTimeout: durationpb.New(15 * time.Second),
958 ScheduleToStartTimeout: durationpb.New(10 * time.Second),
959 StartToCloseTimeout: durationpb.New(15 * time.Second),
960 HeartbeatTimeout: durationpb.New(0 * time.Second),
961 }},
962 }}, nil
963 }
964
965 if requestCancellation {
966 return []*commandpb.Command{{
967 CommandType: enumspb.COMMAND_TYPE_REQUEST_CANCEL_ACTIVITY_TASK,
968 Attributes: &commandpb.Command_RequestCancelActivityTaskCommandAttributes{RequestCancelActivityTaskCommandAttributes: &commandpb.RequestCancelActivityTaskCommandAttributes{
969 ScheduledEventId: activityScheduledID,
970 }},
971 }}, nil
972 }
973
974 env.Logger.Info("Completing Workflow")
975
976 return []*commandpb.Command{{
977 CommandType: enumspb.COMMAND_TYPE_COMPLETE_WORKFLOW_EXECUTION,
978 Attributes: &commandpb.Command_CompleteWorkflowExecutionCommandAttributes{CompleteWorkflowExecutionCommandAttributes: &commandpb.CompleteWorkflowExecutionCommandAttributes{
979 Result: payloads.EncodeString("Done"),
980 }},
981 }}, nil
982 }
983
984 activityCanceled := false
985 atHandler := func(task *workflowservice.PollActivityTaskQueueResponse) (*commonpb.Payloads, bool, error) {
986 s.Equal(id, task.WorkflowExecution.GetWorkflowId())
987 s.Equal(activityName, task.ActivityType.GetName())
988 for i := range 10 {
989 env.Logger.Info("Heartbeating for activity", tag.ActivityID(task.ActivityId), tag.Counter(i))
990 response, err := env.FrontendClient().RecordActivityTaskHeartbeat(s.Context(),
991 &workflowservice.RecordActivityTaskHeartbeatRequest{
992 Namespace: env.Namespace().String(),
993 TaskToken: task.TaskToken,
994 Details: payloads.EncodeString("details"),
995 })
996 if response != nil && response.CancelRequested {
997 activityCanceled = true
998 return payloads.EncodeString("Activity Cancelled"), true, nil
999 }
1000 s.NoError(err)
1001 time.Sleep(10 * time.Millisecond) //nolint:forbidigo
1002 }
1003 return payloads.EncodeString("Activity Result"), false, nil
1004 }
1005
1006 poller := &testcore.TaskPoller{
1007 Client: env.FrontendClient(),
1008 Namespace: env.Namespace().String(),
1009 TaskQueue: taskQueue,
1010 Identity: identity,
1011 WorkflowTaskHandler: wtHandler,
1012 ActivityTaskHandler: atHandler,
1013 Logger: env.Logger,
1014 T: s.T(),
1015 }
1016
1017 _, err := poller.PollAndProcessWorkflowTask()
1018 s.True(err == nil || errors.Is(err, testcore.ErrNoTasks))
1019
1020 cancelCh := make(chan struct{})
1021 go func() {
1022 env.Logger.Info("Trying to cancel the task in a different thread")
1023 // Send signal so that worker can send an activity cancel
1024 _, err1 := env.FrontendClient().SignalWorkflowExecution(s.Context(), &workflowservice.SignalWorkflowExecutionRequest{
1025 Namespace: env.Namespace().String(),
1026 WorkflowExecution: &commonpb.WorkflowExecution{
1027 WorkflowId: id,
1028 RunId: we.RunId,
1029 },
1030 SignalName: "my signal",
1031 Input: nil,
1032 Identity: identity,
1033 })
1034 s.NoError(err1)
1035
1036 scheduleActivity = false
1037 requestCancellation = true
1038 _, err2 := poller.PollAndProcessWorkflowTask()
1039 s.NoError(err2)
1040 close(cancelCh)
1041 }()
1042
1043 env.Logger.Info("Start activity.")
1044 err = poller.PollAndProcessActivityTask(false)
1045 s.True(err == nil || errors.Is(err, testcore.ErrNoTasks))
1046
1047 env.Logger.Info("Waiting for cancel to complete.", tag.WorkflowRunID(we.RunId))
1048 select {
1049 case <-cancelCh:
1050 case <-s.Context().Done():
1051 s.Fail("Test timed out for activity cancellation", s.Context().Err())
1052 return
1053 }
1054 s.True(activityCanceled, "Activity was not cancelled.")
1055 env.Logger.Info("Activity cancelled.", tag.WorkflowRunID(we.RunId))
1056 }
1057
1058 func (s *ActivityTestSuite) TestActivityCancellationNotStarted() {
1059 env := testcore.NewEnv(s.T())
1060 id := "functional-activity-notstarted-cancellation-test"
1061 wt := "functional-activity-notstarted-cancellation-test-type"
1062 tl := "functional-activity-notstarted-cancellation-test-taskqueue"
1063 identity := "worker1"
1064 activityName := "activity_notstarted"
1065
1066 workflowType := &commonpb.WorkflowType{Name: wt}
1067
1068 taskQueue := &taskqueuepb.TaskQueue{Name: tl, Kind: enumspb.TASK_QUEUE_KIND_NORMAL}
1069
1070 request := &workflowservice.StartWorkflowExecutionRequest{
1071 RequestId: uuid.NewString(),
1072 Namespace: env.Namespace().String(),
1073 WorkflowId: id,
1074 WorkflowType: workflowType,
1075 TaskQueue: taskQueue,
1076 Input: nil,
1077 WorkflowRunTimeout: durationpb.New(100 * time.Second),
1078 WorkflowTaskTimeout: durationpb.New(1 * time.Second),
1079 Identity: identity,
1080 }
1081
1082 we, err0 := env.FrontendClient().StartWorkflowExecution(s.Context(), request)
1083 s.NoError(err0)
1084
1085 env.Logger.Info("StartWorkflowExecutionn", tag.WorkflowRunID(we.GetRunId()))
1086
1087 activityCounter := int32(0)
1088 scheduleActivity := true
1089 requestCancellation := false
1090 activityScheduledID := int64(0)
1091
1092 wtHandler := func(task *workflowservice.PollWorkflowTaskQueueResponse) ([]*commandpb.Command, error) {
1093 if scheduleActivity {
1094 activityCounter++
1095 buf := new(bytes.Buffer)
1096 s.NoError(binary.Write(buf, binary.LittleEndian, activityCounter))
1097 env.Logger.Info("Scheduling activity")
1098 activityScheduledID = task.StartedEventId + 2
1099 return []*commandpb.Command{{
1100 CommandType: enumspb.COMMAND_TYPE_SCHEDULE_ACTIVITY_TASK,
1101 Attributes: &commandpb.Command_ScheduleActivityTaskCommandAttributes{ScheduleActivityTaskCommandAttributes: &commandpb.ScheduleActivityTaskCommandAttributes{
1102 ActivityId: convert.Int32ToString(activityCounter),
1103 ActivityType: &commonpb.ActivityType{Name: activityName},
1104 TaskQueue: &taskqueuepb.TaskQueue{Name: tl, Kind: enumspb.TASK_QUEUE_KIND_NORMAL},
1105 Input: payloads.EncodeBytes(buf.Bytes()),
1106 ScheduleToCloseTimeout: durationpb.New(15 * time.Second),
1107 ScheduleToStartTimeout: durationpb.New(2 * time.Second),
1108 StartToCloseTimeout: durationpb.New(15 * time.Second),
1109 HeartbeatTimeout: durationpb.New(0 * time.Second),
1110 }},
1111 }}, nil
1112 }
1113
1114 if requestCancellation {
1115 env.Logger.Info("Requesting cancellation")
1116 return []*commandpb.Command{{
1117 CommandType: enumspb.COMMAND_TYPE_REQUEST_CANCEL_ACTIVITY_TASK,
1118 Attributes: &commandpb.Command_RequestCancelActivityTaskCommandAttributes{RequestCancelActivityTaskCommandAttributes: &commandpb.RequestCancelActivityTaskCommandAttributes{
1119 ScheduledEventId: activityScheduledID,
1120 }},
1121 }}, nil
1122 }
1123
1124 env.Logger.Info("Completing Workflow")
1125 return []*commandpb.Command{{
1126 CommandType: enumspb.COMMAND_TYPE_COMPLETE_WORKFLOW_EXECUTION,
1127 Attributes: &commandpb.Command_CompleteWorkflowExecutionCommandAttributes{CompleteWorkflowExecutionCommandAttributes: &commandpb.CompleteWorkflowExecutionCommandAttributes{
1128 Result: payloads.EncodeString("Done"),
1129 }},
1130 }}, nil
1131 }
1132
1133 // dummy activity handler
1134 atHandler := func(task *workflowservice.PollActivityTaskQueueResponse) (*commonpb.Payloads, bool, error) {
1135 s.Fail("activity should not run")
1136 return nil, false, nil
1137 }
1138
1139 poller := &testcore.TaskPoller{
1140 Client: env.FrontendClient(),
1141 Namespace: env.Namespace().String(),
1142 TaskQueue: taskQueue,
1143 Identity: identity,
1144 WorkflowTaskHandler: wtHandler,
1145 ActivityTaskHandler: atHandler,
1146 Logger: env.Logger,
1147 T: s.T(),
1148 }
1149
1150 _, err := poller.PollAndProcessWorkflowTask()
1151 s.True(err == nil || errors.Is(err, testcore.ErrNoTasks))
1152
1153 // Send signal so that worker can send an activity cancel
1154 signalName := "my signal"
1155 signalInput := payloads.EncodeString("my signal input")
1156 _, err = env.FrontendClient().SignalWorkflowExecution(s.Context(), &workflowservice.SignalWorkflowExecutionRequest{
1157 Namespace: env.Namespace().String(),
1158 WorkflowExecution: &commonpb.WorkflowExecution{
1159 WorkflowId: id,
1160 RunId: we.RunId,
1161 },
1162 SignalName: signalName,
1163 Input: signalInput,
1164 Identity: identity,
1165 })
1166 s.NoError(err)
1167
1168 // Process signal in workflow and send request cancellation
1169 scheduleActivity = false
1170 requestCancellation = true
1171 _, err = poller.PollAndProcessWorkflowTask(testcore.WithDumpHistory)
1172 s.NoError(err)
1173
1174 scheduleActivity = false
1175 requestCancellation = false
1176 _, err = poller.PollAndProcessWorkflowTask()
1177 s.True(err == nil || errors.Is(err, testcore.ErrNoTasks))
1178 }
1179
1180 func (s *ActivityClientTestSuite) TestActivityHeartbeatDetailsDuringRetry() {
1181 env := testcore.NewEnv(s.T())
1182 // Latest reported heartbeat on activity should be available throughout workflow execution or until activity succeeds.
1183 // 1. Start workflow with single activity
1184 // 2. First invocation of activity sets heartbeat details and times out.
1185 // 3. Second invocation triggers retriable error.
1186 // 4. Next invocations succeed.
1187 // 5. Test should start polling for heartbeat details once first heartbeat was reported.
1188 // 6. Once workflow completes -- we're done.
1189
1190 activityTimeout := time.Second
1191
1192 activityExecutedCount := 0
1193 heartbeatDetails := 7771 // any value
1194 heartbeatDetailsPayload, err := payloads.Encode(heartbeatDetails)
1195 s.NoError(err)
1196 activityFn := func(ctx context.Context) error {
1197 var err error
1198 if activityExecutedCount == 0 {
1199 activity.RecordHeartbeat(ctx, heartbeatDetails)
1200 time.Sleep(activityTimeout + time.Second) //nolint:forbidigo
1201 } else if activityExecutedCount == 1 {
1202 time.Sleep(activityTimeout / 2) //nolint:forbidigo
1203 err = errors.New("retryable-error") //nolint:err113
1204 }
1205
1206 if activityExecutedCount > 0 {
1207 s.True(activity.HasHeartbeatDetails(ctx))
1208 var details int
1209 s.NoError(activity.GetHeartbeatDetails(ctx, &details))
1210 s.Equal(details, heartbeatDetails)
1211 }
1212
1213 activityExecutedCount++
1214 return err
1215 }
1216
1217 var err1 error
1218
1219 activityId := "heartbeat_retry"
1220 workflowFn := func(ctx workflow.Context) error {
1221 activityRetryPolicy := &temporal.RetryPolicy{
1222 InitialInterval: time.Second * 2,
1223 BackoffCoefficient: 1,
1224 MaximumInterval: time.Second * 2,
1225 MaximumAttempts: 3,
1226 }
1227
1228 ctx1 := workflow.WithActivityOptions(ctx, workflow.ActivityOptions{
1229 ActivityID: activityId,
1230 ScheduleToStartTimeout: 2 * time.Second,
1231 StartToCloseTimeout: 2 * time.Second,
1232 RetryPolicy: activityRetryPolicy,
1233 })
1234 f1 := workflow.ExecuteActivity(ctx1, activityFn)
1235
1236 err1 = f1.Get(ctx1, nil)
1237
1238 return nil
1239 }
1240
1241 env.SdkWorker().RegisterActivity(activityFn)
1242 env.SdkWorker().RegisterWorkflow(workflowFn)
1243
1244 wfId := "functional-test-heartbeat-details-during-retry"
1245 workflowOptions := sdkclient.StartWorkflowOptions{
1246 ID: wfId,
1247 TaskQueue: env.WorkerTaskQueue(),
1248 WorkflowRunTimeout: 20 * time.Second,
1249 }
1250 workflowRun, err := env.SdkClient().ExecuteWorkflow(s.Context(), workflowOptions, workflowFn)
1251 s.NoError(err)
1252
1253 s.NotNil(workflowRun)
1254 s.NotEmpty(workflowRun.GetRunID())
1255
1256 runId := workflowRun.GetRunID()
1257
1258 describeWorkflowExecution := func() (*workflowservice.DescribeWorkflowExecutionResponse, error) {
1259 return env.FrontendClient().DescribeWorkflowExecution(s.Context(), &workflowservice.DescribeWorkflowExecutionRequest{
1260 Namespace: env.Namespace().String(),
1261 Execution: &commonpb.WorkflowExecution{
1262 WorkflowId: wfId,
1263 RunId: runId,
1264 },
1265 })
1266 }
1267
1268 //nolint:forbidigo
1269 time.Sleep(time.Second) // wait for the timeout to trigger
1270
1271 for dweResult, dweErr := describeWorkflowExecution(); dweResult.GetWorkflowExecutionInfo().GetCloseTime() == nil; dweResult, dweErr = describeWorkflowExecution() {
1272 s.NoError(dweErr)
1273 s.NotNil(dweResult.GetWorkflowExecutionInfo())
1274 s.LessOrEqual(len(dweResult.PendingActivities), 1)
1275
1276 if dweResult.PendingActivities != nil && len(dweResult.PendingActivities) == 1 {
1277 details := dweResult.PendingActivities[0].GetHeartbeatDetails()
1278 s.Equal(heartbeatDetailsPayload, details)
1279 }
1280
1281 time.Sleep(time.Millisecond * 100) //nolint:forbidigo
1282 }
1283
1284 err = workflowRun.Get(s.Context(), nil)
1285 s.NoError(err)
1286
1287 s.NoError(err1)
1288 }
1289
1290 // TestActivityHeartBeat_RecordIdentity verifies that the identity of the worker sending the heartbeat
1291 // is recorded in pending activity info and returned in describe workflow API response. This happens
1292 // only when the worker identity is not sent when a poller picks the task.
1293 func (s *ActivityTestSuite) TestActivityHeartBeat_RecordIdentity() {
1294 env := testcore.NewEnv(s.T())
1295 id := "functional-heartbeat-identity-record"
1296 workerIdentity := "70df788a-b0b2-4113-a0d5-130f13889e35"
1297 activityName := "activity_timer"
1298
1299 taskQueue := &taskqueuepb.TaskQueue{Name: "functional-heartbeat-identity-record-taskqueue", Kind: enumspb.TASK_QUEUE_KIND_NORMAL}
1300
1301 header := &commonpb.Header{
1302 Fields: map[string]*commonpb.Payload{"tracing": payload.EncodeString("sample data")},
1303 }
1304
1305 request := &workflowservice.StartWorkflowExecutionRequest{
1306 RequestId: uuid.NewString(),
1307 Namespace: env.Namespace().String(),
1308 WorkflowId: id,
1309 WorkflowType: &commonpb.WorkflowType{Name: "functional-heartbeat-identity-record-type"},
1310 TaskQueue: taskQueue,
1311 Input: nil,
1312 Header: header,
1313 WorkflowRunTimeout: durationpb.New(100 * time.Second),
1314 WorkflowTaskTimeout: durationpb.New(60 * time.Second),
1315 Identity: workerIdentity,
1316 }
1317
1318 we, err := env.FrontendClient().StartWorkflowExecution(s.Context(), request)
1319 s.NoError(err)
1320
1321 workflowComplete := false
1322 workflowNextCmd := enumspb.COMMAND_TYPE_SCHEDULE_ACTIVITY_TASK
1323 wtHandler := func(task *workflowservice.PollWorkflowTaskQueueResponse) ([]*commandpb.Command, error) {
1324 switch workflowNextCmd { // nolint:exhaustive
1325 case enumspb.COMMAND_TYPE_SCHEDULE_ACTIVITY_TASK:
1326 workflowNextCmd = enumspb.COMMAND_TYPE_COMPLETE_WORKFLOW_EXECUTION
1327 return []*commandpb.Command{{
1328 CommandType: enumspb.COMMAND_TYPE_SCHEDULE_ACTIVITY_TASK,
1329 Attributes: &commandpb.Command_ScheduleActivityTaskCommandAttributes{ScheduleActivityTaskCommandAttributes: &commandpb.ScheduleActivityTaskCommandAttributes{
1330 ActivityId: convert.IntToString(rand.Int()),
1331 ActivityType: &commonpb.ActivityType{Name: activityName},
1332 TaskQueue: taskQueue,
1333 Input: payloads.EncodeString("activity-input"),
1334 Header: header,
1335 ScheduleToCloseTimeout: durationpb.New(15 * time.Second),
1336 ScheduleToStartTimeout: durationpb.New(60 * time.Second),
1337 StartToCloseTimeout: durationpb.New(15 * time.Second),
1338 HeartbeatTimeout: durationpb.New(60 * time.Second),
1339 }},
1340 }}, nil
1341 case enumspb.COMMAND_TYPE_COMPLETE_WORKFLOW_EXECUTION:
1342 workflowComplete = true
1343 return []*commandpb.Command{{
1344 CommandType: enumspb.COMMAND_TYPE_COMPLETE_WORKFLOW_EXECUTION,
1345 Attributes: &commandpb.Command_CompleteWorkflowExecutionCommandAttributes{CompleteWorkflowExecutionCommandAttributes: &commandpb.CompleteWorkflowExecutionCommandAttributes{
1346 Result: payloads.EncodeString("Done"),
1347 }},
1348 }}, nil
1349 }
1350 panic("Unexpected workflow state")
1351 }
1352
1353 activityStartedSignal := make(chan bool) // Used by activity channel to signal the start so that the test can verify empty identity.
1354 heartbeatSignalChan := make(chan bool) // Activity task sends heartbeat when signaled on this chan. It also signals back on the same chan after sending the heartbeat.
1355 endActivityTask := make(chan bool) // Activity task completes when signaled on this chan. This is to force the task to be in pending state.
1356 atHandler := func(task *workflowservice.PollActivityTaskQueueResponse) (*commonpb.Payloads, bool, error) {
1357 activityStartedSignal <- true // signal the start of activity task.
1358 <-heartbeatSignalChan // wait for signal before sending heartbeat.
1359 _, err := env.FrontendClient().RecordActivityTaskHeartbeat(s.Context(), &workflowservice.RecordActivityTaskHeartbeatRequest{
1360 Namespace: env.Namespace().String(),
1361 TaskToken: task.TaskToken,
1362 Details: payloads.EncodeString("details"),
1363 Identity: workerIdentity, // explicitly set the worker identity in the heartbeat request
1364 })
1365 s.NoError(err)
1366 heartbeatSignalChan <- true // signal that the heartbeat is sent.
1367
1368 <-endActivityTask // wait for signal before completing the task
1369 return payloads.EncodeString("Activity Result"), false, nil
1370 }
1371
1372 poller := &testcore.TaskPoller{
1373 Client: env.FrontendClient(),
1374 Namespace: env.Namespace().String(),
1375 TaskQueue: taskQueue,
1376 Identity: "", // Do not send the worker identity.
1377 WorkflowTaskHandler: wtHandler,
1378 ActivityTaskHandler: atHandler,
1379 Logger: env.Logger,
1380 T: s.T(),
1381 }
1382
1383 // execute workflow task so that an activity can be enqueued.
1384 _, err = poller.PollAndProcessWorkflowTask()
1385 s.True(err == nil || errors.Is(err, testcore.ErrNoTasks))
1386
1387 // execute activity task which waits for signal before sending heartbeat.
1388 go func() {
1389 err := poller.PollAndProcessActivityTask(false)
1390 s.True(err == nil || errors.Is(err, testcore.ErrNoTasks))
1391 }()
1392
1393 describeWorkflowExecution := func() (*workflowservice.DescribeWorkflowExecutionResponse, error) {
1394 return env.FrontendClient().DescribeWorkflowExecution(s.Context(), &workflowservice.DescribeWorkflowExecutionRequest{
1395 Namespace: env.Namespace().String(),
1396 Execution: &commonpb.WorkflowExecution{
1397 WorkflowId: id,
1398 RunId: we.RunId,
1399 },
1400 })
1401 }
1402 <-activityStartedSignal // wait for the activity to start
1403
1404 // Verify that the worker identity is empty.
1405 descRespBeforeHeartbeat, err := describeWorkflowExecution()
1406 s.NoError(err)
1407 s.Empty(descRespBeforeHeartbeat.PendingActivities[0].LastWorkerIdentity)
1408
1409 heartbeatSignalChan <- true // ask the activity to send a heartbeat.
1410 <-heartbeatSignalChan // wait for the heartbeat to be sent (to prevent the test from racing to describe the workflow before the heartbeat is sent)
1411
1412 // Verify that the worker identity is set now.
1413 descRespAfterHeartbeat, err := describeWorkflowExecution()
1414 s.NoError(err)
1415 s.Equal(workerIdentity, descRespAfterHeartbeat.PendingActivities[0].LastWorkerIdentity)
1416
1417 // unblock the activity task
1418 endActivityTask <- true
1419
1420 // ensure that the workflow is complete.
1421 _, err = poller.PollAndProcessWorkflowTask(testcore.WithDumpHistory)
1422 s.NoError(err)
1423 s.True(workflowComplete)
1424 }
1425
1426 func (s *ActivityTestSuite) TestActivityTaskCompleteForceCompletion() {
1427 env := testcore.NewEnv(s.T())
1428 activityInfo := make(chan activity.Info, 1)
1429 taskQueue := testcore.RandomizeStr(s.T().Name())
1430 w, wf := s.mockWorkflowWithErrorActivity(activityInfo, env.SdkClient(), taskQueue)
1431 s.NoError(w.Start())
1432 defer w.Stop()
1433
1434 workflowOptions := sdkclient.StartWorkflowOptions{
1435 ID: uuid.NewString(),
1436 TaskQueue: taskQueue,
1437 }
1438 run, err := env.SdkClient().ExecuteWorkflow(s.Context(), workflowOptions, wf)
1439 s.NoError(err)
1440 ai := <-activityInfo
1441 s.EventuallyWithT(func(t *assert.CollectT) {
1442 description, err := env.SdkClient().DescribeWorkflowExecution(s.Context(), run.GetID(), run.GetRunID())
1443 require.NoError(t, err)
1444 require.Len(t, description.PendingActivities, 1)
1445 require.NotNil(t, description.PendingActivities[0].LastFailure)
1446 require.Equal(t, "mock error of an activity", description.PendingActivities[0].LastFailure.Message)
1447 },
1448 10*time.Second,
1449 500*time.Millisecond)
1450
1451 err = env.SdkClient().CompleteActivityByID(s.Context(), env.Namespace().String(), run.GetID(), run.GetRunID(), ai.ActivityID, nil, nil)
1452 s.NoError(err)
1453
1454 // Ensure the activity is completed and the workflow is unblcked.
1455 s.NoError(run.Get(s.Context(), nil))
1456 }
1457
1458 func (s *ActivityTestSuite) TestActivityTaskCompleteRejectCompletion() {
1459 env := testcore.NewEnv(s.T())
1460 activityInfo := make(chan activity.Info, 1)
1461 taskQueue := testcore.RandomizeStr(s.T().Name())
1462 w, wf := s.mockWorkflowWithErrorActivity(activityInfo, env.SdkClient(), taskQueue)
1463 s.NoError(w.Start())
1464 defer w.Stop()
1465
1466 workflowOptions := sdkclient.StartWorkflowOptions{
1467 ID: uuid.NewString(),
1468 TaskQueue: taskQueue,
1469 }
1470 run, err := env.SdkClient().ExecuteWorkflow(s.Context(), workflowOptions, wf)
1471 s.NoError(err)
1472 ai := <-activityInfo
1473 s.EventuallyWithT(func(t *assert.CollectT) {
1474 description, err := env.SdkClient().DescribeWorkflowExecution(s.Context(), run.GetID(), run.GetRunID())
1475 require.NoError(t, err)
1476 require.Len(t, description.PendingActivities, 1)
1477 require.NotNil(t, description.PendingActivities[0].LastFailure)
1478 require.Equal(t, "mock error of an activity", description.PendingActivities[0].LastFailure.Message)
1479 },
1480 10*time.Second,
1481 500*time.Millisecond)
1482
1483 err = env.SdkClient().CompleteActivity(s.Context(), ai.TaskToken, nil, nil)
1484 var svcErr *serviceerror.NotFound
1485 s.ErrorAs(err, &svcErr, "invalid activityID or activity already timed out or invoking workflow is completed")
1486 }
1487
1488 func (s *ActivityTestSuite) mockWorkflowWithErrorActivity(activityInfo chan<- activity.Info, sdkClient sdkclient.Client, taskQueue string) (worker.Worker, func(ctx workflow.Context) error) {
1489 mockErrorActivity := func(ctx context.Context) error {
1490 ai := activity.GetInfo(ctx)
1491 activityInfo <- ai
1492 return errors.New("mock error of an activity") //nolint:err113
1493 }
1494 wf := func(ctx workflow.Context) error {
1495 ao := workflow.ActivityOptions{
1496 StartToCloseTimeout: 3 * time.Minute,
1497 RetryPolicy: &temporal.RetryPolicy{
1498 // Add long initial interval to make sure the next attempt is not scheduled
1499 // before the test gets a chance to complete the activity via API call.
1500 InitialInterval: 2 * time.Minute,
1501 },
1502 }
1503 ctx2 := workflow.WithActivityOptions(ctx, ao)
1504 var mockErrorResult error
1505 err := workflow.ExecuteActivity(ctx2, mockErrorActivity).Get(ctx2, &mockErrorResult)
1506 if err != nil {
1507 return err
1508 }
1509 return mockErrorResult
1510 }
1511
1512 workflowType := "test"
1513 w := worker.New(sdkClient, taskQueue, worker.Options{})
1514 w.RegisterWorkflowWithOptions(wf, workflow.RegisterOptions{Name: workflowType})
1515 w.RegisterActivity(mockErrorActivity)
1516 return w, wf
1517 }
1518
1519 func (s *ActivityClientTestSuite) TestActivity_AttemptsExceeded() {
1520 env := testcore.NewEnv(s.T())
1521 activityFunction := func(ctx context.Context) error {
1522 return errors.New("non-retryable-error") //nolint:err113
1523 }
1524
1525 workflowFn := func(ctx workflow.Context) error {
1526 err := workflow.ExecuteActivity(workflow.WithActivityOptions(ctx, workflow.ActivityOptions{
1527 StartToCloseTimeout: 10 * time.Second,
1528 RetryPolicy: &temporal.RetryPolicy{
1529 MaximumAttempts: 1,
1530 },
1531 }), activityFunction).Get(ctx, nil)
1532 return err
1533 }
1534
1535 env.SdkWorker().RegisterWorkflow(workflowFn)
1536 env.SdkWorker().RegisterActivity(activityFunction)
1537
1538 wfID := testcore.RandomizeStr(s.T().Name())
1539 workflowOptions := sdkclient.StartWorkflowOptions{
1540 ID: wfID,
1541 TaskQueue: env.WorkerTaskQueue(),
1542 }
1543 workflowRun, err := env.SdkClient().ExecuteWorkflow(s.Context(), workflowOptions, workflowFn)
1544 s.NoError(err)
1545
1546 err = workflowRun.Get(s.Context(), nil)
1547 var wfExecutionError *temporal.WorkflowExecutionError
1548 s.ErrorAs(err, &wfExecutionError)
1549 var activityError *temporal.ActivityError
1550 s.ErrorAs(wfExecutionError, &activityError)
1551 s.Equal(enumspb.RETRY_STATE_MAXIMUM_ATTEMPTS_REACHED, activityError.RetryState())
1552 var applicationErr *temporal.ApplicationError
1553 s.ErrorAs(activityError, &applicationErr)
1554 s.Equal("non-retryable-error", applicationErr.Message())
1555
1556 history := env.GetHistory(string(env.Namespace()), &commonpb.WorkflowExecution{WorkflowId: workflowRun.GetID()})
1557 s.ContainsHistory(`ActivityTaskFailed`, &historypb.History{Events: history})
1558 }