66
)
67
68
>
func validateParams(params *ReclaimResourcesParams) error {
workflow.go
69
>
if params.NamespaceID.IsEmpty() {
70
return errors.NewInvalidArgument("namespace ID is required", nil)
71
}
73
return errors.NewInvalidArgument("namespace is required", nil)
74
}
76
>
return nil
77
}
78
79
>
func ReclaimResourcesWorkflow(ctx workflow.Context, params ReclaimResourcesParams) (ReclaimResourcesResult, error) {
workflow.go
80
>
logger := log.With(
81
>
workflow.GetLogger(ctx),
82
>
tag.WorkflowType(WorkflowName),
83
>
tag.WorkflowNamespace(params.Namespace.String()),
84
>
tag.WorkflowNamespaceID(params.NamespaceID.String()))
85
>
logger.Info("Workflow started.")
86
>
87
>
var result ReclaimResourcesResult
88
>
if err := validateParams(¶ms); err != nil {
89
return result, err
90
}
91
92
>
mh := workflow.GetMetricsHandler(ctx).WithTags(map[string]string{"namespace": params.Namespace.String()})
workflow.go
93
>
defer func() {
94
>
if result.NamespaceDeleted {
95
>
mh.Counter(metrics.ReclaimResourcesNamespaceDeleteSuccessCount.Name()).Inc(1)
workflow.go
97
mh.Counter(metrics.ReclaimResourcesNamespaceDeleteFailureCount.Name()).Inc(1)
98
}
100
>
mh.Counter(metrics.ReclaimResourcesDeleteExecutionsSuccessCount.Name()).Inc(int64(result.DeleteSuccessCount))
101
>
}
102
>
if result.DeleteErrorCount > 0 {
103
mh.Counter(metrics.ReclaimResourcesDeleteExecutionsFailureCount.Name()).Inc(int64(result.DeleteErrorCount))
104
}
105
}()
106
107
>
ctx = workflow.WithTaskQueue(ctx, primitives.DeleteNamespaceActivityTQ)
workflow.go
108
>
109
>
var (
110
>
namespaceDeleteDelay = params.NamespaceDeleteDelay
111
>
cancelDeleteDelay workflow.CancelFunc
112
>
)
113
>
err := workflow.SetUpdateHandlerWithOptions(ctx, "update_namespace_delete_delay", func(ctx workflow.Context, newNamespaceDeleteDelayStr string) (string, error) {
114
>
// This must succeed because Update validator already validated the input.
workflow.go
115
>
namespaceDeleteDelay, _ = time.ParseDuration(newNamespaceDeleteDelayStr)
116
>
117
>
var updateResult string
118
>
if namespaceDeleteDelay == 0 {
119
>
logger.Info("Namespace delete delay is removed. Namespace will be deleted immediately after all workflow executions are deleted.")
120
>
updateResult = "Namespace delete delay is removed."
121
>
} else {
122
>
logger.Info("Namespace delete delay is updated.", "new-delete-delay", namespaceDeleteDelay)
123
>
updateResult = fmt.Sprintf("Namespace delete delay is updated to %s.", namespaceDeleteDelay)
124
>
}
125
127
>
cancelDeleteDelay()
128
>
logger.Info("Existing namespace delete delay timer is cancelled.")
129
>
updateResult = "Existing namespace delete delay timer is cancelled. " + updateResult
130
>
}
131
>
return updateResult, nil
132
}, workflow.UpdateHandlerOptions{
133
>
Validator: func(_ workflow.Context, newNamespaceDeleteDelayStr string) error {
workflow.go
134
>
if newNamespaceDeleteDelayStr == "" {
135
return errors.NewInvalidArgument("delay duration is required", nil)
136
}
137
>
newDuration, err := time.ParseDuration(newNamespaceDeleteDelayStr)
workflow.go
138
>
if err != nil {
139
return errors.NewInvalidArgument("unable to parse delay duration", err)
140
}
142
return errors.NewInvalidArgument("delay duration must be positive", nil)
143
}
145
return errors.NewInvalidArgument("delay duration must be less than 30 days", nil)
146
}
148
},
149
})
151
return result, err
152
}
153
155
>
156
>
// TODO: remove version check and 9s const after v1.27 release.
157
>
namespaceCacheRefreshDelay := 9 * time.Second
158
>
v := workflow.GetVersion(ctx, "namespace-refresh-dc", workflow.DefaultVersion, 0)
159
>
if v != workflow.DefaultVersion {
160
>
// Step 0. This workflow is started right after the namespace is marked as DELETED and renamed.
161
>
// Wait for namespace cache refresh to make sure no new executions are created. 2 seconds are a random buffer.
162
>
ctx0 := workflow.WithLocalActivityOptions(ctx, localActivityOptions)
163
>
err = workflow.ExecuteLocalActivity(ctx0, la.GetNamespaceCacheRefreshInterval).Get(ctx, &namespaceCacheRefreshDelay)
164
>
if err != nil {
165
return result, err
166
}
167
}
168
169
>
err = workflow.Sleep(ctx, namespaceCacheRefreshDelay+2*time.Second)
workflow.go
170
>
if err != nil {
171
return result, err
172
}
173
174
// Step 1. Delete workflow executions.
175
>
result, err = deleteWorkflowExecutions(ctx, logger, params)
workflow.go
176
>
if err != nil {
177
return result, err
178
}
179
180
// Step 2. Sleep before deleting namespace from a database.
183
>
cancelableCtx, cancelDeleteDelay = workflow.WithCancel(ctx)
184
>
logger.Info("Delaying namespace delete. Send 'update_namespace_delete_delay' update to change or clear the delay.",
185
>
"duration", namespaceDeleteDelay.String())
186
>
187
>
ndd := namespaceDeleteDelay
188
>
namespaceDeleteDelay = 0
189
>
_ = workflow.Sleep(cancelableCtx, ndd)
190
>
}
191
192
// Step 3. Delete namespace from database.
193
>
ctx5 := workflow.WithLocalActivityOptions(ctx, localActivityOptions)
workflow.go
194
>
err = workflow.ExecuteLocalActivity(ctx5, la.DeleteNamespaceActivity, params.NamespaceID, params.Namespace).Get(ctx, nil)
195
>
if err != nil {
196
return result, err
197
}
198
200
>
logger.Info("Workflow finished successfully.")
201
>
return result, nil
202
}
203
204
>
func deleteWorkflowExecutions(ctx workflow.Context, logger log.Logger, params ReclaimResourcesParams) (ReclaimResourcesResult, error) {
workflow.go
205
>
var a *Activities
206
>
var la *LocalActivities
207
>
208
>
var result ReclaimResourcesResult
209
>
210
>
ctx1 := workflow.WithLocalActivityOptions(ctx, localActivityOptions)
211
>
212
>
// TODO: remove this code branch after v1.26 release
213
>
v := workflow.GetVersion(ctx, "remove-std-vis", workflow.DefaultVersion, 0)
214
>
if v == workflow.DefaultVersion {
215
// Standard visibility was removed from server codebase since v1.24 release. We don't need to call this local
216
// activity to know if it is advanced visibility, we know it is true. However, we need to keep it here so that