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
}