73
)
74
75
>
func validateParams(ctx workflow.Context, params *DeleteExecutionsParams) error {
workflow.go
76
>
if params.NamespaceID.IsEmpty() {
77
return errors.NewInvalidArgument("namespace ID is required", nil)
78
}
80
return errors.NewInvalidArgument("namespace is required", nil)
81
}
82
>
if params.FirstRunStartTime.IsZero() {
workflow.go
83
params.FirstRunStartTime = workflow.Now(ctx).UTC()
84
}
86
>
return nil
87
}
88
89
>
func DeleteExecutionsWorkflow(ctx workflow.Context, params DeleteExecutionsParams) (DeleteExecutionsResult, error) {
workflow.go
90
>
logger := log.With(
91
>
workflow.GetLogger(ctx),
92
>
tag.WorkflowType(WorkflowName),
93
>
tag.WorkflowNamespace(params.Namespace.String()),
94
>
tag.WorkflowNamespaceID(params.NamespaceID.String()))
95
>
96
>
logger.Info("Workflow started.")
97
>
result := DeleteExecutionsResult{
98
>
SuccessCount: params.PreviousSuccessCount,
99
>
ErrorCount: params.PreviousErrorCount,
100
>
}
101
>
102
>
if err := validateParams(ctx, ¶ms); err != nil {
103
return result, err
104
}
105
>
logger.Info("Effective config.", tag.Value(params.Config.String()))
workflow.go
106
>
107
>
if err := workflow.SetQueryHandler(ctx, StatsQuery, func() (DeleteExecutionsStats, error) {
109
>
des := DeleteExecutionsStats{
110
>
DeleteExecutionsResult: result,
111
>
ContinueAsNewCount: params.ContinueAsNewCount,
112
>
TotalExecutionsCount: params.TotalExecutionsCount,
113
>
StartTime: params.FirstRunStartTime,
114
>
}
115
>
if params.TotalExecutionsCount > 0 {
116
>
des.RemainingExecutionsCount = params.TotalExecutionsCount - (result.SuccessCount + result.ErrorCount)
117
>
}
118
>
secondsSinceStart := int(now.Sub(params.FirstRunStartTime).Seconds())
119
>
if secondsSinceStart > 0 {
120
>
des.AverageRPS = (result.SuccessCount + result.ErrorCount) / secondsSinceStart
121
>
}
122
>
if des.AverageRPS > 0 {
123
>
des.ApproximateTimeLeft = time.Duration(des.RemainingExecutionsCount/des.AverageRPS) * time.Second
124
>
}
125
>
des.ApproximateEndTime = now.Add(des.ApproximateTimeLeft)
126
>
return des, nil
127
}); err != nil {
128
return result, err
129
}
130
132
>
var la *LocalActivities
133
>
134
>
ctx = workflow.WithTaskQueue(ctx, primitives.DeleteNamespaceActivityTQ)
135
>
136
>
nextPageToken := params.NextPageToken
137
>
runningDeleteExecutionsActivityCount := 0
138
>
runningDeleteExecutionsSelector := workflow.NewSelector(ctx)
139
>
var lastDeleteExecutionsActivityErr error
140
>
141
>
// Two activities DeleteExecutionsActivity and GetNextPageTokenActivity are executed here essentially in reverse order
142
>
// because Get is called immediately for GetNextPageTokenActivity but not for DeleteExecutionsActivity.
143
>
// These activities scan visibility storage independently but GetNextPageTokenActivity considered to be quick and can be done synchronously.
144
>
// It reads nextPageToken and pass it DeleteExecutionsActivity. This allocates block of workflow executions to delete for
145
>
// DeleteExecutionsActivity which takes much longer to complete. This is why this workflow starts
146
>
// ConcurrentDeleteExecutionsActivities number of them and executes them concurrently on available workers.
147
>
for i := 0; i < params.Config.PagesPerExecution; i++ {
148
>
ctx1 := workflow.WithActivityOptions(ctx, deleteWorkflowExecutionsActivityOptions)
149
>
deleteExecutionsFuture := workflow.ExecuteActivity(ctx1, a.DeleteExecutionsActivity, &DeleteExecutionsActivityParams{
150
>
Namespace: params.Namespace,
151
>
NamespaceID: params.NamespaceID,
152
>
RPS: params.Config.DeleteActivityRPS,
153
>
ListPageSize: params.Config.PageSize,
154
>
NextPageToken: nextPageToken,
155
>
})
156
>
157
>
ctx2 := workflow.WithLocalActivityOptions(ctx, localActivityOptions)
158
>
err := workflow.ExecuteLocalActivity(ctx2, la.GetNextPageTokenActivity, GetNextPageTokenParams{
159
>
NamespaceID: params.NamespaceID,
160
>
Namespace: params.Namespace,
161
>
PageSize: params.Config.PageSize,
162
>
NextPageToken: nextPageToken,
163
>
}).Get(ctx, &nextPageToken)
164
>
if err != nil {
165
return result, err
166
}
167
168
>
runningDeleteExecutionsActivityCount++
workflow.go
169
>
runningDeleteExecutionsSelector.AddFuture(deleteExecutionsFuture, func(f workflow.Future) {
170
>
runningDeleteExecutionsActivityCount--
171
>
var der DeleteExecutionsActivityResult
172
>
deErr := f.Get(ctx, &der)
173
>
if deErr != nil {
174
lastDeleteExecutionsActivityErr = deErr
175
return
176
}
177
>
result.SuccessCount += der.SuccessCount
workflow.go
178
>
result.ErrorCount += der.ErrorCount
179
})
180
181
>
if runningDeleteExecutionsActivityCount >= params.Config.ConcurrentDeleteExecutionsActivities {
workflow.go
182
>
// Wait for one of running activities to complete.
workflow.go
183
>
runningDeleteExecutionsSelector.Select(ctx)
184
>
if lastDeleteExecutionsActivityErr != nil {
185
return result, lastDeleteExecutionsActivityErr
186
}
187
}
188
191
}
192
}
193
194
// Wait for all running activities to complete.
195
>
for runningDeleteExecutionsActivityCount > 0 {
workflow.go
196
runningDeleteExecutionsSelector.Select(ctx)
197
if lastDeleteExecutionsActivityErr != nil {