Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions docs/concepts.md
Original file line number Diff line number Diff line change
Expand Up @@ -125,6 +125,12 @@ Defines how Drained versions are cleaned up:
- **scaledownDelay**: How long to wait after a version has been Drained before scaling pods to zero
- **deleteDelay**: How long to wait after a version has been Drained before deleting the Kubernetes `Deployment`

> **NOTE**: Versions superseded before ever becoming Current or Ramping remain Inactive in Temporal;
> they never acquire a drainage timestamp. The controller retires these versions after
> their Deployment has fully scaled to zero, visibility reports no running pinned workflows,
> and Temporal accepts the normal version deletion request. The drainage-based sunset delays
> do not apply to these unused versions.

### Template
The pod template used for the target version of this worker deployment. Similar to the pod template used in a standar Kubernetes `Deployment`, but managed by the controller.

Expand Down
51 changes: 46 additions & 5 deletions internal/controller/execplan.go
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@ import (
commonpb "go.temporal.io/api/common/v1"
enumspb "go.temporal.io/api/enums/v1"
"go.temporal.io/api/serviceerror"
"go.temporal.io/api/workflowservice/v1"
sdkclient "go.temporal.io/sdk/client"
"go.temporal.io/sdk/converter"
"go.temporal.io/sdk/worker"
Expand Down Expand Up @@ -664,13 +665,20 @@ func (r *WorkerDeploymentReconciler) executeWRTOperations(
return errors.Join(append(applyErrs, statusErrs...)...)
}

// deleteDrainedVersions prunes the Temporal server-side Worker Deployment Version
// deleteDeprecatedVersions prunes the Temporal server-side Worker Deployment Version
// record for each k8s Deployment in DeleteDeployments, before executeK8sOperations
// deletes them. It is mutated to remove the k8s Deployments that should not be
// deleted because their Temporal server-side WDV record could not be removed or
// because the k8s Deployment has no build ID label. The removed k8s deployments
// stay in the cluster so that a later reconcile retries the pruning.
//
// Because DeleteVersion does not check to see if there are open workflows using
// pinned execution for Inactive versions, we query the Temporal visibility service
// for open pinned workflows before deleting Inactive versions.
// Like Temporal drainage, visibility is eventually consistent: callers must stop
// sending new pinned overrides to a version being retired. This is not an atomic
// exclusion against concurrent workflow starts or override updates.
//
// The planner only adds a drained version to DeleteDeployments once it is
// EligibleForDeletion (see planner.getDeleteDeployments): drained past the sunset
// delays with no active worker pods. Deleting the Kubernetes Deployment alone would
Expand All @@ -689,11 +697,12 @@ func (r *WorkerDeploymentReconciler) executeWRTOperations(
// phase without reaching back into k8sState. NotRegistered Deployments are also carried
// in DeleteDeployments; they have no server-side version, so they skip DeleteVersion and
// are retained for deletion.
func (r *WorkerDeploymentReconciler) deleteDrainedVersions(
func (r *WorkerDeploymentReconciler) deleteDeprecatedVersions(
ctx context.Context,
l logr.Logger,
workerDeploy *temporaliov1alpha1.WorkerDeployment,
depHandle sdkclient.WorkerDeploymentHandle,
temporalClient sdkclient.Client,
p *plan,
) {
identity := getControllerIdentity()
Expand All @@ -712,6 +721,19 @@ func (r *WorkerDeploymentReconciler) deleteDrainedVersions(
markedForDeletion = append(markedForDeletion, d)
continue
}
if slices.ContainsFunc(workerDeploy.Status.DeprecatedVersions, func(v *temporaliov1alpha1.DeprecatedWorkerDeploymentVersion) bool {
return v.BuildID == buildID && v.Status == temporaliov1alpha1.VersionStatusInactive
}) {
count, err := getOpenPinnedWorkflowExecutions(ctx, temporalClient, p.WorkerDeploymentName, buildID)
if err != nil || count == nil {
l.Info("could not confirm inactive version has no running pinned workflows, keeping its Deployment", "buildID", buildID, "error", err)
Comment thread
jaypipes marked this conversation as resolved.
continue
}
if count.Count != 0 {
l.Info("inactive version has running pinned workflows, keeping its Deployment", "buildID", buildID, "count", count.Count)
continue
}
}
backoffKey := k8s.ComputeWorkerDeploymentName(workerDeploy) + "/" + buildID
if r.skipVersionDelete(backoffKey) {
continue
Expand All @@ -733,14 +755,33 @@ func (r *WorkerDeploymentReconciler) deleteDrainedVersions(
}
l.Info("worker deployment version already deleted", "buildID", buildID)
} else {
l.Info("deleted drained worker deployment version", "buildID", buildID)
l.Info("deleted deprecated worker deployment version", "buildID", buildID)
}
r.noteVersionDeleteSuccess(backoffKey)
markedForDeletion = append(markedForDeletion, d)
}
p.DeleteDeployments = markedForDeletion
}

func getOpenPinnedWorkflowExecutions(
ctx context.Context,
temporalClient sdkclient.Client,
deploymentName string,
buildID string,
) (*workflowservice.CountWorkflowExecutionsResponse, error) {
// Visibility can contain either the legacy dot separator or the newer colon form.
legacyVersion := strings.ReplaceAll(deploymentName+"."+buildID, "'", "''")
version := strings.ReplaceAll(deploymentName+":"+buildID, "'", "''")
qs := fmt.Sprintf(
"TemporalWorkerDeploymentVersion IN ('%s', '%s') "+
"AND TemporalWorkflowVersioningBehavior = 'Pinned' "+
"AND ExecutionStatus = 'Running'",
legacyVersion, version,
)
req := &workflowservice.CountWorkflowExecutionsRequest{Query: qs}
return temporalClient.CountWorkflow(ctx, req)
}

// versionDeleteBackoff returns the per-version DeleteVersion backoff.
func (r *WorkerDeploymentReconciler) versionDeleteBackoff() *flowcontrol.Backoff {
r.deleteBackoffOnce.Do(func() {
Expand Down Expand Up @@ -795,8 +836,8 @@ func (r *WorkerDeploymentReconciler) executePlan(
// Prune the Temporal server-side version records before their k8s Deployments are
// deleted, and narrow the plan to the versions the server confirmed gone. A Deployment
// held back here keeps its version nominated for deletion, so a failed deletion is retried
// on the next reconcile instead of orphaning the record; see deleteDrainedVersions.
r.deleteDrainedVersions(ctx, l, workerDeploy, deploymentHandler, p)
// on the next reconcile instead of orphaning the record; see deleteDeprecatedVersions.
r.deleteDeprecatedVersions(ctx, l, workerDeploy, deploymentHandler, temporalClient, p)
deletedWorkerResources, err := r.executeK8sOperations(ctx, l, workerDeploy, p)
if err != nil {
return err
Expand Down
82 changes: 82 additions & 0 deletions internal/controller/execplan_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@ import (
"github.com/temporalio/temporal-worker-controller/internal/k8s"
"github.com/temporalio/temporal-worker-controller/internal/temporal"
"go.temporal.io/api/serviceerror"
"go.temporal.io/api/workflowservice/v1"
sdkclient "go.temporal.io/sdk/client"
"go.temporal.io/sdk/converter"
appsv1 "k8s.io/api/apps/v1"
Expand Down Expand Up @@ -663,3 +664,84 @@ func TestGeneratePlan_CarriesEncodingAndMessageType(t *testing.T) {
require.Equal(t, "my.package.DeployRequest", wf.messageType)
require.Equal(t, []byte(`{"service":"checkout"}`), wf.input)
}

// A visibility failure must retain the Deployment so the next reconciliation can retry.
type inactivePruneClient struct {
*stubTemporalClient
response *workflowservice.CountWorkflowExecutionsResponse
err error
query string
}

func (c *inactivePruneClient) CountWorkflow(_ context.Context, request *workflowservice.CountWorkflowExecutionsRequest) (*workflowservice.CountWorkflowExecutionsResponse, error) {
c.query = request.Query
return c.response, c.err
}

func TestExecutePlan_InactiveVersionDeletion(t *testing.T) {
for _, tc := range []struct {
name string
response *workflowservice.CountWorkflowExecutionsResponse
countErr error
deleteErr error
wantAttempt bool
wantDelete bool
}{
{name: "unused version", response: &workflowservice.CountWorkflowExecutionsResponse{}, wantAttempt: true, wantDelete: true},
{name: "running pinned workflow", response: &workflowservice.CountWorkflowExecutionsResponse{Count: 1}},
{name: "visibility failure", countErr: errors.New("visibility unavailable")},
{name: "missing response"},
{name: "server rejects deletion", response: &workflowservice.CountWorkflowExecutionsResponse{}, deleteErr: serviceerror.NewFailedPrecondition("active pollers"), wantAttempt: true},
{name: "Temporal version already deleted", response: &workflowservice.CountWorkflowExecutionsResponse{}, deleteErr: serviceerror.NewNotFound("version"), wantAttempt: true, wantDelete: true},
} {
t.Run(tc.name, func(t *testing.T) {
connection := temporaliov1alpha1.ConnectionSpec{HostPort: "test:7233"}
twd := makeExecplanTWD("my-worker", "default")
old := makeVersionedDeployment(twd, "old", 0, connection)
current := makeVersionedDeployment(twd, "current", 1, connection)
wrt := makeExecplanWRT("config", twd)
wrt.Spec.Template.Raw = []byte(`{"apiVersion":"v1","kind":"ConfigMap","data":{"test":"present"}}`)
oldConfig, oldHash := renderWRT(t, wrt, old, "old", testTemporalNamespace)
currentConfig, currentHash := renderWRT(t, wrt, current, "current", testTemporalNamespace)
// This resource outlived its Deployment and must still be cleaned up,
// even when the old version's Deployment and resources are retained.
orphan := makeVersionedDeployment(twd, "orphan", 0, connection)
orphanConfig, orphanHash := renderWRT(t, wrt, orphan, "orphan", testTemporalNamespace)
wrt.Status.Versions = []temporaliov1alpha1.WorkerResourceTemplateVersionStatus{
k8s.WorkerResourceTemplateVersionStatusForBuildID("old", oldConfig.GetName(), 1, oldHash, ""),
k8s.WorkerResourceTemplateVersionStatusForBuildID("current", currentConfig.GetName(), 1, currentHash, ""),
k8s.WorkerResourceTemplateVersionStatusForBuildID("orphan", orphanConfig.GetName(), 1, orphanHash, ""),
}
r, _ := newTestReconciler([]client.Object{twd, old, current, wrt, oldConfig, currentConfig, orphanConfig})
handle := newPruneStubHandle(tc.deleteErr)
c := &inactivePruneClient{stubTemporalClient: newStubTemporalClientWithHandle(handle), response: tc.response, err: tc.countErr}
status := statusWithDeprecated("current", current, &temporaliov1alpha1.DeprecatedWorkerDeploymentVersion{
BaseWorkerDeploymentVersion: baseVersion("old", old, temporaliov1alpha1.VersionStatusInactive),
})
twd.Status = status
state := &temporal.TemporalWorkerState{Versions: map[string]*temporal.VersionInfo{
"old": {Status: temporaliov1alpha1.VersionStatusInactive},
"current": {Status: temporaliov1alpha1.VersionStatusCurrent},
}}
p, err := r.generatePlan(context.Background(), logr.Discard(), twd, connection, state)
require.NoError(t, err)
require.NoError(t, r.executePlan(context.Background(), logr.Discard(), twd, c, p))
require.Equal(t, tc.wantAttempt, len(handle.deletedVersions) == 1)
require.Equal(t, tc.wantDelete, len(p.DeleteDeployments) == 1)
require.Equal(t, !tc.wantDelete, deploymentExists(t, r, "default", old.Name))
require.True(t, deploymentExists(t, r, "default", current.Name))
configErr := r.Get(context.Background(), types.NamespacedName{Namespace: "default", Name: oldConfig.GetName()}, &corev1.ConfigMap{})
if tc.wantDelete {
require.True(t, apierrors.IsNotFound(configErr))
} else {
require.NoError(t, configErr, "retained version must keep its rendered resources")
}
require.NoError(t, r.Get(context.Background(), types.NamespacedName{Namespace: "default", Name: currentConfig.GetName()}, &corev1.ConfigMap{}))
orphanErr := r.Get(context.Background(), types.NamespacedName{Namespace: "default", Name: orphanConfig.GetName()}, &corev1.ConfigMap{})
require.True(t, apierrors.IsNotFound(orphanErr), "orphan cleanup must continue when another version is retained")
require.Contains(t, c.query, "TemporalWorkerDeploymentVersion IN ('default/my-worker.old', 'default/my-worker:old')")
require.Contains(t, c.query, "TemporalWorkflowVersioningBehavior = 'Pinned'")
require.Contains(t, c.query, "ExecutionStatus = 'Running'")
})
}
}
12 changes: 11 additions & 1 deletion internal/planner/planner.go
Original file line number Diff line number Diff line change
Expand Up @@ -744,6 +744,16 @@ func getDeleteDeployments(
}

switch version.Status {
case temporaliov1alpha1.VersionStatusInactive:
// Superseded versions that never received routed traffic never become Drained.
// Wait for scale-down to finish; execution checks pinned workflows before pruning.
if foundDeploymentInTemporal && status.TargetVersion.BuildID != version.BuildID &&
(status.CurrentVersion == nil || status.CurrentVersion.BuildID != version.BuildID) &&
d.Spec.Replicas != nil && *d.Spec.Replicas == 0 &&
d.Status.ObservedGeneration >= d.Generation && d.Status.Replicas == 0 &&
(d.Status.TerminatingReplicas == nil || *d.Status.TerminatingReplicas == 0) {
deleteDeployments = append(deleteDeployments, d)
}
case temporaliov1alpha1.VersionStatusDrained:
// Deleting a deployment is only possible when:
// 1. The deployment has been drained for deleteDelay + scaledownDelay.
Expand All @@ -754,7 +764,7 @@ func getDeleteDeployments(
// reconcile as the Deployment delete: EligibleForDeletion is only
// computable while the Deployment (and thus this DeprecatedVersions
// entry) still exists, so this is the only point that can reliably
// prune it. See execplan.deleteDrainedVersions.
// prune it. See execplan.deleteDeprecatedVersions.
if version.DrainedSince != nil &&
(time.Since(version.DrainedSince.Time) > spec.SunsetStrategy.DeleteDelay.Duration+spec.SunsetStrategy.ScaledownDelay.Duration) &&
d.Spec.Replicas != nil && *d.Spec.Replicas == 0 &&
Expand Down
45 changes: 44 additions & 1 deletion internal/planner/planner_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -734,7 +734,7 @@ func TestGetDeleteDeployments(t *testing.T) {
// EligibleForDeletion: worker pods have not fully terminated (Status.Replicas > 0),
// so versioned pollers may still be registered and the Temporal-side DeleteVersion
// would fail. Deleting the Deployment now would strand that server-side version
// record with no way to retry (see execplan.deleteDrainedVersions), so hold off.
// record with no way to retry (see execplan.deleteDeprecatedVersions), so hold off.
name: "drained long enough and scaled to zero in spec, but not eligible for deletion - not deleted",
k8sState: &k8s.DeploymentState{
Deployments: map[string]*appsv1.Deployment{
Expand Down Expand Up @@ -4793,3 +4793,46 @@ func TestGetSunsetScaleDownBuildIDs(t *testing.T) {
})
}
}

func TestGetDeleteDeployments_Inactive(t *testing.T) {
for _, tc := range []struct {
name string
mutate func(*appsv1.Deployment, *temporaliov1alpha1.WorkerDeploymentStatus)
wantDelete bool
}{
{name: "superseded and fully scaled down", wantDelete: true},
{name: "scale down requested but pods remain", mutate: func(d *appsv1.Deployment, _ *temporaliov1alpha1.WorkerDeploymentStatus) { d.Status.Replicas = 1 }},
{name: "terminating pods remain", mutate: func(d *appsv1.Deployment, _ *temporaliov1alpha1.WorkerDeploymentStatus) {
n := int32(1)
d.Status.TerminatingReplicas = &n
}},
{name: "scale down not observed", mutate: func(d *appsv1.Deployment, _ *temporaliov1alpha1.WorkerDeploymentStatus) { d.Generation = 1 }},
{name: "replicas unspecified", mutate: func(d *appsv1.Deployment, _ *temporaliov1alpha1.WorkerDeploymentStatus) { d.Spec.Replicas = nil }},
{name: "replicas positive", mutate: func(d *appsv1.Deployment, _ *temporaliov1alpha1.WorkerDeploymentStatus) { *d.Spec.Replicas = 1 }},
{name: "target retained", mutate: func(_ *appsv1.Deployment, s *temporaliov1alpha1.WorkerDeploymentStatus) {
s.TargetVersion.BuildID = "old"
}},
{name: "current retained", mutate: func(_ *appsv1.Deployment, s *temporaliov1alpha1.WorkerDeploymentStatus) {
s.CurrentVersion = &temporaliov1alpha1.CurrentWorkerDeploymentVersion{BaseWorkerDeploymentVersion: temporaliov1alpha1.BaseWorkerDeploymentVersion{BuildID: "old"}}
}},
} {
t.Run(tc.name, func(t *testing.T) {
d := createDeploymentWithDefaultConnectionSpecHash(0)
s := &temporaliov1alpha1.WorkerDeploymentStatus{
DeprecatedVersions: []*temporaliov1alpha1.DeprecatedWorkerDeploymentVersion{{
BaseWorkerDeploymentVersion: temporaliov1alpha1.BaseWorkerDeploymentVersion{
BuildID: "old", Status: temporaliov1alpha1.VersionStatusInactive,
Deployment: &corev1.ObjectReference{Name: "old"},
},
}},
}
if tc.mutate != nil {
tc.mutate(d, s)
}
state := &k8s.DeploymentState{Deployments: map[string]*appsv1.Deployment{"old": d}}
deleted := getDeleteDeployments(state, s, &temporaliov1alpha1.WorkerDeploymentSpec{}, true)
assert.Equal(t, tc.wantDelete, len(deleted) == 1)
assert.Empty(t, getDeleteDeployments(state, s, &temporaliov1alpha1.WorkerDeploymentSpec{}, false))
})
}
}
6 changes: 6 additions & 0 deletions internal/tests/internal/deletion_integration_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -49,6 +49,12 @@ func runDeletionTests(
t.Run("drained-version-pruned-from-temporal-on-sunset", func(t *testing.T) {
testDrainedVersionPrunedOnSunset(t, k8sClient, ts, testNamespace)
})

for _, pinned := range []bool{false, true} {
t.Run(fmt.Sprintf("inactive-version-retirement/pinned-%t", pinned), func(t *testing.T) {
testInactiveVersionRetirement(t, k8sClient, ts, testNamespace, pinned)
})
}
}

// testDeletionSetsCurrentToUnversioned verifies the core fix: when a WD is deleted,
Expand Down
1 change: 1 addition & 0 deletions internal/tests/internal/deployment_controller.go
Original file line number Diff line number Diff line change
Expand Up @@ -177,6 +177,7 @@ func scaleDeploymentToZero(t *testing.T, ctx context.Context, k8sClient client.C
dep.Status.ReadyReplicas = 0
dep.Status.AvailableReplicas = 0
dep.Status.UpdatedReplicas = 0
dep.Status.ObservedGeneration = dep.Generation
return k8sClient.Status().Update(ctx, &dep)
}); err != nil {
t.Fatalf("failed to zero status of deployment %s: %v", name, err)
Expand Down
Loading
Loading