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