diff --git a/internal/controller/worker_controller.go b/internal/controller/worker_controller.go index a657fc36..361e95a5 100644 --- a/internal/controller/worker_controller.go +++ b/internal/controller/worker_controller.go @@ -534,8 +534,15 @@ func (r *WorkerDeploymentReconciler) markWRTsWDNotFound(ctx context.Context, wd // The cleanup sequence: // 1. Clear the ramping version (must happen first to avoid a split-traffic window) // 2. Set the current version to "unversioned" (empty BuildID) so new tasks route to unversioned workers -// 3. Delete all registered versions (with SkipDrainage since the WD is being removed entirely) -// 4. Delete the deployment record itself once all versions are gone +// 3. Delete the worker Deployments so their pods stop polling. This is required +// before a version can be deleted: the server rejects DeleteVersion (even with +// SkipDrainage) while a version still has active pollers. Normally this teardown +// happens in executePlan, but that path never runs during deletion, so without +// it the pods keep polling and cleanup deadlocks. Pollers linger in the server's +// cache for a few minutes after the pods die, so steps 4-5 requeue until they +// age out. +// 4. Delete all registered versions (with SkipDrainage since the WD is being removed entirely) +// 5. Delete the deployment record itself once all versions are gone func (r *WorkerDeploymentReconciler) handleDeletion( ctx context.Context, l logr.Logger, @@ -643,11 +650,38 @@ func (r *WorkerDeploymentReconciler) handleDeletion( l.Info("No current version set, skipping unversioned redirect") } - // Step 3: Delete versions that are eligible. Versions that are still draining + // Step 3: Tear down the worker Deployments backing these versions so their pods + // terminate and stop polling. This is required before a version can be deleted: + // the server rejects DeleteVersion (even with SkipDrainage=true) while a version + // still has active pollers. During normal reconciliation this teardown happens in + // executePlan, but that path is downstream of the deletion bail-out and never runs + // here. Without this step the pods keep polling, DeleteVersion keeps failing, and + // the finalizer is never removed (deadlock). Deleting them explicitly breaks it, + // rather than waiting for ownerRef GC, which is itself blocked on the finalizer. + k8sState, err := k8s.GetDeploymentState( + ctx, + r.Client, + workerDeploy.Namespace, + workerDeploy.Name, + workerDeploymentName) + if err != nil { + return fmt.Errorf("unable to list child deployments during cleanup: %w", err) + } + // Deletion order is irrelevant, but iterating over the DeploymentsByTime slice instead + // of the Deployments map gives deterministic ordering, so the log lines below stay + // stable across the 10s retries. + for _, d := range k8sState.DeploymentsByTime { + l.Info("Deleting worker deployment during cleanup", "deployment", d.Name) + if err := r.Delete(ctx, d); err != nil && !apierrors.IsNotFound(err) { + return fmt.Errorf("unable to delete worker deployment %s during cleanup (will retry): %w", d.Name, err) + } + } + + // Step 4: Delete versions that are eligible. Versions that are still draining // are force-deleted with SkipDrainage since the WD is being removed entirely. - // If any version fails to delete (e.g. active pollers), return an error so the - // reconciler requeues. Pollers disappear once pods terminate and the next - // reconciliation will succeed. + // If any version fails to delete (e.g. active pollers still draining after the + // Deployment delete above), return an error so the reconciler requeues. Pollers + // disappear once pods terminate and a subsequent reconciliation will succeed. for _, version := range resp.Info.VersionSummaries { buildID := version.Version.BuildID l.Info("Deleting worker deployment version", "buildID", buildID) @@ -660,7 +694,7 @@ func (r *WorkerDeploymentReconciler) handleDeletion( } } - // Step 4: Delete the deployment itself. This only succeeds if all versions are gone. + // Step 5: Delete the deployment itself. This only succeeds if all versions are gone. l.Info("Attempting to delete worker deployment from Temporal server", "name", workerDeploymentName) if _, err := temporalClient.WorkerDeploymentClient().Delete(ctx, sdkclient.WorkerDeploymentDeleteOptions{ Name: workerDeploymentName,