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 }