go.temporal.io/server/tests/describe_test.go
309 LOC · 0 covered · 309 uncovered · 0 ranges · 0 concepts · 0 introducers · 0 tests
1
package tests
2
3
import (
4
"bytes"
5
"encoding/binary"
6
"testing"
7
"time"
8
9
"github.com/google/uuid"
10
commandpb "go.temporal.io/api/command/v1"
11
commonpb "go.temporal.io/api/common/v1"
12
enumspb "go.temporal.io/api/enums/v1"
13
taskqueuepb "go.temporal.io/api/taskqueue/v1"
14
workflowpb "go.temporal.io/api/workflow/v1"
15
"go.temporal.io/api/workflowservice/v1"
16
"go.temporal.io/server/common"
17
"go.temporal.io/server/common/convert"
18
"go.temporal.io/server/common/log/tag"
19
"go.temporal.io/server/common/payloads"
20
"go.temporal.io/server/common/primitives/timestamp"
21
"go.temporal.io/server/common/testing/parallelsuite"
22
"go.temporal.io/server/tests/testcore"
23
"google.golang.org/protobuf/types/known/durationpb"
24
)
25
26
type DescribeTestSuite struct {
27
parallelsuite.Suite[*DescribeTestSuite]
28
}
29
30
func TestDescribeTestSuite(t *testing.T) {
31
parallelsuite.Run(t, &DescribeTestSuite{})
32
}
33
34
func (s *DescribeTestSuite) TestDescribeWorkflowExecution() {
35
env := testcore.NewEnv(s.T())
36
id := "functional-describe-wfe-test"
37
wt := "functional-describe-wfe-test-type"
38
tq := "functional-describe-wfe-test-taskqueue"
39
identity := "worker1"
40
41
// Start workflow execution
42
requestID := uuid.NewString()
43
request := &workflowservice.StartWorkflowExecutionRequest{
44
RequestId: requestID,
45
Namespace: env.Namespace().String(),
46
WorkflowId: id,
47
WorkflowType: &commonpb.WorkflowType{Name: wt},
48
TaskQueue: &taskqueuepb.TaskQueue{Name: tq, Kind: enumspb.TASK_QUEUE_KIND_NORMAL},
49
Input: nil,
50
WorkflowRunTimeout: durationpb.New(100 * time.Second),
51
WorkflowTaskTimeout: durationpb.New(1 * time.Second),
52
Identity: identity,
53
}
54
55
we, err0 := env.FrontendClient().StartWorkflowExecution(s.Context(), request)
56
s.NoError(err0)
57
58
env.Logger.Info("StartWorkflowExecution", tag.WorkflowRunID(we.RunId))
59
60
describeWorkflowExecution := func() (*workflowservice.DescribeWorkflowExecutionResponse, error) {
61
return env.FrontendClient().DescribeWorkflowExecution(s.Context(), &workflowservice.DescribeWorkflowExecutionRequest{
62
Namespace: env.Namespace().String(),
63
Execution: &commonpb.WorkflowExecution{
64
WorkflowId: id,
65
RunId: we.RunId,
66
},
67
})
68
}
69
dweResponse, err := describeWorkflowExecution()
70
s.NoError(err)
71
wfInfo := dweResponse.WorkflowExecutionInfo
72
s.Nil(wfInfo.CloseTime)
73
s.Nil(wfInfo.ExecutionDuration)
74
s.Equal(int64(2), wfInfo.HistoryLength) // WorkflowStarted, WorkflowTaskScheduled
75
s.Equal(wfInfo.GetStartTime(), wfInfo.GetExecutionTime())
76
s.Equal(tq, wfInfo.TaskQueue)
77
s.Positive(wfInfo.GetHistorySizeBytes())
78
s.Empty(wfInfo.GetParentNamespaceId())
79
s.Nil(wfInfo.GetParentExecution())
80
s.NotNil(wfInfo.GetRootExecution())
81
s.Equal(id, wfInfo.RootExecution.GetWorkflowId())
82
s.Equal(we.RunId, wfInfo.RootExecution.GetRunId())
83
s.Equal(we.RunId, wfInfo.GetFirstRunId())
84
s.NotNil(dweResponse.WorkflowExtendedInfo)
85
s.Nil(dweResponse.WorkflowExtendedInfo.LastResetTime) // workflow was not reset
86
s.Nil(dweResponse.WorkflowExtendedInfo.ExecutionExpirationTime)
87
s.NotNil(dweResponse.WorkflowExtendedInfo.RunExpirationTime)
88
s.NotNil(dweResponse.WorkflowExtendedInfo.OriginalStartTime)
89
s.NotNil(dweResponse.WorkflowExtendedInfo.RequestIdInfos)
90
s.Contains(dweResponse.WorkflowExtendedInfo.RequestIdInfos, requestID)
91
s.NotNil(dweResponse.WorkflowExtendedInfo.RequestIdInfos[requestID])
92
s.ProtoEqual(
93
&workflowpb.RequestIdInfo{
94
EventType: enumspb.EVENT_TYPE_WORKFLOW_EXECUTION_STARTED,
95
EventId: common.FirstEventID,
96
Buffered: false,
97
},
98
dweResponse.WorkflowExtendedInfo.RequestIdInfos[requestID],
99
)
100
101
// workflow logic
102
workflowComplete := false
103
signalSent := false
104
wtHandler := func(task *workflowservice.PollWorkflowTaskQueueResponse) ([]*commandpb.Command, error) {
105
if !signalSent {
106
signalSent = true
107
108
s.NoError(err)
109
return []*commandpb.Command{{
110
CommandType: enumspb.COMMAND_TYPE_SCHEDULE_ACTIVITY_TASK,
111
Attributes: &commandpb.Command_ScheduleActivityTaskCommandAttributes{ScheduleActivityTaskCommandAttributes: &commandpb.ScheduleActivityTaskCommandAttributes{
112
ActivityId: "1",
113
ActivityType: &commonpb.ActivityType{Name: "test-activity-type"},
114
TaskQueue: &taskqueuepb.TaskQueue{Name: tq, Kind: enumspb.TASK_QUEUE_KIND_NORMAL},
115
Input: payloads.EncodeString("test-input"),
116
ScheduleToCloseTimeout: durationpb.New(100 * time.Second),
117
ScheduleToStartTimeout: durationpb.New(2 * time.Second),
118
StartToCloseTimeout: durationpb.New(50 * time.Second),
119
HeartbeatTimeout: durationpb.New(5 * time.Second),
120
}},
121
}}, nil
122
}
123
124
workflowComplete = true
125
return []*commandpb.Command{{
126
CommandType: enumspb.COMMAND_TYPE_COMPLETE_WORKFLOW_EXECUTION,
127
Attributes: &commandpb.Command_CompleteWorkflowExecutionCommandAttributes{CompleteWorkflowExecutionCommandAttributes: &commandpb.CompleteWorkflowExecutionCommandAttributes{
128
Result: payloads.EncodeString("Done"),
129
}},
130
}}, nil
131
}
132
133
atHandler := func(task *workflowservice.PollActivityTaskQueueResponse) (*commonpb.Payloads, bool, error) {
134
return payloads.EncodeString("Activity Result"), false, nil
135
}
136
137
poller := &testcore.TaskPoller{
138
Client: env.FrontendClient(),
139
Namespace: env.Namespace().String(),
140
TaskQueue: &taskqueuepb.TaskQueue{Name: tq, Kind: enumspb.TASK_QUEUE_KIND_NORMAL},
141
Identity: identity,
142
WorkflowTaskHandler: wtHandler,
143
ActivityTaskHandler: atHandler,
144
Logger: env.Logger,
145
T: s.T(),
146
}
147
148
// first workflow task to schedule new activity
149
_, err = poller.PollAndProcessWorkflowTask()
150
env.Logger.Info("PollAndProcessWorkflowTask", tag.Error(err))
151
s.NoError(err)
152
153
dweResponse, err = describeWorkflowExecution()
154
s.NoError(err)
155
wfInfo = dweResponse.WorkflowExecutionInfo
156
s.Equal(enumspb.WORKFLOW_EXECUTION_STATUS_RUNNING, wfInfo.GetStatus())
157
s.Nil(wfInfo.CloseTime)
158
s.Nil(wfInfo.ExecutionDuration)
159
s.Equal(int64(5), wfInfo.HistoryLength) // WorkflowTaskStarted, WorkflowTaskCompleted, ActivityScheduled
160
s.Len(dweResponse.PendingActivities, 1)
161
s.Equal("test-activity-type", dweResponse.PendingActivities[0].ActivityType.GetName())
162
s.True(timestamp.TimeValue(dweResponse.PendingActivities[0].GetLastHeartbeatTime()).IsZero())
163
164
// process activity task
165
err = poller.PollAndProcessActivityTask(false)
166
167
dweResponse, err = describeWorkflowExecution()
168
s.NoError(err)
169
wfInfo = dweResponse.WorkflowExecutionInfo
170
s.Equal(enumspb.WORKFLOW_EXECUTION_STATUS_RUNNING, wfInfo.GetStatus())
171
s.Equal(int64(8), wfInfo.HistoryLength) // ActivityTaskStarted, ActivityTaskCompleted, WorkflowTaskScheduled
172
s.Empty(dweResponse.PendingActivities)
173
174
// Process signal in workflow
175
_, err = poller.PollAndProcessWorkflowTask(testcore.WithDumpHistory)
176
s.NoError(err)
177
s.True(workflowComplete)
178
179
dweResponse, err = describeWorkflowExecution()
180
s.NoError(err)
181
wfInfo = dweResponse.WorkflowExecutionInfo
182
s.Equal(enumspb.WORKFLOW_EXECUTION_STATUS_COMPLETED, wfInfo.GetStatus())
183
s.NotNil(wfInfo.CloseTime)
184
s.NotNil(wfInfo.ExecutionDuration)
185
s.Equal(
186
wfInfo.GetCloseTime().AsTime().Sub(wfInfo.ExecutionTime.AsTime()),
187
wfInfo.ExecutionDuration.AsDuration(),
188
)
189
s.Equal(int64(11), wfInfo.HistoryLength) // WorkflowTaskStarted, WorkflowTaskCompleted, WorkflowCompleted
190
}
191
192
func (s *DescribeTestSuite) TestDescribeTaskQueue() {
193
env := testcore.NewEnv(s.T())
194
workflowID := "functional-get-poller-history"
195
wt := "functional-get-poller-history-type"
196
tl := "functional-get-poller-history-taskqueue"
197
identity := "worker1"
198
activityName := "activity_type1"
199
200
// Start workflow execution
201
request := &workflowservice.StartWorkflowExecutionRequest{
202
RequestId: uuid.NewString(),
203
Namespace: env.Namespace().String(),
204
WorkflowId: workflowID,
205
WorkflowType: &commonpb.WorkflowType{Name: wt},
206
TaskQueue: &taskqueuepb.TaskQueue{Name: tl, Kind: enumspb.TASK_QUEUE_KIND_NORMAL},
207
Input: nil,
208
WorkflowRunTimeout: durationpb.New(100 * time.Second),
209
WorkflowTaskTimeout: durationpb.New(1 * time.Second),
210
Identity: identity,
211
}
212
213
we, err0 := env.FrontendClient().StartWorkflowExecution(s.Context(), request)
214
s.NoError(err0)
215
216
env.Logger.Info("StartWorkflowExecution", tag.WorkflowRunID(we.RunId))
217
218
// workflow logic
219
activityScheduled := false
220
activityData := int32(1)
221
// var signalEvent *historypb.HistoryEvent
222
wtHandler := func(task *workflowservice.PollWorkflowTaskQueueResponse) ([]*commandpb.Command, error) {
223
if !activityScheduled {
224
activityScheduled = true
225
buf := new(bytes.Buffer)
226
s.NoError(binary.Write(buf, binary.LittleEndian, activityData))
227
228
return []*commandpb.Command{{
229
CommandType: enumspb.COMMAND_TYPE_SCHEDULE_ACTIVITY_TASK,
230
Attributes: &commandpb.Command_ScheduleActivityTaskCommandAttributes{ScheduleActivityTaskCommandAttributes: &commandpb.ScheduleActivityTaskCommandAttributes{
231
ActivityId: convert.Int32ToString(1),
232
ActivityType: &commonpb.ActivityType{Name: activityName},
233
TaskQueue: &taskqueuepb.TaskQueue{Name: tl, Kind: enumspb.TASK_QUEUE_KIND_NORMAL},
234
Input: payloads.EncodeBytes(buf.Bytes()),
235
ScheduleToCloseTimeout: durationpb.New(100 * time.Second),
236
ScheduleToStartTimeout: durationpb.New(25 * time.Second),
237
StartToCloseTimeout: durationpb.New(50 * time.Second),
238
HeartbeatTimeout: durationpb.New(25 * time.Second),
239
}},
240
}}, nil
241
}
242
243
return []*commandpb.Command{{
244
CommandType: enumspb.COMMAND_TYPE_COMPLETE_WORKFLOW_EXECUTION,
245
Attributes: &commandpb.Command_CompleteWorkflowExecutionCommandAttributes{CompleteWorkflowExecutionCommandAttributes: &commandpb.CompleteWorkflowExecutionCommandAttributes{
246
Result: payloads.EncodeString("Done"),
247
}},
248
}}, nil
249
}
250
251
atHandler := func(task *workflowservice.PollActivityTaskQueueResponse) (*commonpb.Payloads, bool, error) {
252
return payloads.EncodeString("Activity Result"), false, nil
253
}
254
255
poller := &testcore.TaskPoller{
256
Client: env.FrontendClient(),
257
Namespace: env.Namespace().String(),
258
TaskQueue: &taskqueuepb.TaskQueue{Name: tl, Kind: enumspb.TASK_QUEUE_KIND_NORMAL},
259
Identity: identity,
260
WorkflowTaskHandler: wtHandler,
261
ActivityTaskHandler: atHandler,
262
Logger: env.Logger,
263
T: s.T(),
264
}
265
266
// this function poll events from history side
267
testDescribeTaskQueue := func(namespace string, taskqueue *taskqueuepb.TaskQueue, taskqueueType enumspb.TaskQueueType) []*taskqueuepb.PollerInfo {
268
responseInner, errInner := env.FrontendClient().DescribeTaskQueue(s.Context(), &workflowservice.DescribeTaskQueueRequest{
269
Namespace: namespace,
270
TaskQueue: taskqueue,
271
TaskQueueType: taskqueueType,
272
})
273
274
s.NoError(errInner)
275
return responseInner.Pollers
276
}
277
278
before := time.Now().UTC()
279
280
// when no one polling on the taskqueue (activity or workflow), there shall be no poller information
281
tq := &taskqueuepb.TaskQueue{Name: tl, Kind: enumspb.TASK_QUEUE_KIND_NORMAL}
282
pollerInfos := testDescribeTaskQueue(env.Namespace().String(), tq, enumspb.TASK_QUEUE_TYPE_ACTIVITY)
283
s.Empty(pollerInfos)
284
pollerInfos = testDescribeTaskQueue(env.Namespace().String(), tq, enumspb.TASK_QUEUE_TYPE_WORKFLOW)
285
s.Empty(pollerInfos)
286
287
_, errWorkflowTask := poller.PollAndProcessWorkflowTask()
288
s.NoError(errWorkflowTask)
289
pollerInfos = testDescribeTaskQueue(env.Namespace().String(), tq, enumspb.TASK_QUEUE_TYPE_ACTIVITY)
290
s.Empty(pollerInfos)
291
pollerInfos = testDescribeTaskQueue(env.Namespace().String(), tq, enumspb.TASK_QUEUE_TYPE_WORKFLOW)
292
s.Len(pollerInfos, 1)
293
s.Equal(identity, pollerInfos[0].GetIdentity())
294
s.True(pollerInfos[0].GetLastAccessTime().AsTime().After(before))
295
s.NotEmpty(pollerInfos[0].GetLastAccessTime())
296
297
errActivity := poller.PollAndProcessActivityTask(false)
298
s.NoError(errActivity)
299
pollerInfos = testDescribeTaskQueue(env.Namespace().String(), tq, enumspb.TASK_QUEUE_TYPE_ACTIVITY)
300
s.Len(pollerInfos, 1)
301
s.Equal(identity, pollerInfos[0].GetIdentity())
302
s.True(pollerInfos[0].GetLastAccessTime().AsTime().After(before))
303
s.NotEmpty(pollerInfos[0].GetLastAccessTime())
304
pollerInfos = testDescribeTaskQueue(env.Namespace().String(), tq, enumspb.TASK_QUEUE_TYPE_WORKFLOW)
305
s.Len(pollerInfos, 1)
306
s.Equal(identity, pollerInfos[0].GetIdentity())
307
s.True(pollerInfos[0].GetLastAccessTime().AsTime().After(before))
308
s.NotEmpty(pollerInfos[0].GetLastAccessTime())
309
}