event_generator.go ×48

Frontier kind: Code frontier

unlabeled · c_df43c9e0fe36

8 tests · 1869 LOC · 57 files · introduces 0 tests · 643 LOC · 2 files

Introduces — evidence that enters the hierarchy at this concept

Code
91 ranges643 lines · 2 files
Tests
0 tests

Contains — complete concept membership

All code (extent)
200 ranges1869 lines · 57 files · Browse complete extent
All tests (intent)
8 testsBrowse complete intent

Neighbourhood graph

The orange circle is the focus. Violet and green circles are every ancestor and descendant, broader and narrower, at any distance; blue squares and pink diamonds are the introduced files and exact introduced tests of every visible concept, not only the focus's. Arrows point from broader to narrower concepts and bridge only concepts omitted from this view. Undirected links show source or test introduction. Concept and file size follows LOC; exact test nodes use test-count units.

Introduced files, introduced tests, and structurally relevant concept specialization

In the embedded map, ordinary wheel input scrolls the page; use the visible controls to zoom and drag to pan. Open the full-screen map for canvas navigation: wheel pans, Ctrl/Command plus wheel zooms, and arrow keys pan when this region is focused. On touch screens, open the full-screen map to pan or pinch. If JavaScript or WebGL is unavailable, use the native relationship evidence on this page.

Introduced files, introduced tests, and structurally relevant concept specializationvalid_SA_alias · 0 introduced LOCvalid_SA_aliasvalid_SA_field · 0 introduced LOCvalid_SA_fieldadmin_handler.go ×3 · 6 introduced LOCadmin_handler.go ×3history_event_util.go ×2 · 38 introduced LOChistory_event_util.go ×2invalid_SA_field · 0 introduced LOCinvalid_SA_fieldinvalid_SA_field · 0 introduced LOCinvalid_SA_fieldhistory_event_util.go ×1 · 19 introduced LOChistory_event_util.go ×1mapper_mock.go ×2 · 15 introduced LOCmapper_mock.go ×2admin_handler.go ×2 · 3 introduced LOCadmin_handler.go ×2event_gen.go ×1 · 3 introduced LOCevent_gen.go ×1event_gen.go ×1 · 24 introduced LOCevent_gen.go ×1event_gen.go ×2 · 26 introduced LOCevent_gen.go ×2TestImportWorkflowExecution_NoSearchAttributes · 0 introduced LOCTestImportWorkflowExecut…history_event_util.go ×1 · 10 introduced LOChistory_event_util.go ×1event_gen.go ×1 · 2 introduced LOCevent_gen.go ×1history_event_util.go ×1 · 20 introduced LOChistory_event_util.go ×1admin_handler.go ×5 · 31 introduced LOCadmin_handler.go ×5admin_handler.go ×6 · 32 introduced LOCadmin_handler.go ×6history_require.go ×13 · 72 introduced LOChistory_require.go ×13history_event_util.go ×1 · 12 introduced LOChistory_event_util.go ×1history_event_util.go ×3 · 52 introduced LOChistory_event_util.go ×3history_event_util.go ×2 · 33 introduced LOChistory_event_util.go ×2history_event_util.go ×1 · 19 introduced LOChistory_event_util.go ×1history_event_util.go ×1 · 17 introduced LOChistory_event_util.go ×1history_event_util.go ×2 · 21 introduced LOChistory_event_util.go ×2event_generator.go ×2 · 24 introduced LOCevent_generator.go ×2event_generator.go ×4 · 59 introduced LOCevent_generator.go ×4history_event_util.go ×1 · 14 introduced LOChistory_event_util.go ×1event_generator.go ×1 · 15 introduced LOCevent_generator.go ×1history_event_util.go ×1 · 17 introduced LOChistory_event_util.go ×1history_event_util.go ×2 · 34 introduced LOChistory_event_util.go ×2namespace.go ×1 · 3 introduced LOCnamespace.go ×1namespace.go ×1 · 3 introduced LOCnamespace.go ×1service.pb.go ×2 · 342 introduced LOCservice.pb.go ×2message.pb.go ×2 · 20 introduced LOCmessage.pb.go ×2message.pb.go ×2 · 442 introduced LOCmessage.pb.go ×2defs.go ×5 · 91 introduced LOCdefs.go ×5tags.go ×9 · 37 introduced LOCtags.go ×9message.pb.go ×2 · 49 introduced LOCmessage.pb.go ×2cluster.pb.go ×2 · 200 introduced LOCcluster.pb.go ×2retrypolicy.go ×3 · 19 introduced LOCretrypolicy.go ×3message.pb.go ×2 · 20 introduced LOCmessage.pb.go ×2TestOperatorServiceMetadata, TestWorkflowServiceMetadata · 0 introduced LOCTestOperatorServiceMetad…go.temporal.io/server/api/adminservice/v1/request_response.pb.go · 6661 LOCv1/request_response.pb.g…go.temporal.io/server/api/adminservice/v1/service.pb.go · 297 LOCv1/service.pb.gogo.temporal.io/server/api/clock/v1/message.pb.go · 217 LOCv1/message.pb.gogo.temporal.io/server/api/cluster/v1/message.pb.go · 366 LOCv1/message.pb.gogo.temporal.io/server/api/common/v1/api_category.pb.go · 230 LOCv1/api_category.pb.gogo.temporal.io/server/api/common/v1/dlq.pb.go · 319 LOCv1/dlq.pb.gogo.temporal.io/server/api/deployment/v1/message.pb.go · 4650 LOCv1/message.pb.gogo.temporal.io/server/api/enums/v1/cluster.pb.go · 233 LOCv1/cluster.pb.gogo.temporal.io/server/api/enums/v1/common.pb.go · 288 LOCv1/common.pb.gogo.temporal.io/server/api/enums/v1/dlq.pb.go · 211 LOCv1/dlq.pb.gogo.temporal.io/server/api/enums/v1/fairness_state.pb.go · 147 LOCv1/fairness_state.pb.gogo.temporal.io/server/api/enums/v1/nexus.pb.go · 182 LOCv1/nexus.pb.gogo.temporal.io/server/api/enums/v1/predicate.pb.go · 194 LOCv1/predicate.pb.gogo.temporal.io/server/api/enums/v1/replication.pb.go · 338 LOCv1/replication.pb.gogo.temporal.io/server/api/enums/v1/task.pb.go · 480 LOCv1/task.pb.gogo.temporal.io/server/api/enums/v1/workflow.pb.go · 299 LOCv1/workflow.pb.gogo.temporal.io/server/api/enums/v1/workflow_task_type.pb.go · 148 LOCv1/workflow_task_type.pb…go.temporal.io/server/api/errordetails/v1/message.pb.go · 649 LOCv1/message.pb.gogo.temporal.io/server/api/health/v1/message.pb.go · 324 LOCv1/message.pb.gogo.temporal.io/server/api/history/v1/message.pb.go · 522 LOCv1/message.pb.gogo.temporal.io/server/api/historyservice/v1/request_response.pb.go · 12003 LOCv1/request_response.pb.g…go.temporal.io/server/api/historyservice/v1/service.pb.go · 452 LOCv1/service.pb.gogo.temporal.io/server/api/historyservicemock/v1/service_grpc.pb.mock.go · 3083 LOCv1/service_grpc.pb.mock.…go.temporal.io/server/api/matchingservice/v1/request_response.pb.go · 6878 LOCv1/request_response.pb.g…go.temporal.io/server/api/matchingservice/v1/service.pb.go · 274 LOCv1/service.pb.gogo.temporal.io/server/api/metrics/v1/message.pb.go · 128 LOCv1/message.pb.gogo.temporal.io/server/api/namespace/v1/message.pb.go · 138 LOCv1/message.pb.gogo.temporal.io/server/api/persistence/v1/chasm.pb.go · 1215 LOCv1/chasm.pb.gogo.temporal.io/server/api/persistence/v1/chasm_visibility.pb.go · 170 LOCv1/chasm_visibility.pb.g…go.temporal.io/server/api/persistence/v1/cluster_metadata.pb.go · 313 LOCv1/cluster_metadata.pb.g…go.temporal.io/server/api/persistence/v1/executions.pb.go · 5768 LOCv1/executions.pb.gogo.temporal.io/server/api/persistence/v1/history_tree.pb.go · 298 LOCv1/history_tree.pb.gogo.temporal.io/server/api/persistence/v1/hsm.pb.go · 929 LOCv1/hsm.pb.gogo.temporal.io/server/api/persistence/v1/namespaces.pb.go · 560 LOCv1/namespaces.pb.gogo.temporal.io/server/api/persistence/v1/nexus.pb.go · 509 LOCv1/nexus.pb.gogo.temporal.io/server/api/persistence/v1/predicates.pb.go · 864 LOCv1/predicates.pb.gogo.temporal.io/server/api/persistence/v1/queue_metadata.pb.go · 129 LOCv1/queue_metadata.pb.gogo.temporal.io/server/api/persistence/v1/queues.pb.go · 641 LOCv1/queues.pb.gogo.temporal.io/server/api/persistence/v1/task_queues.pb.go · 924 LOCv1/task_queues.pb.gogo.temporal.io/server/api/persistence/v1/tasks.pb.go · 870 LOCv1/tasks.pb.gogo.temporal.io/server/api/persistence/v1/update.pb.go · 455 LOCv1/update.pb.gogo.temporal.io/server/api/persistence/v1/workflow_mutable_state.pb.go · 560 LOCv1/workflow_mutable_stat…go.temporal.io/server/api/replication/v1/message.pb.go · 2481 LOCv1/message.pb.gogo.temporal.io/server/api/taskqueue/v1/message.pb.go · 1363 LOCv1/message.pb.gogo.temporal.io/server/api/token/v1/message.pb.go · 809 LOCv1/message.pb.gogo.temporal.io/server/api/workflow/v1/message.pb.go · 302 LOCv1/message.pb.gogo.temporal.io/server/common/backoff/retrypolicy.go · 347 LOCbackoff/retrypolicy.gogo.temporal.io/server/common/build/build.go · 56 LOCbuild/build.gogo.temporal.io/server/common/log/tag/tags.go · 1039 LOCtag/tags.gogo.temporal.io/server/common/log/tag/zap_tag.go · 242 LOCtag/zap_tag.gogo.temporal.io/server/common/metrics/defs.go · 76 LOCmetrics/defs.gogo.temporal.io/server/common/metrics/defs_base.go · 30 LOCmetrics/defs_base.gogo.temporal.io/server/common/metrics/noop_impl.go · 58 LOCmetrics/noop_impl.gogo.temporal.io/server/common/metrics/option.go · 21 LOCmetrics/option.gogo.temporal.io/server/common/metrics/registry.go · 75 LOCmetrics/registry.gogo.temporal.io/server/common/namespace/namespace.go · 377 LOCnamespace/namespace.gogo.temporal.io/server/common/searchattribute/event_gen.go · 42 LOCsearchattribute/event_ge…go.temporal.io/server/common/searchattribute/mapper_mock.go · 110 LOCsearchattribute/mapper_m…go.temporal.io/server/common/testing/event_generator.go · 542 LOCtesting/event_generator.…go.temporal.io/server/common/testing/history_event_util.go · 980 LOCtesting/history_event_ut…go.temporal.io/server/common/testing/historyrequire/history_require.go · 612 LOChistoryrequire/history_r…go.temporal.io/server/common/testing/testvars/any.go · 85 LOCtestvars/any.gogo.temporal.io/server/common/testing/testvars/test_vars.go · 460 LOCtestvars/test_vars.gogo.temporal.io/server/service/frontend/admin_handler.go · 2629 LOCfrontend/admin_handler.g…TestOperatorServiceMetadata · introduced test · go.temporal.io/server/common/api/TestOperatorServiceMetadataTestOperatorServiceMetad…TestWorkflowServiceMetadata · introduced test · go.temporal.io/server/common/api/TestWorkflowServiceMetadataTestWorkflowServiceMetad…TestMutexMapBaggage · introduced test · go.temporal.io/server/common/metrics/TestBaggageBenchSuite/TestMutexMapBaggageTestMutexMapBaggageTestSyncMapBaggage · introduced test · go.temporal.io/server/common/metrics/TestBaggageBenchSuite/TestSyncMapBaggageTestSyncMapBaggageTest_HistoryEvent_Generator · introduced test · go.temporal.io/server/common/testing/TestHistoryEventTestSuite/Test_HistoryEvent_GeneratorTest_HistoryEvent_Genera…TestPrintHistoryEvents · introduced test · go.temporal.io/server/common/testing/historyrequire/TestPrintHistoryEventsTestPrintHistoryEventsTestImportWorkflowExecution_NoSearchAttributes · introduced test · go.temporal.io/server/service/frontend/TestAdminHandlerSuite/TestImportWorkflowExecution_NoSearchAttributesTestImportWorkflowExecut…invalid_SA_alias · introduced test · go.temporal.io/server/service/frontend/TestAdminHandlerSuite/TestImportWorkflowExecution_WithAliasedSearchAttributes/invalid_SA_aliasinvalid_SA_aliasinvalid_SA_field · introduced test · go.temporal.io/server/service/frontend/TestAdminHandlerSuite/TestImportWorkflowExecution_WithAliasedSearchAttributes/invalid_SA_fieldinvalid_SA_fieldvalid_SA_alias · introduced test · go.temporal.io/server/service/frontend/TestAdminHandlerSuite/TestImportWorkflowExecution_WithAliasedSearchAttributes/valid_SA_aliasvalid_SA_aliasinvalid_SA_field · introduced test · go.temporal.io/server/service/frontend/TestAdminHandlerSuite/TestImportWorkflowExecution_WithNonAliasedSearchAttributes/invalid_SA_fieldinvalid_SA_fieldvalid_SA_field · introduced test · go.temporal.io/server/service/frontend/TestAdminHandlerSuite/TestImportWorkflowExecution_WithNonAliasedSearchAttributes/valid_SA_fieldvalid_SA_fieldFocused concept · event_generator.go ×48 · 643 introduced LOCevent_generator.go ×48

Graph controls are ready.

Interactive rendering requires JavaScript and WebGL. Use the native relationship evidence on this page while the interactive map is unavailable.

Native relationship evidence

Every exact file and test below is linked only from the concept that introduces it.

Introduced tests

Every collected test enters the hierarchy at exactly one concept.

No tests are introduced at this concept. Its intent tests are introduced by other concepts.

Introduced code

Every collected source range enters the hierarchy at exactly one concept.

2 files ranked by introduced lines: 643 introduced LOC across 91 ranges. Expand a file to inspect source; the > gutter marks introduced lines.

go.temporal.io/server/common/testing/history_event_util.go 413 introduced LOC · 43 ranges

Open complete file

40 nsID namespace.ID,
41 defaultVersion int64,
42 > ) Generator { history_event_util.go
43 >
44 > generator := NewEventGenerator(time.Now().UnixNano())
45 > generator.SetVersion(defaultVersion)
46 > // Functions
47 > notPendingWorkflowTask := func(input ...any) bool {
48 > count := 0
49 > history := input[0].([]Vertex)
50 > for _, e := range history {
51 > switch e.GetName() {
52 > case enumspb.EVENT_TYPE_WORKFLOW_TASK_SCHEDULED.String():
53 > count++
54 case enumspb.EVENT_TYPE_WORKFLOW_TASK_COMPLETED.String(),
55 enumspb.EVENT_TYPE_WORKFLOW_TASK_FAILED.String(),
56 > enumspb.EVENT_TYPE_WORKFLOW_TASK_TIMED_OUT.String(): history_event_util.go
57 > count--
58 }
59 }
60 > return count <= 0 history_event_util.go
61 }
62 > containActivityComplete := func(input ...any) bool { history_event_util.go
63 > history := input[0].([]Vertex)
64 > for _, e := range history {
65 > if e.GetName() == enumspb.EVENT_TYPE_ACTIVITY_TASK_COMPLETED.String() {
66 > return true
67 > }
68 }
69 > return false history_event_util.go
70 }
71 > hasPendingActivity := func(input ...any) bool { history_event_util.go
72 > count := 0
73 > history := input[0].([]Vertex)
74 > for _, e := range history {
75 > switch e.GetName() {
76 > case enumspb.EVENT_TYPE_ACTIVITY_TASK_SCHEDULED.String():
77 > count++
78 case enumspb.EVENT_TYPE_ACTIVITY_TASK_CANCELED.String(),
79 enumspb.EVENT_TYPE_ACTIVITY_TASK_FAILED.String(),
83 }
84 }
85 > return count > 0 history_event_util.go
86 }
87 > canDoBatch := func(currentBatch []Vertex, history []Vertex) bool { history_event_util.go
88 > if len(currentBatch) == 0 {
89 > return true
90 > }
91
92 > hasPendingWorkflowTask := false history_event_util.go
93 > for _, event := range history {
94 > switch event.GetName() {
95 > case enumspb.EVENT_TYPE_WORKFLOW_TASK_SCHEDULED.String():
96 > hasPendingWorkflowTask = true
97 case enumspb.EVENT_TYPE_WORKFLOW_TASK_COMPLETED.String(),
98 enumspb.EVENT_TYPE_WORKFLOW_TASK_FAILED.String(),
99 > enumspb.EVENT_TYPE_WORKFLOW_TASK_TIMED_OUT.String(): history_event_util.go
100 > hasPendingWorkflowTask = false
101 }
102 }
103 > if hasPendingWorkflowTask { history_event_util.go
104 > return false
105 > }
106 > if currentBatch[len(currentBatch)-1].GetName() == enumspb.EVENT_TYPE_WORKFLOW_TASK_SCHEDULED.String() {
107 return false
108 }
109 > if currentBatch[0].GetName() == enumspb.EVENT_TYPE_WORKFLOW_TASK_COMPLETED.String() { history_event_util.go
110 > return len(currentBatch) == 1
111 > }
112 > return true
113 }
114
115 // Setup workflow task model
116 > historyEventModel := NewHistoryEventModel() history_event_util.go
117 > workflowTaskSchedule := NewHistoryEventVertex(enumspb.EVENT_TYPE_WORKFLOW_TASK_SCHEDULED.String())
118 > workflowTaskSchedule.SetDataFunc(func(input ...any) any {
119 > lastGeneratedEvent := input[1].(*historypb.HistoryEvent)
120 > eventID := lastGeneratedEvent.GetEventId() + 1
121 > version := input[2].(int64)
122 > historyEvent := getDefaultHistoryEvent(eventID, version)
123 > historyEvent.EventType = enumspb.EVENT_TYPE_WORKFLOW_TASK_SCHEDULED
124 > historyEvent.Attributes = &historypb.HistoryEvent_WorkflowTaskScheduledEventAttributes{WorkflowTaskScheduledEventAttributes: &historypb.WorkflowTaskScheduledEventAttributes{
125 > TaskQueue: &taskqueuepb.TaskQueue{
126 > Name: taskQueue,
127 > Kind: enumspb.TASK_QUEUE_KIND_NORMAL,
128 > },
129 > StartToCloseTimeout: durationpb.New(timeout),
130 > Attempt: workflowTaskAttempts,
131 > }}
132 > return historyEvent
133 > })
134 > workflowTaskStart := NewHistoryEventVertex(enumspb.EVENT_TYPE_WORKFLOW_TASK_STARTED.String())
135 > workflowTaskStart.SetIsStrictOnNextVertex(true)
136 > workflowTaskStart.SetDataFunc(func(input ...any) any {
137 > lastEvent := input[0].(*historypb.HistoryEvent)
138 > lastGeneratedEvent := input[1].(*historypb.HistoryEvent)
139 > eventID := lastGeneratedEvent.GetEventId() + 1
140 > version := input[2].(int64)
141 > historyEvent := getDefaultHistoryEvent(eventID, version)
142 > historyEvent.EventType = enumspb.EVENT_TYPE_WORKFLOW_TASK_STARTED
143 > historyEvent.Attributes = &historypb.HistoryEvent_WorkflowTaskStartedEventAttributes{WorkflowTaskStartedEventAttributes: &historypb.WorkflowTaskStartedEventAttributes{
144 > ScheduledEventId: lastEvent.EventId,
145 > Identity: identity,
146 > RequestId: uuid.NewString(),
147 > }}
148 > return historyEvent
149 > })
150 > workflowTaskFail := NewHistoryEventVertex(enumspb.EVENT_TYPE_WORKFLOW_TASK_FAILED.String())
151 > workflowTaskFail.SetDataFunc(func(input ...any) any {
152 > lastEvent := input[0].(*historypb.HistoryEvent)
153 > lastGeneratedEvent := input[1].(*historypb.HistoryEvent)
154 > eventID := lastGeneratedEvent.GetEventId() + 1
155 > version := input[2].(int64)
156 > historyEvent := getDefaultHistoryEvent(eventID, version)
157 > historyEvent.EventType = enumspb.EVENT_TYPE_WORKFLOW_TASK_FAILED
158 > historyEvent.Attributes = &historypb.HistoryEvent_WorkflowTaskFailedEventAttributes{WorkflowTaskFailedEventAttributes: &historypb.WorkflowTaskFailedEventAttributes{
159 > ScheduledEventId: lastEvent.GetWorkflowTaskStartedEventAttributes().ScheduledEventId,
160 > StartedEventId: lastEvent.EventId,
161 > Cause: enumspb.WORKFLOW_TASK_FAILED_CAUSE_UNHANDLED_COMMAND,
162 > Identity: identity,
163 > ForkEventVersion: version,
164 > }}
165 > return historyEvent
166 > })
167 > workflowTaskTimedOut := NewHistoryEventVertex(enumspb.EVENT_TYPE_WORKFLOW_TASK_TIMED_OUT.String())
168 > workflowTaskTimedOut.SetDataFunc(func(input ...any) any {
169 > lastEvent := input[0].(*historypb.HistoryEvent)
170 > lastGeneratedEvent := input[1].(*historypb.HistoryEvent)
171 > eventID := lastGeneratedEvent.GetEventId() + 1
172 > version := input[2].(int64)
173 > historyEvent := getDefaultHistoryEvent(eventID, version)
174 > historyEvent.EventType = enumspb.EVENT_TYPE_WORKFLOW_TASK_TIMED_OUT
175 > historyEvent.Attributes = &historypb.HistoryEvent_WorkflowTaskTimedOutEventAttributes{WorkflowTaskTimedOutEventAttributes: &historypb.WorkflowTaskTimedOutEventAttributes{
176 > ScheduledEventId: lastEvent.GetWorkflowTaskStartedEventAttributes().ScheduledEventId,
177 > StartedEventId: lastEvent.EventId,
178 > TimeoutType: enumspb.TIMEOUT_TYPE_SCHEDULE_TO_START,
179 > }}
180 > return historyEvent
181 > })
182 > workflowTaskComplete := NewHistoryEventVertex(enumspb.EVENT_TYPE_WORKFLOW_TASK_COMPLETED.String())
183 > workflowTaskComplete.SetDataFunc(func(input ...any) any {
184 > lastEvent := input[0].(*historypb.HistoryEvent)
185 > lastGeneratedEvent := input[1].(*historypb.HistoryEvent)
186 > eventID := lastGeneratedEvent.GetEventId() + 1
187 > version := input[2].(int64)
188 > historyEvent := getDefaultHistoryEvent(eventID, version)
189 > historyEvent.EventType = enumspb.EVENT_TYPE_WORKFLOW_TASK_COMPLETED
190 > historyEvent.Attributes = &historypb.HistoryEvent_WorkflowTaskCompletedEventAttributes{WorkflowTaskCompletedEventAttributes: &historypb.WorkflowTaskCompletedEventAttributes{
191 > ScheduledEventId: lastEvent.GetWorkflowTaskStartedEventAttributes().ScheduledEventId,
192 > StartedEventId: lastEvent.EventId,
193 > Identity: identity,
194 > BinaryChecksum: checksum,
195 > }}
196 > return historyEvent
197 > })
198 > workflowTaskComplete.SetIsStrictOnNextVertex(true)
199 > workflowTaskComplete.SetMaxNextVertex(2)
200 > workflowTaskScheduleToStart := NewHistoryEventEdge(workflowTaskSchedule, workflowTaskStart)
201 > workflowTaskStartToComplete := NewHistoryEventEdge(workflowTaskStart, workflowTaskComplete)
202 > workflowTaskStartToFail := NewHistoryEventEdge(workflowTaskStart, workflowTaskFail)
203 > workflowTaskStartToTimedOut := NewHistoryEventEdge(workflowTaskStart, workflowTaskTimedOut)
204 > workflowTaskFailToSchedule := NewHistoryEventEdge(workflowTaskFail, workflowTaskSchedule)
205 > workflowTaskFailToSchedule.SetCondition(notPendingWorkflowTask)
206 > workflowTaskTimedOutToSchedule := NewHistoryEventEdge(workflowTaskTimedOut, workflowTaskSchedule)
207 > workflowTaskTimedOutToSchedule.SetCondition(notPendingWorkflowTask)
208 > historyEventModel.AddEdge(workflowTaskScheduleToStart, workflowTaskStartToComplete, workflowTaskStartToFail, workflowTaskStartToTimedOut,
209 > workflowTaskFailToSchedule, workflowTaskTimedOutToSchedule)
210 >
211 > // Setup workflow model
212 > workflowModel := NewHistoryEventModel()
213 >
214 > workflowStart := NewHistoryEventVertex(enumspb.EVENT_TYPE_WORKFLOW_EXECUTION_STARTED.String())
215 > workflowStart.SetDataFunc(func(input ...any) any {
216 > historyEvent := getDefaultHistoryEvent(1, defaultVersion)
217 > historyEvent.EventType = enumspb.EVENT_TYPE_WORKFLOW_EXECUTION_STARTED
218 > historyEvent.Attributes = &historypb.HistoryEvent_WorkflowExecutionStartedEventAttributes{WorkflowExecutionStartedEventAttributes: &historypb.WorkflowExecutionStartedEventAttributes{
219 > WorkflowType: &commonpb.WorkflowType{
220 > Name: workflowType,
221 > },
222 > TaskQueue: &taskqueuepb.TaskQueue{
223 > Name: taskQueue,
224 > Kind: enumspb.TASK_QUEUE_KIND_NORMAL,
225 > },
226 > WorkflowExecutionTimeout: durationpb.New(timeout),
227 > WorkflowRunTimeout: durationpb.New(timeout),
228 > WorkflowTaskTimeout: durationpb.New(timeout),
229 > Identity: identity,
230 > FirstExecutionRunId: uuid.NewString(),
231 > Attempt: 1,
232 > }}
233 > return historyEvent
234 > })
235 > workflowSignal := NewHistoryEventVertex(enumspb.EVENT_TYPE_WORKFLOW_EXECUTION_SIGNALED.String())
236 > workflowSignal.SetDataFunc(func(input ...any) any {
237 lastGeneratedEvent := input[1].(*historypb.HistoryEvent)
238 eventID := lastGeneratedEvent.GetEventId() + 1
246 return historyEvent
247 })
248 > workflowComplete := NewHistoryEventVertex(enumspb.EVENT_TYPE_WORKFLOW_EXECUTION_COMPLETED.String()) history_event_util.go
249 > workflowComplete.SetDataFunc(func(input ...any) any {
250 lastEvent := input[0].(*historypb.HistoryEvent)
251 eventID := lastEvent.GetEventId() + 1
258 return historyEvent
259 })
260 > continueAsNew := NewHistoryEventVertex(enumspb.EVENT_TYPE_WORKFLOW_EXECUTION_CONTINUED_AS_NEW.String()) history_event_util.go
261 > continueAsNew.SetDataFunc(func(input ...any) any {
262 lastEvent := input[0].(*historypb.HistoryEvent)
263 eventID := lastEvent.GetEventId() + 1
281 return historyEvent
282 })
283 > workflowFail := NewHistoryEventVertex(enumspb.EVENT_TYPE_WORKFLOW_EXECUTION_FAILED.String()) history_event_util.go
284 > workflowFail.SetDataFunc(func(input ...any) any {
285 lastEvent := input[0].(*historypb.HistoryEvent)
286 eventID := lastEvent.GetEventId() + 1
293 return historyEvent
294 })
295 > workflowCancel := NewHistoryEventVertex(enumspb.EVENT_TYPE_WORKFLOW_EXECUTION_CANCELED.String()) history_event_util.go
296 > workflowCancel.SetDataFunc(func(input ...any) any {
297 lastEvent := input[0].(*historypb.HistoryEvent)
298 lastGeneratedEvent := input[1].(*historypb.HistoryEvent)
306 return historyEvent
307 })
308 > workflowCancelRequest := NewHistoryEventVertex(enumspb.EVENT_TYPE_WORKFLOW_EXECUTION_CANCEL_REQUESTED.String()) history_event_util.go
309 > workflowCancelRequest.SetDataFunc(func(input ...any) any {
310 lastGeneratedEvent := input[1].(*historypb.HistoryEvent)
311 eventID := lastGeneratedEvent.GetEventId() + 1
324 return historyEvent
325 })
326 > workflowTerminate := NewHistoryEventVertex(enumspb.EVENT_TYPE_WORKFLOW_EXECUTION_TERMINATED.String()) history_event_util.go
327 > workflowTerminate.SetDataFunc(func(input ...any) any {
328 lastEvent := input[0].(*historypb.HistoryEvent)
329 eventID := lastEvent.GetEventId() + 1
337 return historyEvent
338 })
339 > workflowTimedOut := NewHistoryEventVertex(enumspb.EVENT_TYPE_WORKFLOW_EXECUTION_TIMED_OUT.String()) history_event_util.go
340 > workflowTimedOut.SetDataFunc(func(input ...any) any {
341 lastEvent := input[0].(*historypb.HistoryEvent)
342 eventID := lastEvent.GetEventId() + 1
349 return historyEvent
350 })
351 > workflowStartToSignal := NewHistoryEventEdge(workflowStart, workflowSignal) history_event_util.go
352 > workflowStartToWorkflowTaskSchedule := NewHistoryEventEdge(workflowStart, workflowTaskSchedule)
353 > workflowStartToWorkflowTaskSchedule.SetCondition(notPendingWorkflowTask)
354 > workflowSignalToWorkflowTaskSchedule := NewHistoryEventEdge(workflowSignal, workflowTaskSchedule)
355 > workflowSignalToWorkflowTaskSchedule.SetCondition(notPendingWorkflowTask)
356 > workflowTaskCompleteToWorkflowComplete := NewHistoryEventEdge(workflowTaskComplete, workflowComplete)
357 > workflowTaskCompleteToWorkflowComplete.SetCondition(containActivityComplete)
358 > workflowTaskCompleteToWorkflowFailed := NewHistoryEventEdge(workflowTaskComplete, workflowFail)
359 > workflowTaskCompleteToWorkflowFailed.SetCondition(containActivityComplete)
360 > workflowTaskCompleteToCAN := NewHistoryEventEdge(workflowTaskComplete, continueAsNew)
361 > workflowTaskCompleteToCAN.SetCondition(containActivityComplete)
362 > workflowCancelRequestToCancel := NewHistoryEventEdge(workflowCancelRequest, workflowCancel)
363 > workflowModel.AddEdge(workflowStartToSignal, workflowStartToWorkflowTaskSchedule, workflowSignalToWorkflowTaskSchedule,
364 > workflowTaskCompleteToCAN, workflowTaskCompleteToWorkflowComplete, workflowTaskCompleteToWorkflowFailed, workflowCancelRequestToCancel)
365 >
366 > // Setup activity model
367 > activityModel := NewHistoryEventModel()
368 > activitySchedule := NewHistoryEventVertex(enumspb.EVENT_TYPE_ACTIVITY_TASK_SCHEDULED.String())
369 > activitySchedule.SetDataFunc(func(input ...any) any {
370 > lastEvent := input[0].(*historypb.HistoryEvent)
371 > lastGeneratedEvent := input[1].(*historypb.HistoryEvent)
372 > eventID := lastGeneratedEvent.GetEventId() + 1
373 > version := input[2].(int64)
374 > historyEvent := getDefaultHistoryEvent(eventID, version)
375 > historyEvent.EventType = enumspb.EVENT_TYPE_ACTIVITY_TASK_SCHEDULED
376 > historyEvent.Attributes = &historypb.HistoryEvent_ActivityTaskScheduledEventAttributes{ActivityTaskScheduledEventAttributes: &historypb.ActivityTaskScheduledEventAttributes{
377 > ActivityId: uuid.NewString(),
378 > ActivityType: &commonpb.ActivityType{Name: "activity"},
379 > TaskQueue: &taskqueuepb.TaskQueue{
380 > Name: taskQueue,
381 > Kind: enumspb.TASK_QUEUE_KIND_NORMAL,
382 > },
383 > ScheduleToCloseTimeout: durationpb.New(timeout),
384 > ScheduleToStartTimeout: durationpb.New(timeout),
385 > StartToCloseTimeout: durationpb.New(timeout),
386 > WorkflowTaskCompletedEventId: lastEvent.EventId,
387 > }}
388 > return historyEvent
389 > })
390 > activityStart := NewHistoryEventVertex(enumspb.EVENT_TYPE_ACTIVITY_TASK_STARTED.String())
391 > activityStart.SetDataFunc(func(input ...any) any {
392 > lastEvent := input[0].(*historypb.HistoryEvent)
393 > lastGeneratedEvent := input[1].(*historypb.HistoryEvent)
394 > eventID := lastGeneratedEvent.GetEventId() + 1
395 > version := input[2].(int64)
396 > historyEvent := getDefaultHistoryEvent(eventID, version)
397 > historyEvent.EventType = enumspb.EVENT_TYPE_ACTIVITY_TASK_STARTED
398 > historyEvent.Attributes = &historypb.HistoryEvent_ActivityTaskStartedEventAttributes{ActivityTaskStartedEventAttributes: &historypb.ActivityTaskStartedEventAttributes{
399 > ScheduledEventId: lastEvent.EventId,
400 > Identity: identity,
401 > RequestId: uuid.NewString(),
402 > Attempt: 1,
403 > }}
404 > return historyEvent
405 > })
406 > activityComplete := NewHistoryEventVertex(enumspb.EVENT_TYPE_ACTIVITY_TASK_COMPLETED.String())
407 > activityComplete.SetDataFunc(func(input ...any) any {
408 > lastEvent := input[0].(*historypb.HistoryEvent)
409 > lastGeneratedEvent := input[1].(*historypb.HistoryEvent)
410 > eventID := lastGeneratedEvent.GetEventId() + 1
411 > version := input[2].(int64)
412 > historyEvent := getDefaultHistoryEvent(eventID, version)
413 > historyEvent.EventType = enumspb.EVENT_TYPE_ACTIVITY_TASK_COMPLETED
414 > historyEvent.Attributes = &historypb.HistoryEvent_ActivityTaskCompletedEventAttributes{ActivityTaskCompletedEventAttributes: &historypb.ActivityTaskCompletedEventAttributes{
415 > ScheduledEventId: lastEvent.GetActivityTaskStartedEventAttributes().ScheduledEventId,
416 > StartedEventId: lastEvent.EventId,
417 > Identity: identity,
418 > }}
419 > return historyEvent
420 > })
421 > activityFail := NewHistoryEventVertex(enumspb.EVENT_TYPE_ACTIVITY_TASK_FAILED.String())
422 > activityFail.SetDataFunc(func(input ...any) any {
423 lastEvent := input[0].(*historypb.HistoryEvent)
424 lastGeneratedEvent := input[1].(*historypb.HistoryEvent)
435 return historyEvent
436 })
437 > activityTimedOut := NewHistoryEventVertex(enumspb.EVENT_TYPE_ACTIVITY_TASK_TIMED_OUT.String()) history_event_util.go
438 > activityTimedOut.SetDataFunc(func(input ...any) any {
439 lastEvent := input[0].(*historypb.HistoryEvent)
440 lastGeneratedEvent := input[1].(*historypb.HistoryEvent)
454 return historyEvent
455 })
456 > activityCancelRequest := NewHistoryEventVertex(enumspb.EVENT_TYPE_ACTIVITY_TASK_CANCEL_REQUESTED.String()) history_event_util.go
457 > activityCancelRequest.SetDataFunc(func(input ...any) any {
458 lastEvent := input[0].(*historypb.HistoryEvent)
459 lastGeneratedEvent := input[1].(*historypb.HistoryEvent)
468 return historyEvent
469 })
470 > activityCancel := NewHistoryEventVertex(enumspb.EVENT_TYPE_ACTIVITY_TASK_CANCELED.String()) history_event_util.go
471 > activityCancel.SetDataFunc(func(input ...any) any {
472 lastEvent := input[0].(*historypb.HistoryEvent)
473 lastGeneratedEvent := input[1].(*historypb.HistoryEvent)
484 return historyEvent
485 })
486 > workflowTaskCompleteToATSchedule := NewHistoryEventEdge(workflowTaskComplete, activitySchedule) history_event_util.go
487 >
488 > activityScheduleToStart := NewHistoryEventEdge(activitySchedule, activityStart)
489 > activityScheduleToStart.SetCondition(hasPendingActivity)
490 >
491 > activityStartToComplete := NewHistoryEventEdge(activityStart, activityComplete)
492 > activityStartToComplete.SetCondition(hasPendingActivity)
493 >
494 > activityStartToFail := NewHistoryEventEdge(activityStart, activityFail)
495 > activityStartToFail.SetCondition(hasPendingActivity)
496 >
497 > activityStartToTimedOut := NewHistoryEventEdge(activityStart, activityTimedOut)
498 > activityStartToTimedOut.SetCondition(hasPendingActivity)
499 >
500 > activityCompleteToWorkflowTaskSchedule := NewHistoryEventEdge(activityComplete, workflowTaskSchedule)
501 > activityCompleteToWorkflowTaskSchedule.SetCondition(notPendingWorkflowTask)
502 > activityFailToWorkflowTaskSchedule := NewHistoryEventEdge(activityFail, workflowTaskSchedule)
503 > activityFailToWorkflowTaskSchedule.SetCondition(notPendingWorkflowTask)
504 > activityTimedOutToWorkflowTaskSchedule := NewHistoryEventEdge(activityTimedOut, workflowTaskSchedule)
505 > activityTimedOutToWorkflowTaskSchedule.SetCondition(notPendingWorkflowTask)
506 > activityCancelToWorkflowTaskSchedule := NewHistoryEventEdge(activityCancel, workflowTaskSchedule)
507 > activityCancelToWorkflowTaskSchedule.SetCondition(notPendingWorkflowTask)
508 >
509 > // TODO: bypass activity cancel request event. Support this event later.
510 > // activityScheduleToActivityCancelRequest := NewHistoryEventEdge(activitySchedule, activityCancelRequest)
511 > // activityScheduleToActivityCancelRequest.SetCondition(hasPendingActivity)
512 > activityCancelReqToCancel := NewHistoryEventEdge(activityCancelRequest, activityCancel)
513 > activityCancelReqToCancel.SetCondition(hasPendingActivity)
514 >
515 > activityModel.AddEdge(workflowTaskCompleteToATSchedule, activityScheduleToStart, activityStartToComplete,
516 > activityStartToFail, activityStartToTimedOut, workflowTaskCompleteToATSchedule, activityCompleteToWorkflowTaskSchedule,
517 > activityFailToWorkflowTaskSchedule, activityTimedOutToWorkflowTaskSchedule, activityCancelReqToCancel,
518 > activityCancelToWorkflowTaskSchedule)
519 >
520 > // Setup timer model
521 > timerModel := NewHistoryEventModel()
522 > timerStart := NewHistoryEventVertex(enumspb.EVENT_TYPE_TIMER_STARTED.String())
523 > timerStart.SetDataFunc(func(input ...any) any {
524 lastEvent := input[0].(*historypb.HistoryEvent)
525 lastGeneratedEvent := input[1].(*historypb.HistoryEvent)
535 return historyEvent
536 })
537 > timerFired := NewHistoryEventVertex(enumspb.EVENT_TYPE_TIMER_FIRED.String()) history_event_util.go
538 > timerFired.SetDataFunc(func(input ...any) any {
539 lastEvent := input[0].(*historypb.HistoryEvent)
540 lastGeneratedEvent := input[1].(*historypb.HistoryEvent)
549 return historyEvent
550 })
551 > timerCancel := NewHistoryEventVertex(enumspb.EVENT_TYPE_TIMER_CANCELED.String()) history_event_util.go
552 > timerCancel.SetDataFunc(func(input ...any) any {
553 lastEvent := input[0].(*historypb.HistoryEvent)
554 lastGeneratedEvent := input[1].(*historypb.HistoryEvent)
565 return historyEvent
566 })
567 > timerStartToFire := NewHistoryEventEdge(timerStart, timerFired) history_event_util.go
568 > timerStartToCancel := NewHistoryEventEdge(timerStart, timerCancel)
569 >
570 > workflowTaskCompleteToTimerStart := NewHistoryEventEdge(workflowTaskComplete, timerStart)
571 > timerFiredToWorkflowTaskSchedule := NewHistoryEventEdge(timerFired, workflowTaskSchedule)
572 > timerFiredToWorkflowTaskSchedule.SetCondition(notPendingWorkflowTask)
573 > timerCancelToWorkflowTaskSchedule := NewHistoryEventEdge(timerCancel, workflowTaskSchedule)
574 > timerCancelToWorkflowTaskSchedule.SetCondition(notPendingWorkflowTask)
575 > timerModel.AddEdge(timerStartToFire, timerStartToCancel, workflowTaskCompleteToTimerStart, timerFiredToWorkflowTaskSchedule, timerCancelToWorkflowTaskSchedule)
576 >
577 > // Setup child workflow model
578 > childWorkflowModel := NewHistoryEventModel()
579 > childWorkflowInitial := NewHistoryEventVertex(enumspb.EVENT_TYPE_START_CHILD_WORKFLOW_EXECUTION_INITIATED.String())
580 > childWorkflowInitial.SetDataFunc(func(input ...any) any {
581 lastEvent := input[0].(*historypb.HistoryEvent)
582 lastGeneratedEvent := input[1].(*historypb.HistoryEvent)
602 return historyEvent
603 })
604 > childWorkflowInitialFail := NewHistoryEventVertex(enumspb.EVENT_TYPE_START_CHILD_WORKFLOW_EXECUTION_FAILED.String()) history_event_util.go
605 > childWorkflowInitialFail.SetDataFunc(func(input ...any) any {
606 lastEvent := input[0].(*historypb.HistoryEvent)
607 lastGeneratedEvent := input[1].(*historypb.HistoryEvent)
621 return historyEvent
622 })
623 > childWorkflowStart := NewHistoryEventVertex(enumspb.EVENT_TYPE_CHILD_WORKFLOW_EXECUTION_STARTED.String()) history_event_util.go
624 > childWorkflowStart.SetDataFunc(func(input ...any) any {
625 lastEvent := input[0].(*historypb.HistoryEvent)
626 lastGeneratedEvent := input[1].(*historypb.HistoryEvent)
641 return historyEvent
642 })
643 > childWorkflowCancel := NewHistoryEventVertex(enumspb.EVENT_TYPE_CHILD_WORKFLOW_EXECUTION_CANCELED.String()) history_event_util.go
644 > childWorkflowCancel.SetDataFunc(func(input ...any) any {
645 lastEvent := input[0].(*historypb.HistoryEvent)
646 lastGeneratedEvent := input[1].(*historypb.HistoryEvent)
662 return historyEvent
663 })
664 > childWorkflowComplete := NewHistoryEventVertex(enumspb.EVENT_TYPE_CHILD_WORKFLOW_EXECUTION_COMPLETED.String()) history_event_util.go
665 > childWorkflowComplete.SetDataFunc(func(input ...any) any {
666 lastEvent := input[0].(*historypb.HistoryEvent)
667 lastGeneratedEvent := input[1].(*historypb.HistoryEvent)
683 return historyEvent
684 })
685 > childWorkflowFail := NewHistoryEventVertex(enumspb.EVENT_TYPE_CHILD_WORKFLOW_EXECUTION_FAILED.String()) history_event_util.go
686 > childWorkflowFail.SetDataFunc(func(input ...any) any {
687 lastEvent := input[0].(*historypb.HistoryEvent)
688 lastGeneratedEvent := input[1].(*historypb.HistoryEvent)
704 return historyEvent
705 })
706 > childWorkflowTerminate := NewHistoryEventVertex(enumspb.EVENT_TYPE_CHILD_WORKFLOW_EXECUTION_TERMINATED.String()) history_event_util.go
707 > childWorkflowTerminate.SetDataFunc(func(input ...any) any {
708 lastEvent := input[0].(*historypb.HistoryEvent)
709 lastGeneratedEvent := input[1].(*historypb.HistoryEvent)
725 return historyEvent
726 })
727 > childWorkflowTimedOut := NewHistoryEventVertex(enumspb.EVENT_TYPE_CHILD_WORKFLOW_EXECUTION_TIMED_OUT.String()) history_event_util.go
728 > childWorkflowTimedOut.SetDataFunc(func(input ...any) any {
729 lastEvent := input[0].(*historypb.HistoryEvent)
730 lastGeneratedEvent := input[1].(*historypb.HistoryEvent)
747 return historyEvent
748 })
749 > workflowTaskCompleteToChildWorkflowInitial := NewHistoryEventEdge(workflowTaskComplete, childWorkflowInitial) history_event_util.go
750 > childWorkflowInitialToFail := NewHistoryEventEdge(childWorkflowInitial, childWorkflowInitialFail)
751 > childWorkflowInitialToStart := NewHistoryEventEdge(childWorkflowInitial, childWorkflowStart)
752 > childWorkflowStartToCancel := NewHistoryEventEdge(childWorkflowStart, childWorkflowCancel)
753 > childWorkflowStartToFail := NewHistoryEventEdge(childWorkflowStart, childWorkflowFail)
754 > childWorkflowStartToComplete := NewHistoryEventEdge(childWorkflowStart, childWorkflowComplete)
755 > childWorkflowStartToTerminate := NewHistoryEventEdge(childWorkflowStart, childWorkflowTerminate)
756 > childWorkflowStartToTimedOut := NewHistoryEventEdge(childWorkflowStart, childWorkflowTimedOut)
757 > childWorkflowCancelToWorkflowTaskSchedule := NewHistoryEventEdge(childWorkflowCancel, workflowTaskSchedule)
758 > childWorkflowCancelToWorkflowTaskSchedule.SetCondition(notPendingWorkflowTask)
759 > childWorkflowFailToWorkflowTaskSchedule := NewHistoryEventEdge(childWorkflowFail, workflowTaskSchedule)
760 > childWorkflowFailToWorkflowTaskSchedule.SetCondition(notPendingWorkflowTask)
761 > childWorkflowCompleteToWorkflowTaskSchedule := NewHistoryEventEdge(childWorkflowComplete, workflowTaskSchedule)
762 > childWorkflowCompleteToWorkflowTaskSchedule.SetCondition(notPendingWorkflowTask)
763 > childWorkflowTerminateToWorkflowTaskSchedule := NewHistoryEventEdge(childWorkflowTerminate, workflowTaskSchedule)
764 > childWorkflowTerminateToWorkflowTaskSchedule.SetCondition(notPendingWorkflowTask)
765 > childWorkflowTimedOutToWorkflowTaskSchedule := NewHistoryEventEdge(childWorkflowTimedOut, workflowTaskSchedule)
766 > childWorkflowTimedOutToWorkflowTaskSchedule.SetCondition(notPendingWorkflowTask)
767 > childWorkflowInitialFailToWorkflowTaskSchedule := NewHistoryEventEdge(childWorkflowInitialFail, workflowTaskSchedule)
768 > childWorkflowInitialFailToWorkflowTaskSchedule.SetCondition(notPendingWorkflowTask)
769 > childWorkflowModel.AddEdge(workflowTaskCompleteToChildWorkflowInitial, childWorkflowInitialToFail, childWorkflowInitialToStart,
770 > childWorkflowStartToCancel, childWorkflowStartToFail, childWorkflowStartToComplete, childWorkflowStartToTerminate,
771 > childWorkflowStartToTimedOut, childWorkflowCancelToWorkflowTaskSchedule, childWorkflowFailToWorkflowTaskSchedule,
772 > childWorkflowCompleteToWorkflowTaskSchedule, childWorkflowTerminateToWorkflowTaskSchedule, childWorkflowTimedOutToWorkflowTaskSchedule,
773 > childWorkflowInitialFailToWorkflowTaskSchedule)
774 >
775 > // Setup external workflow model
776 > externalWorkflowModel := NewHistoryEventModel()
777 > externalWorkflowSignal := NewHistoryEventVertex(enumspb.EVENT_TYPE_SIGNAL_EXTERNAL_WORKFLOW_EXECUTION_INITIATED.String())
778 > externalWorkflowSignal.SetDataFunc(func(input ...any) any {
779 lastEvent := input[0].(*historypb.HistoryEvent)
780 lastGeneratedEvent := input[1].(*historypb.HistoryEvent)
796 return historyEvent
797 })
798 > externalWorkflowSignalFailed := NewHistoryEventVertex(enumspb.EVENT_TYPE_SIGNAL_EXTERNAL_WORKFLOW_EXECUTION_FAILED.String()) history_event_util.go
799 > externalWorkflowSignalFailed.SetDataFunc(func(input ...any) any {
800 lastEvent := input[0].(*historypb.HistoryEvent)
801 lastGeneratedEvent := input[1].(*historypb.HistoryEvent)
817 return historyEvent
818 })
819 > externalWorkflowSignaled := NewHistoryEventVertex(enumspb.EVENT_TYPE_EXTERNAL_WORKFLOW_EXECUTION_SIGNALED.String()) history_event_util.go
820 > externalWorkflowSignaled.SetDataFunc(func(input ...any) any {
821 lastEvent := input[0].(*historypb.HistoryEvent)
822 lastGeneratedEvent := input[1].(*historypb.HistoryEvent)
836 return historyEvent
837 })
838 > externalWorkflowCancel := NewHistoryEventVertex(enumspb.EVENT_TYPE_REQUEST_CANCEL_EXTERNAL_WORKFLOW_EXECUTION_INITIATED.String()) history_event_util.go
839 > externalWorkflowCancel.SetDataFunc(func(input ...any) any {
840 lastEvent := input[0].(*historypb.HistoryEvent)
841 lastGeneratedEvent := input[1].(*historypb.HistoryEvent)
857 return historyEvent
858 })
859 > externalWorkflowCancelFail := NewHistoryEventVertex(enumspb.EVENT_TYPE_REQUEST_CANCEL_EXTERNAL_WORKFLOW_EXECUTION_FAILED.String()) history_event_util.go
860 > externalWorkflowCancelFail.SetDataFunc(func(input ...any) any {
861 lastEvent := input[0].(*historypb.HistoryEvent)
862 lastGeneratedEvent := input[1].(*historypb.HistoryEvent)
878 return historyEvent
879 })
880 > externalWorkflowCanceled := NewHistoryEventVertex(enumspb.EVENT_TYPE_EXTERNAL_WORKFLOW_EXECUTION_CANCEL_REQUESTED.String()) history_event_util.go
881 > externalWorkflowCanceled.SetDataFunc(func(input ...any) any {
882 lastEvent := input[0].(*historypb.HistoryEvent)
883 lastGeneratedEvent := input[1].(*historypb.HistoryEvent)
897 return historyEvent
898 })
899 > workflowTaskCompleteToExternalWorkflowSignal := NewHistoryEventEdge(workflowTaskComplete, externalWorkflowSignal) history_event_util.go
900 > workflowTaskCompleteToExternalWorkflowCancel := NewHistoryEventEdge(workflowTaskComplete, externalWorkflowCancel)
901 > externalWorkflowSignalToFail := NewHistoryEventEdge(externalWorkflowSignal, externalWorkflowSignalFailed)
902 > externalWorkflowSignalToSignaled := NewHistoryEventEdge(externalWorkflowSignal, externalWorkflowSignaled)
903 > externalWorkflowCancelToFail := NewHistoryEventEdge(externalWorkflowCancel, externalWorkflowCancelFail)
904 > externalWorkflowCancelToCanceled := NewHistoryEventEdge(externalWorkflowCancel, externalWorkflowCanceled)
905 > externalWorkflowSignaledToWorkflowTaskSchedule := NewHistoryEventEdge(externalWorkflowSignaled, workflowTaskSchedule)
906 > externalWorkflowSignaledToWorkflowTaskSchedule.SetCondition(notPendingWorkflowTask)
907 > externalWorkflowSignalFailedToWorkflowTaskSchedule := NewHistoryEventEdge(externalWorkflowSignalFailed, workflowTaskSchedule)
908 > externalWorkflowSignalFailedToWorkflowTaskSchedule.SetCondition(notPendingWorkflowTask)
909 > externalWorkflowCanceledToWorkflowTaskSchedule := NewHistoryEventEdge(externalWorkflowCanceled, workflowTaskSchedule)
910 > externalWorkflowCanceledToWorkflowTaskSchedule.SetCondition(notPendingWorkflowTask)
911 > externalWorkflowCancelFailToWorkflowTaskSchedule := NewHistoryEventEdge(externalWorkflowCancelFail, workflowTaskSchedule)
912 > externalWorkflowCancelFailToWorkflowTaskSchedule.SetCondition(notPendingWorkflowTask)
913 > externalWorkflowModel.AddEdge(workflowTaskCompleteToExternalWorkflowSignal, workflowTaskCompleteToExternalWorkflowCancel,
914 > externalWorkflowSignalToFail, externalWorkflowSignalToSignaled, externalWorkflowCancelToFail, externalWorkflowCancelToCanceled,
915 > externalWorkflowSignaledToWorkflowTaskSchedule, externalWorkflowSignalFailedToWorkflowTaskSchedule,
916 > externalWorkflowCanceledToWorkflowTaskSchedule, externalWorkflowCancelFailToWorkflowTaskSchedule)
917 >
918 > // Config event generator
919 > generator.SetBatchGenerationRule(canDoBatch)
920 > generator.AddInitialEntryVertex(workflowStart)
921 > generator.AddExitVertex(workflowComplete, workflowFail, workflowTerminate, workflowTimedOut, continueAsNew)
922 > // generator.AddRandomEntryVertex(workflowSignal, workflowTerminate, workflowTimedOut)
923 > generator.AddModel(historyEventModel)
924 > generator.AddModel(workflowModel)
925 > generator.AddModel(activityModel)
926 > generator.AddModel(timerModel)
927 > generator.AddModel(childWorkflowModel)
928 > generator.AddModel(externalWorkflowModel)
929 > return generator
930 }
931
933 eventID int64,
934 version int64,
935 > ) *historypb.HistoryEvent { history_event_util.go
936 >
937 > globalTaskID++
938 > return &historypb.HistoryEvent{
939 > EventId: eventID,
940 > EventTime: timestamppb.New(time.Now().UTC()),
941 > TaskId: globalTaskID,
942 > Version: version,
943 > }
944 > }
945
946 func copyConnections(
go.temporal.io/server/common/testing/event_generator.go 230 introduced LOC · 48 ranges

Open complete file

64 func NewEventGenerator(
65 seed int64,
66 > ) Generator { event_generator.go
67 >
68 > return &EventGenerator{
69 > connections: make(map[string][]Edge),
70 > previousVertices: make([]Vertex, 0),
71 > leafVertices: make([]Vertex, 0),
72 > entryVertices: make([]Vertex, 0),
73 > exitVertices: make(map[string]bool),
74 > randomEntryVertices: make([]Vertex, 0),
75 > dice: rand.New(rand.NewSource(seed)),
76 > seed: seed,
77 > canDoBatch: defaultBatchFunc,
78 > version: defaultVersion,
79 > }
80 > }
81
82 // AddInitialEntryVertex adds the initial history event vertices
84 func (g *EventGenerator) AddInitialEntryVertex(
85 entry ...Vertex,
87 >
88 > g.entryVertices = append(
89 > g.entryVertices,
90 > entry...)
91 > }
92
93 // AddExitVertex adds the terminate history event vertex
94 func (g *EventGenerator) AddExitVertex(
95 exit ...Vertex,
97 >
98 > for _, v := range exit {
99 > g.exitVertices[v.GetName()] = true
100 > }
101 }
102
112 func (g *EventGenerator) AddModel(
113 model Model,
114 > ) { event_generator.go
115 >
116 > for _, e := range model.ListEdges() {
117 > if _, ok := g.connections[e.GetStartVertex().GetName()]; !ok {
118 > g.connections[e.GetStartVertex().GetName()] = make([]Edge, 0)
119 > }
120 > g.connections[e.GetStartVertex().GetName()] = append(g.connections[e.GetStartVertex().GetName()], e)
121 }
122 }
129
130 // HasNextVertex checks if there is accessible history event vertex
131 > func (g *EventGenerator) HasNextVertex() bool { event_generator.go
132 >
133 > for _, prev := range g.previousVertices {
134 > if _, ok := g.exitVertices[prev.GetName()]; ok {
135 return false
136 }
137 }
138 > return len(g.leafVertices) > 0 || (len(g.previousVertices) == 0 && len(g.entryVertices) > 0) event_generator.go
139 }
140
141 // GetNextVertices generates a batch of history events happened in the same transaction
142 > func (g *EventGenerator) GetNextVertices() []Vertex { event_generator.go
143 >
144 > if !g.HasNextVertex() {
145 panic("Generator reached to a terminate state.")
146 }
147
148 > batch := make([]Vertex, 0) event_generator.go
149 > for g.HasNextVertex() && g.canDoBatch(batch, g.previousVertices) {
150 > res := g.generateNextEventBatch()
151 > g.updateContext(res)
152 > batch = append(batch, res...)
153 > }
154 > return batch
155 }
156
183 func (g *EventGenerator) SetBatchGenerationRule(
184 canDoBatchFunc func([]Vertex, []Vertex) bool,
185 > ) { event_generator.go
186 >
187 > g.canDoBatch = canDoBatchFunc
188 > }
189
190 // SetVersion sets the event version
191 > func (g *EventGenerator) SetVersion(version int64) { event_generator.go
192 >
193 > g.version = version
194 > }
195
196 // GetVersion returns event version
200 }
201
202 > func (g *EventGenerator) generateNextEventBatch() []Vertex { event_generator.go
203 >
204 > batch := make([]Vertex, 0)
205 > switch {
206 > case len(g.previousVertices) == 0:
207 > // Generate for the first time, get the event candidates from entry vertex group
208 > batch = append(batch, g.getEntryVertex())
209 case len(g.randomEntryVertices) > 0 && g.dice.Intn(len(g.connections)) == 0:
210 // Get the event candidate from random vertex group
211 batch = append(batch, g.getRandomVertex())
212 > default: event_generator.go
213 > // Get the event candidates based on context
214 > idx := g.getVertexCandidate()
215 > batch = append(batch, g.randomNextVertex(idx)...)
216 > g.leafVertices = append(g.leafVertices[:idx], g.leafVertices[idx+1:]...)
217 }
218 > return batch event_generator.go
219 }
220
221 func (g *EventGenerator) updateContext(
222 batch []Vertex,
223 > ) { event_generator.go
224 >
225 > g.leafVertices = append(g.leafVertices, batch...)
226 > g.previousVertices = append(g.previousVertices, batch...)
227 > }
228
229 > func (g *EventGenerator) getEntryVertex() Vertex { event_generator.go
230 >
231 > if len(g.entryVertices) == 0 {
232 panic("No possible start vertex to go to next step")
233 }
234 > nextRange := len(g.entryVertices) event_generator.go
235 > nextIdx := g.dice.Intn(nextRange)
236 > vertex := g.entryVertices[nextIdx].DeepCopy()
237 > vertex.GenerateData(nil)
238 > return vertex
239 }
240
253 }
254
255 > func (g *EventGenerator) getVertexCandidate() int { event_generator.go
256 >
257 > if len(g.leafVertices) == 0 {
258 panic("No possible vertex to go to next step")
259 }
260 > nextRange := len(g.leafVertices) event_generator.go
261 > notAvailable := make(map[int]bool)
262 > var nextVertexIdx int
263 > for len(notAvailable) < nextRange {
264 > nextVertexIdx = g.dice.Intn(nextRange)
265 > // If the vertex is not accessible at this state, skip it
266 > if _, ok := notAvailable[nextVertexIdx]; ok {
267 continue
268 }
269 > isAccessible, nextVertexIdx := g.findAccessibleVertex(nextVertexIdx) event_generator.go
270 > if isAccessible {
271 > return nextVertexIdx
272 > }
273 > notAvailable[nextVertexIdx] = true
274 }
275 // If all history event cannot be accessible, which means the model is incorrect
279 func (g *EventGenerator) findAccessibleVertex(
280 vertexIndex int,
281 > ) (bool, int) { event_generator.go
282 >
283 > candidate := g.leafVertices[vertexIndex]
284 > if g.leafVertices[len(g.leafVertices)-1].IsStrictOnNextVertex() {
285 > vertexIndex = len(g.leafVertices) - 1
286 > candidate = g.leafVertices[vertexIndex]
287 > }
288 > neighbors := g.connections[candidate.GetName()]
289 > for _, nextV := range neighbors {
290 > if nextV.GetCondition() == nil || nextV.GetCondition()(g.previousVertices) {
291 > return true, vertexIndex
292 > }
293 }
294 > return false, emptyCandidateIndex event_generator.go
295 }
296
297 func (g *EventGenerator) randomNextVertex(
298 nextVertexIdx int,
299 > ) []Vertex { event_generator.go
300 >
301 > nextVertex := g.leafVertices[nextVertexIdx]
302 >
303 > count := g.dice.Intn(nextVertex.GetMaxNextVertex()) + 1
304 > res := make([]Vertex, 0)
305 > latestVertex := g.previousVertices[len(g.previousVertices)-1]
306 > for range count {
307 > endVertex := g.pickRandomVertex(nextVertex)
308 > endVertex.GenerateData(nextVertex.GetData(), latestVertex.GetData(), g.version)
309 > latestVertex = endVertex
310 > res = append(res, endVertex)
311 > if _, ok := g.exitVertices[endVertex.GetName()]; ok {
312 res = []Vertex{endVertex}
313 return res
314 }
315 }
316 > return res event_generator.go
317 }
318
319 func (g *EventGenerator) pickRandomVertex(
320 nextVertex Vertex,
321 > ) Vertex { event_generator.go
322 >
323 > neighbors := g.connections[nextVertex.GetName()]
324 > neighborsRange := len(neighbors)
325 > nextIdx := g.dice.Intn(neighborsRange)
326 > for neighbors[nextIdx].GetCondition() != nil && !neighbors[nextIdx].GetCondition()(g.previousVertices) {
327 nextIdx = g.dice.Intn(neighborsRange)
328 }
329 > newConnection := neighbors[nextIdx] event_generator.go
330 > endVertex := newConnection.GetEndVertex()
331 > if newConnection.GetAction() != nil {
332 newConnection.GetAction()()
333 }
334 > return endVertex.DeepCopy() event_generator.go
335 }
336
339 start Vertex,
340 end Vertex,
341 > ) Edge { event_generator.go
342 >
343 > return &HistoryEventEdge{
344 > startVertex: start,
345 > endVertex: end,
346 > }
347 > }
348
349 // SetStartVertex sets the start vertex
356
357 // GetStartVertex returns the start vertex
358 > func (c HistoryEventEdge) GetStartVertex() Vertex { event_generator.go
359 >
360 > return c.startVertex
361 > }
362
363 // SetEndVertex sets the end vertex
368
369 // GetEndVertex returns the end vertex
370 > func (c HistoryEventEdge) GetEndVertex() Vertex { event_generator.go
371 >
372 > return c.endVertex
373 > }
374
375 // SetCondition sets the condition to access this edge
376 func (c *HistoryEventEdge) SetCondition(
377 condition func(...any) bool,
378 > ) { event_generator.go
379 >
380 > c.condition = condition
381 > }
382
383 // GetCondition returns the condition
384 > func (c HistoryEventEdge) GetCondition() func(...any) bool { event_generator.go
385 >
386 > return c.condition
387 > }
388
389 // SetAction sets an action to perform when the end vertex hits
394
395 // GetAction returns the action
396 > func (c HistoryEventEdge) GetAction() func() { event_generator.go
397 >
398 > return c.action
399 > }
400
401 // DeepCopy copies a new edge
413 func NewHistoryEventVertex(
414 name string,
415 > ) Vertex { event_generator.go
416 >
417 > return &HistoryEventVertex{
418 > name: name,
419 > isStrictOnNextVertex: false,
420 > maxNextGeneration: 1,
421 > }
422 > }
423
424 // GetName returns the name
425 > func (he HistoryEventVertex) GetName() string { event_generator.go
426 >
427 > return he.name
428 > }
429
430 // SetName sets the name
447 func (he *HistoryEventVertex) SetIsStrictOnNextVertex(
448 isStrict bool,
449 > ) { event_generator.go
450 >
451 > he.isStrictOnNextVertex = isStrict
452 > }
453
454 // IsStrictOnNextVertex returns the isStrict flag
455 > func (he HistoryEventVertex) IsStrictOnNextVertex() bool { event_generator.go
456 >
457 > return he.isStrictOnNextVertex
458 > }
459
460 // SetMaxNextVertex sets the max concurrent path can be generated from this vertex
461 func (he *HistoryEventVertex) SetMaxNextVertex(
462 maxNextGeneration int,
463 > ) { event_generator.go
464 >
465 > if maxNextGeneration < 1 {
466 panic("max next vertex number cannot less than 1")
467 }
468 > he.maxNextGeneration = maxNextGeneration event_generator.go
469 }
470
471 // GetMaxNextVertex returns the max concurrent path
472 > func (he HistoryEventVertex) GetMaxNextVertex() int { event_generator.go
473 >
474 > return he.maxNextGeneration
475 > }
476
477 // SetDataFunc sets the data generation function
478 func (he *HistoryEventVertex) SetDataFunc(
479 dataFunc func(...any) any,
480 > ) { event_generator.go
481 >
482 > he.dataFunc = dataFunc
483 > }
484
485 // GetDataFunc returns the data generation function
486 > func (he HistoryEventVertex) GetDataFunc() func(...any) any { event_generator.go
487 >
488 > return he.dataFunc
489 > }
490
491 // GenerateData generates the data and return
492 func (he *HistoryEventVertex) GenerateData(
493 input ...any,
494 > ) any { event_generator.go
495 >
496 > if he.dataFunc == nil {
497 return nil
498 }
499
500 > he.data = he.dataFunc(input...) event_generator.go
501 > return he.data
502 }
503
504 // GetData returns the vertex data
505 > func (he HistoryEventVertex) GetData() any { event_generator.go
506 >
507 > return he.data
508 > }
509
510 // DeepCopy returns the a deep copy of vertex
511 > func (he HistoryEventVertex) DeepCopy() Vertex { event_generator.go
512 >
513 > return &HistoryEventVertex{
514 > name: he.GetName(),
515 > isStrictOnNextVertex: he.IsStrictOnNextVertex(),
516 > maxNextGeneration: he.GetMaxNextVertex(),
517 > dataFunc: he.GetDataFunc(),
518 > data: he.GetData(),
519 > }
520 > }
521
522 // NewHistoryEventModel initials new history event model
523 > func NewHistoryEventModel() Model { event_generator.go
524 >
525 > return &HistoryEventModel{
526 > edges: make([]Edge, 0),
527 > }
528 > }
529
530 // AddEdge adds an edge to the model
531 func (m *HistoryEventModel) AddEdge(
532 edge ...Edge,
533 > ) { event_generator.go
534 >
535 > m.edges = append(m.edges, edge...)
536 > }
537
538 // ListEdges returns all added edges
539 > func (m HistoryEventModel) ListEdges() []Edge { event_generator.go
540 >
541 > return m.edges
542 > }