Files
temporal/tests/worker_deployment_version_test.go
Shivam 15f3532ea1 Add Worker Deployment and BuildID labels to (workflow,activity) task completion metrics (#11348)
## What changed

- Add `worker_deployment_name` and `worker_build_id` labels to Workflow
Task completion and failure metrics.
- Add the split labels to Activity success, failure, cancellation,
timeout, and completion-latency metrics.
- Add the split labels to Workflow Task and Activity schedule-to-start
latency and poll-time task-dispatch latency.
- Continue emitting the existing combined `worker_version` label on
task-dispatch latency for compatibility.
- Reuse `metrics.breakdownByBuildID` as the task-queue-scoped
cardinality gate for the new labels.

## Which metrics changed

Every metric below now emits both `worker_deployment_name` and
`worker_build_id` when `metrics.breakdownByBuildID` is enabled and
deployment attribution is available. When the gate is disabled or
attribution is unavailable—including unversioned tasks and timeouts
before a worker starts—both labels remain present with empty values.

### Workflow Task outcomes

- `workflow_tasks_completed`
- `failed_workflow_tasks`

### Activity outcomes

- `activity_success`
- `activity_fail`
- `activity_task_fail`
- `activity_cancel`
- `activity_timeout`
- `activity_task_timeout`

### Activity latency

- `activity_end_to_end_latency` (deprecated; use
`activity_start_to_close_latency` instead)
- `activity_start_to_close_latency`
- `activity_schedule_to_close_latency`

### Task-routing latency

- `task_schedule_to_start_latency`
- `task_dispatch_latency` (continues to emit the existing combined
`worker_version` label as well)

## Why

- To improve user-experience for worker-versioning by having more
insights
- Worker Deployment name and build ID need to be independently
filterable. Emitting them separately avoids requiring consumers to parse
the combined `worker_version` value. This PR preserves `worker_version`
on `task_dispatch_latency` for compatibility and does not remove or
deprecate it.

## How each metric was tested

The functional tests run real versioned Workflow and Activity workers
with `metrics.breakdownByBuildID` enabled. They assert that emitted
server metrics contain the actual Worker Deployment name and build
ID—not merely that the label keys exist.

| Functional scenario | Metrics verified | What is asserted |
|---|---|---|
| Workflow and Activity task dispatch |
`task_schedule_to_start_latency`, `task_dispatch_latency` | Both
Workflow and Activity task series contain real deployment/build values
when enabled. Disabled and unversioned cases emit empty values. |
| Successful Workflow and Activity | `workflow_tasks_completed`,
`activity_success`, `activity_end_to_end_latency`,
`activity_start_to_close_latency`, `activity_schedule_to_close_latency`
| Successful completion series contain real deployment/build values. |
| Terminal Activity failure | `activity_task_fail`, `activity_fail`,
`activity_end_to_end_latency`, `activity_start_to_close_latency`,
`activity_schedule_to_close_latency` | Both the failed attempt and
terminal-failure series retain real worker attribution. |
| Activity cancellation | `activity_cancel` | The cancellation series is
attributed to the worker that started the Activity. |
| Terminal Activity timeout | `activity_task_timeout`,
`activity_timeout` | Both attempt-level and terminal timeout series
retain the started worker’s deployment/build values. |
| Failed Workflow Task | `failed_workflow_tasks` | A real versioned
Workflow Task is polled and failed; the resulting series contains the
poller’s deployment/build values. |

### Local end-to-end validation

The PR server was also run locally with the Worker Deployment versioning
canary and a `bench-go` workload.

Prometheus was scraped directly to verify that:

- real `worker_deployment_name` and `worker_build_id` values were
emitted;
- versioned and unversioned series were both accepted without
inconsistent-label errors;
- `activity_success`, `activity_task_fail`,
`activity_end_to_end_latency`, `activity_start_to_close_latency`, and
`activity_schedule_to_close_latency` carried the expected real
deployment/build values.

The remaining failure, cancellation, timeout, and failed-Workflow-Task
paths are covered by the functional tests above.

NEW: Also tested each of these 13 metric families in a cloud test cell
with sample metrics pasted here:
https://grafana.tmprl-internal.cloud/d/shdb6gf/new-dashboard?orgId=1&from=2026-08-27T00:00:00.000Z&to=2026-08-27T23:59:59.000Z&timezone=utc


<!-- CURSOR_SUMMARY -->
---

> [!NOTE]
> **Medium Risk**
> Broad metrics surface area across History and Matching with
cardinality gated by config; deployment attribution on timeouts and
eager starts changes which tag values appear on existing metric names.
> 
> **Overview**
> Adds **`worker_deployment_name`** and **`worker_build_id`** to History
and Matching metrics so worker versioning can be filtered without
parsing the combined **`worker_version`** label.
**`metrics.breakdownByBuildID`** (task-queue scoped) controls whether
those tags carry real values or empty strings;
**`task_dispatch_latency`** still emits **`worker_version`** for
compatibility.
> 
> **Workflow tasks:** Completion and failure counters
(`workflow_tasks_completed`, `failed_workflow_tasks`) now go through
shared helpers that attach versioning behavior plus deployment tags from
the poller’s **`DeploymentOptions`**.
> 
> **Workflow activities:** Success, failure, cancel, timeout, and
related latency metrics get the same split labels via
**`VersioningMetricContext`** on respond paths (from request deployment
options), timer-driven timeouts (from **`LastDeploymentVersion`** when
an attempt had started), and schedule-to-start latency on
**`RecordActivityTaskStarted`**. Eager activities started during WFT
completion now record the completing worker’s deployment on the started
event.
> 
> **Standalone activities (CHASM):** Completion/timeout/cancel paths use
**`completionMetricsHandler`**, which always includes the deployment
label keys with empty values until standalone versioning exists—keeping
Prometheus label-set parity with workflow-embedded activities.
> 
> **Other:** **`AddActivityTaskStartedEvent`** clears stored deployment
when the poller is unversioned; Matching **`task_dispatch_latency`**
adds the new tags alongside **`worker_version`**.
> 
> <sup>Reviewed by [Cursor Bugbot](https://cursor.com/bugbot) for commit
1c60b05f06. Bugbot is set up for automated
code reviews on this repo. Configure
[here](https://www.cursor.com/dashboard/bugbot).</sup>
<!-- /CURSOR_SUMMARY -->
2026-08-27 16:38:22 +00:00

3970 lines
156 KiB
Go

package tests
import (
"context"
"errors"
"fmt"
"slices"
"strings"
"sync"
"testing"
"time"
"github.com/google/uuid"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
batchpb "go.temporal.io/api/batch/v1"
commonpb "go.temporal.io/api/common/v1"
computepb "go.temporal.io/api/compute/v1"
deploymentpb "go.temporal.io/api/deployment/v1"
enumspb "go.temporal.io/api/enums/v1"
"go.temporal.io/api/serviceerror"
workflowpb "go.temporal.io/api/workflow/v1"
"go.temporal.io/api/workflowservice/v1"
computeprovider "go.temporal.io/auto-scaled-workers/wci/workflow/compute_provider"
"go.temporal.io/sdk/activity"
sdkclient "go.temporal.io/sdk/client"
"go.temporal.io/sdk/temporal"
"go.temporal.io/sdk/worker"
"go.temporal.io/sdk/workflow"
deploymentspb "go.temporal.io/server/api/deployment/v1"
"go.temporal.io/server/api/matchingservice/v1"
persistencespb "go.temporal.io/server/api/persistence/v1"
"go.temporal.io/server/common/dynamicconfig"
"go.temporal.io/server/common/metrics"
"go.temporal.io/server/common/metrics/metricstest"
"go.temporal.io/server/common/testing/await"
"go.temporal.io/server/common/testing/parallelsuite"
"go.temporal.io/server/common/testing/testhooks"
"go.temporal.io/server/common/testing/testvars"
"go.temporal.io/server/common/tqid"
"go.temporal.io/server/common/worker_versioning"
"go.temporal.io/server/service/worker/workerdeployment"
"go.temporal.io/server/tests/testcore"
"google.golang.org/protobuf/proto"
"google.golang.org/protobuf/types/known/fieldmaskpb"
"google.golang.org/protobuf/types/known/timestamppb"
)
const (
maxConcurrentBatchOperations = 3
testVersionDrainageRefreshInterval = 3 * time.Second
testVersionDrainageVisibilityGracePeriod = 3 * time.Second
testLongVersionDrainageRefreshInterval = 10 * time.Second
testLongVersionDrainageVisibilityGracePeriod = 10 * time.Second
testMaxVersionsInDeployment = 4
)
type (
DeploymentVersionSuite struct {
parallelsuite.Suite[*DeploymentVersionSuite]
}
)
// TODO: this is always true. cleanup code
const useV32 = true
var (
testRandomMetadataValue = []byte("random metadata value")
)
func TestDeploymentVersionSuite(t *testing.T) {
testcore.UseSuiteScopedCluster(t) //nolint:staticcheck // SA1019: suite reuses one worker-service cluster to avoid per-test cluster churn.
parallelsuite.RunLegacySequential(t, &DeploymentVersionSuite{}) //nolint:staticcheck // SA1019: suite reuses one worker-service cluster to avoid per-test cluster churn.
}
// newTestEnv creates a TestEnv with the dynamic config and test variables this suite needs.
// Additional per-test options may be passed in opts.
func (s *DeploymentVersionSuite) newTestEnv(opts ...testcore.TestOption) *testcore.TestEnv {
baseOpts := []testcore.TestOption{
testcore.WithDynamicConfig(dynamicconfig.MatchingDeploymentWorkflowVersion, int(workerdeployment.VersionDataRevisionNumber)),
// Make sure we don't hit the rate limiter in tests
testcore.WithDynamicConfig(dynamicconfig.FrontendGlobalNamespaceNamespaceReplicationInducingAPIsRPS, 1000),
testcore.WithDynamicConfig(dynamicconfig.FrontendMaxNamespaceNamespaceReplicationInducingAPIsBurstRatioPerInstance, 1),
testcore.WithDynamicConfig(dynamicconfig.MatchingNumTaskqueueReadPartitions, 1),
testcore.WithDynamicConfig(dynamicconfig.MatchingNumTaskqueueWritePartitions, 1),
// Reduce the chance of hitting max batch job limit in tests
testcore.WithDynamicConfig(dynamicconfig.FrontendMaxConcurrentBatchOperationPerNamespace, maxConcurrentBatchOperations),
testcore.WithDynamicConfig(dynamicconfig.VersionDrainageStatusRefreshInterval, testVersionDrainageRefreshInterval),
testcore.WithDynamicConfig(dynamicconfig.VersionDrainageStatusVisibilityGracePeriod, testVersionDrainageVisibilityGracePeriod),
// Keep deployment versions short because worker-deployment system workflow IDs must fit into 255 characters (database constraint).
testcore.WithTestVars(func(tv *testvars.TestVars) *testvars.TestVars {
return tv.WithDeploymentSeries("wd").WithBuildID("b")
}),
}
return testcore.NewEnv(s.T(), append(baseOpts, opts...)...)
}
// pollFromDeployment calls PollWorkflowTaskQueue to start deployment related workflows
func (s *DeploymentVersionSuite) pollFromDeployment(ctx context.Context, env *testcore.TestEnv, tv *testvars.TestVars) {
_, _ = env.FrontendClient().PollWorkflowTaskQueue(ctx, &workflowservice.PollWorkflowTaskQueueRequest{
Namespace: env.Namespace().String(),
TaskQueue: tv.TaskQueue(),
Identity: uuid.NewString(),
DeploymentOptions: tv.WorkerDeploymentOptions(true),
})
}
func (s *DeploymentVersionSuite) pollActivityFromDeployment(ctx context.Context, env *testcore.TestEnv, tv *testvars.TestVars) {
_, _ = env.FrontendClient().PollActivityTaskQueue(ctx, &workflowservice.PollActivityTaskQueueRequest{
Namespace: env.Namespace().String(),
TaskQueue: tv.TaskQueue(),
Identity: uuid.NewString(),
DeploymentOptions: tv.WorkerDeploymentOptions(true),
})
}
func (s *DeploymentVersionSuite) describeVersion(env *testcore.TestEnv, tv *testvars.TestVars) (*workflowservice.DescribeWorkerDeploymentVersionResponse, error) {
ctx, cancel := context.WithTimeout(s.Context(), 10*time.Second)
defer cancel()
req := &workflowservice.DescribeWorkerDeploymentVersionRequest{
Namespace: env.Namespace().String(),
}
if useV32 {
req.DeploymentVersion = tv.ExternalDeploymentVersion()
} else {
req.Version = tv.DeploymentVersionString() //nolint:staticcheck // SA1019: worker versioning v0.31
}
return env.FrontendClient().DescribeWorkerDeploymentVersion(ctx, req)
}
func (s *DeploymentVersionSuite) updateMetadata(env *testcore.TestEnv, tv *testvars.TestVars, upsertEntries map[string]*commonpb.Payload, removeEntries []string) (*workflowservice.UpdateWorkerDeploymentVersionMetadataResponse, error) {
ctx, cancel := context.WithTimeout(s.Context(), 10*time.Second)
defer cancel()
req := &workflowservice.UpdateWorkerDeploymentVersionMetadataRequest{
Namespace: env.Namespace().String(),
UpsertEntries: upsertEntries,
RemoveEntries: removeEntries,
}
if useV32 {
req.DeploymentVersion = tv.ExternalDeploymentVersion()
} else {
req.Version = tv.DeploymentVersionString() //nolint:staticcheck // SA1019: worker versioning v0.31
}
return env.FrontendClient().UpdateWorkerDeploymentVersionMetadata(ctx, req)
}
func (s *DeploymentVersionSuite) startVersionWorkflow(ctx context.Context, env *testcore.TestEnv, tv *testvars.TestVars) {
go s.pollFromDeployment(ctx, env, tv)
s.waitForVersionWorkflow(ctx, env, tv)
}
func (s *DeploymentVersionSuite) startVersionWorkflowAndStopPoll(ctx context.Context, env *testcore.TestEnv, tv *testvars.TestVars) {
pollCtx, cancelPoll := context.WithCancel(ctx)
pollDone := make(chan struct{})
go func() {
defer close(pollDone)
s.pollFromDeployment(pollCtx, env, tv)
}()
s.waitForVersionWorkflow(ctx, env, tv)
cancelPoll()
await.Rcv(s.T(), pollDone)
}
func (s *DeploymentVersionSuite) waitForVersionWorkflow(ctx context.Context, env *testcore.TestEnv, tv *testvars.TestVars) {
s.EventuallyWithT(func(t *assert.CollectT) {
a := assert.New(t)
resp, err := s.describeVersion(env, tv)
if !a.NoError(err) {
return
}
// regardless of s.useV32, we want to read both version formats
a.Equal(tv.DeploymentVersionString(), resp.GetWorkerDeploymentVersionInfo().GetVersion())
a.Equal(tv.ExternalDeploymentVersion().GetDeploymentName(), resp.GetWorkerDeploymentVersionInfo().GetDeploymentVersion().GetDeploymentName())
a.Equal(tv.ExternalDeploymentVersion().GetBuildId(), resp.GetWorkerDeploymentVersionInfo().GetDeploymentVersion().GetBuildId())
a.Equal(enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_INACTIVE, resp.GetWorkerDeploymentVersionInfo().GetStatus())
newResp, err := env.FrontendClient().DescribeWorkerDeployment(ctx, &workflowservice.DescribeWorkerDeploymentRequest{
Namespace: env.Namespace().String(),
DeploymentName: tv.DeploymentSeries(),
})
if !a.NoError(err) {
return
}
var versionSummaryNames []string
var versionSummaryVersions []*deploymentpb.WorkerDeploymentVersion
for _, versionSummary := range newResp.GetWorkerDeploymentInfo().GetVersionSummaries() {
versionSummaryNames = append(versionSummaryNames, versionSummary.GetVersion())
versionSummaryVersions = append(versionSummaryVersions, versionSummary.GetDeploymentVersion())
}
a.Contains(versionSummaryNames, tv.DeploymentVersionString())
contains := slices.ContainsFunc(versionSummaryVersions, func(v *deploymentpb.WorkerDeploymentVersion) bool {
return v.GetDeploymentName() == tv.ExternalDeploymentVersion().GetDeploymentName() &&
v.GetBuildId() == tv.ExternalDeploymentVersion().GetBuildId()
})
a.True(contains)
}, time.Second*5, time.Millisecond*200)
}
func (s *DeploymentVersionSuite) startVersionWorkflowExpectFailAddVersion(ctx context.Context, env *testcore.TestEnv, tv *testvars.TestVars) {
_, err := env.FrontendClient().PollWorkflowTaskQueue(ctx, &workflowservice.PollWorkflowTaskQueueRequest{
Namespace: env.Namespace().String(),
TaskQueue: tv.TaskQueue(),
Identity: "random",
DeploymentOptions: tv.WorkerDeploymentOptions(true),
})
var resourceExhausted *serviceerror.ResourceExhausted
s.ErrorAs(err, &resourceExhausted)
s.Contains(resourceExhausted.Message, "maximum number of versions")
}
func (s *DeploymentVersionSuite) TestForceCAN_NoOpenWFS() {
env := s.newTestEnv()
// Start a version workflow
s.startVersionWorkflow(s.Context(), env, env.Tv())
// Set the version as current
err := s.setCurrent(env, env.Tv(), false)
s.NoError(err)
// ForceCAN
versionWorkflowID := workerdeployment.GenerateVersionWorkflowID(env.Tv().DeploymentSeries(), env.Tv().BuildID())
workflowExecution := &commonpb.WorkflowExecution{
WorkflowId: versionWorkflowID,
}
err = env.SendSignal(env.Namespace().String(), workflowExecution, workerdeployment.ForceCANSignalName, nil, env.Tv().ClientIdentity())
s.NoError(err)
// verifying we see our registered workers in the version deployment even after a CAN
s.EventuallyWithT(func(t *assert.CollectT) {
a := assert.New(t)
resp, err := s.describeVersion(env, env.Tv())
if !a.NoError(err) {
return
}
a.Equal(env.Tv().DeploymentVersionString(), resp.GetWorkerDeploymentVersionInfo().GetVersion()) //nolint:staticcheck // SA1019: worker versioning v0.31
a.Equal(env.Tv().ExternalDeploymentVersion().GetDeploymentName(), resp.GetWorkerDeploymentVersionInfo().GetDeploymentVersion().GetDeploymentName())
a.Equal(env.Tv().ExternalDeploymentVersion().GetBuildId(), resp.GetWorkerDeploymentVersionInfo().GetDeploymentVersion().GetBuildId())
a.Len(resp.GetVersionTaskQueues(), 1)
a.Len(resp.GetWorkerDeploymentVersionInfo().GetTaskQueueInfos(), 1)
// verify that the version state is intact even after a CAN
a.Equal(env.Tv().TaskQueue().GetName(), resp.GetVersionTaskQueues()[0].Name)
a.Equal(env.Tv().TaskQueue().GetName(), resp.GetWorkerDeploymentVersionInfo().GetTaskQueueInfos()[0].Name)
a.NotNil(resp.GetWorkerDeploymentVersionInfo().GetCurrentSinceTime())
a.NotNil(resp.GetWorkerDeploymentVersionInfo().GetRoutingChangedTime())
a.NotNil(resp.GetWorkerDeploymentVersionInfo().GetCurrentSinceTime())
a.Nil(resp.GetWorkerDeploymentVersionInfo().GetDrainageInfo())
a.Equal(enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_CURRENT, resp.GetWorkerDeploymentVersionInfo().GetStatus())
}, time.Second*10, time.Millisecond*1000)
}
func (s *DeploymentVersionSuite) TestForceCAN_WithOverrideState() {
env := s.newTestEnv()
// Start a version workflow
s.startVersionWorkflow(s.Context(), env, env.Tv())
// Create a modified state with metadata to verify override works
overrideState := &deploymentspb.VersionLocalState{
Version: env.Tv().DeploymentVersion(),
CreateTime: timestamppb.New(time.Now()),
Status: enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_INACTIVE,
Metadata: &deploymentpb.VersionMetadata{
Entries: map[string]*commonpb.Payload{
"override-key": {Data: []byte("override-value")},
},
},
TaskQueueFamilies: map[string]*deploymentspb.VersionLocalState_TaskQueueFamilyData{
env.Tv().TaskQueue().GetName(): {
TaskQueues: map[int32]*deploymentspb.TaskQueueVersionData{
int32(enumspb.TASK_QUEUE_TYPE_WORKFLOW): {},
},
},
},
}
// Create signal args with the override state
signalArgs := &deploymentspb.ForceCANVersionSignalArgs{
OverrideState: overrideState,
}
marshaledData, err := proto.Marshal(signalArgs)
s.NoError(err)
signalPayload := &commonpb.Payloads{
Payloads: []*commonpb.Payload{
{
Metadata: map[string][]byte{
"encoding": []byte("binary/protobuf"),
},
Data: marshaledData,
},
},
}
// Send ForceCAN signal with override state
versionWorkflowID := workerdeployment.GenerateVersionWorkflowID(env.Tv().DeploymentSeries(), env.Tv().BuildID())
workflowExecution := &commonpb.WorkflowExecution{
WorkflowId: versionWorkflowID,
}
err = env.SendSignal(env.Namespace().String(), workflowExecution, workerdeployment.ForceCANSignalName, signalPayload, env.Tv().ClientIdentity())
s.NoError(err)
// Verify that the override state is used after CAN (metadata should be present)
s.EventuallyWithT(func(t *assert.CollectT) {
a := assert.New(t)
resp, err := s.describeVersion(env, env.Tv())
if !a.NoError(err) {
return
}
// Verify the metadata from override state is present
entries := resp.GetWorkerDeploymentVersionInfo().GetMetadata().GetEntries()
if !a.Len(entries, 1) {
return
}
a.Equal([]byte("override-value"), entries["override-key"].Data)
}, time.Second*10, time.Millisecond*1000)
}
func (s *DeploymentVersionSuite) TestDescribeVersion_RegisterTaskQueue() {
env := s.newTestEnv()
numberOfDeployments := 1
// Starting a deployment workflow
go s.pollFromDeployment(s.Context(), env, env.Tv())
// Querying the Deployment
s.EventuallyWithT(func(t *assert.CollectT) {
a := require.New(t)
resp, err := s.describeVersion(env, env.Tv())
a.NoError(err)
a.Equal(env.Tv().DeploymentVersionString(), resp.GetWorkerDeploymentVersionInfo().GetVersion()) //nolint:staticcheck // SA1019: worker versioning v0.31
a.Equal(env.Tv().ExternalDeploymentVersion().GetDeploymentName(), resp.GetWorkerDeploymentVersionInfo().GetDeploymentVersion().GetDeploymentName())
a.Equal(env.Tv().ExternalDeploymentVersion().GetBuildId(), resp.GetWorkerDeploymentVersionInfo().GetDeploymentVersion().GetBuildId())
a.Equal(enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_INACTIVE, resp.GetWorkerDeploymentVersionInfo().GetStatus())
a.Len(resp.GetWorkerDeploymentVersionInfo().GetTaskQueueInfos(), numberOfDeployments)
a.Equal(env.Tv().TaskQueue().GetName(), resp.GetWorkerDeploymentVersionInfo().GetTaskQueueInfos()[0].Name)
a.Len(resp.GetVersionTaskQueues(), numberOfDeployments)
a.Equal(env.Tv().TaskQueue().GetName(), resp.GetVersionTaskQueues()[0].Name)
}, time.Second*5, time.Millisecond*200)
}
func (s *DeploymentVersionSuite) TestDescribeVersion_RegisterTaskQueue_ConcurrentPollers() {
env := s.newTestEnv()
root, err := tqid.PartitionFromProto(env.Tv().TaskQueue(), env.Namespace().String(), enumspb.TASK_QUEUE_TYPE_WORKFLOW)
s.NoError(err)
// Making concurrent polls to 4 partitions, 3 polls to each
for p := range 4 {
tv2 := env.Tv().WithTaskQueue(root.TaskQueue().NormalPartition(p).RpcName())
for range 3 {
go s.pollFromDeployment(s.Context(), env, tv2)
go s.pollActivityFromDeployment(s.Context(), env, tv2)
}
}
// Querying the Worker Deployment Version
s.EventuallyWithT(func(t *assert.CollectT) {
a := assert.New(t)
resp, err := s.describeVersion(env, env.Tv())
if !a.NoError(err) {
return
}
a.Equal(env.Tv().DeploymentVersionString(), resp.GetWorkerDeploymentVersionInfo().GetVersion()) //nolint:staticcheck // SA1019: worker versioning v0.31
a.Equal(env.Tv().ExternalDeploymentVersion().GetDeploymentName(), resp.GetWorkerDeploymentVersionInfo().GetDeploymentVersion().GetDeploymentName())
a.Equal(env.Tv().ExternalDeploymentVersion().GetBuildId(), resp.GetWorkerDeploymentVersionInfo().GetDeploymentVersion().GetBuildId())
a.Equal(enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_INACTIVE, resp.GetWorkerDeploymentVersionInfo().GetStatus())
a.Len(resp.GetWorkerDeploymentVersionInfo().GetTaskQueueInfos(), 2)
a.Equal(env.Tv().TaskQueue().GetName(), resp.GetWorkerDeploymentVersionInfo().GetTaskQueueInfos()[0].Name)
a.Len(resp.GetVersionTaskQueues(), 2)
a.Equal(env.Tv().TaskQueue().GetName(), resp.GetVersionTaskQueues()[0].Name)
}, time.Second*10, time.Millisecond*1000)
}
//nolint:forbidigo
func (s *DeploymentVersionSuite) TestDrainageStatus_SetCurrentVersion_NoOpenWFs() {
env := s.newTestEnv()
tv1 := env.Tv().WithBuildIDNumber(1)
tv2 := env.Tv().WithBuildIDNumber(2)
// Start deployment workflow 1 and wait for the deployment version to exist
s.startVersionWorkflow(s.Context(), env, tv1)
// Start deployment workflow 2 and wait for the deployment version to exist
s.startVersionWorkflow(s.Context(), env, tv2)
// non-current deployments have never been used and have no drainage info
s.checkVersionDrainageAndVersionStatus(env, tv1, &deploymentpb.VersionDrainageInfo{}, enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_INACTIVE, 0)
s.checkVersionDrainageAndVersionStatus(env, tv2, &deploymentpb.VersionDrainageInfo{}, enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_INACTIVE, 0)
// SetCurrent tv1
err := s.setCurrent(env, tv1, true)
s.NoError(err)
// Both versions have no drainage info and tv1 has it's status updated to current
s.checkVersionDrainageAndVersionStatus(env, tv1, &deploymentpb.VersionDrainageInfo{}, enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_CURRENT, 0)
s.checkVersionDrainageAndVersionStatus(env, tv2, &deploymentpb.VersionDrainageInfo{}, enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_INACTIVE, 0)
baseTime := time.Now()
// SetCurrent tv2 --> tv1 starts the child drainage workflow
err = s.setCurrent(env, tv2, true)
s.NoError(err)
changed1, checked1 := s.checkVersionDrainageAndVersionStatus(env, tv1, &deploymentpb.VersionDrainageInfo{
Status: enumspb.VERSION_DRAINAGE_STATUS_DRAINING,
}, enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_DRAINING, 0)
s.Greater(changed1, baseTime)
s.GreaterOrEqual(checked1, changed1)
// tv1 should now be "drained"
changed2, checked2 := s.checkVersionDrainageAndVersionStatus(env, tv1, &deploymentpb.VersionDrainageInfo{
Status: enumspb.VERSION_DRAINAGE_STATUS_DRAINED,
}, enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_DRAINED, testVersionDrainageVisibilityGracePeriod)
s.Greater(changed2, changed1)
s.GreaterOrEqual(checked2, changed2)
}
func (s *DeploymentVersionSuite) TestDrainageStatus_SetCurrentVersion_YesOpenWFs() {
env := s.newTestEnv()
tv1 := env.Tv().WithBuildIDNumber(1)
tv2 := env.Tv().WithBuildIDNumber(2)
// start deployment workflow 1 and wait for the deployment version to exist
s.startVersionWorkflow(s.Context(), env, tv1)
// start deployment workflow 2 and wait for the deployment version to exist
s.startVersionWorkflow(s.Context(), env, tv2)
// non-current deployments have never been used and have no drainage info
s.checkVersionDrainageAndVersionStatus(env, tv1, &deploymentpb.VersionDrainageInfo{}, enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_INACTIVE, 0)
s.checkVersionDrainageAndVersionStatus(env, tv2, &deploymentpb.VersionDrainageInfo{}, enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_INACTIVE, 0)
// SetCurrent tv1
err := s.setCurrent(env, tv1, true)
s.NoError(err)
// both versions have no drainage info and tv1 has it's status updated to current
s.checkVersionDrainageAndVersionStatus(env, tv1, &deploymentpb.VersionDrainageInfo{}, enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_CURRENT, 0)
s.checkVersionDrainageAndVersionStatus(env, tv2, &deploymentpb.VersionDrainageInfo{}, enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_INACTIVE, 0)
// start a pinned workflow on v1
run := s.startPinnedWorkflow(s.Context(), env, tv1)
baseTime := time.Now()
// SetCurrent tv2 --> tv1 starts the child drainage workflow
err = s.setCurrent(env, tv2, true)
s.NoError(err)
// tv1 should now be "draining" for visibilityGracePeriod duration
changed1, checked1 := s.checkVersionDrainageAndVersionStatus(env, tv1, &deploymentpb.VersionDrainageInfo{
Status: enumspb.VERSION_DRAINAGE_STATUS_DRAINING,
}, enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_DRAINING, 0)
s.Greater(changed1, baseTime)
s.GreaterOrEqual(checked1, changed1)
// tv1 should still be "draining" for visibilityGracePeriod duration
changed2, checked2 := s.checkVersionDrainageAndVersionStatus(env, tv1, &deploymentpb.VersionDrainageInfo{
Status: enumspb.VERSION_DRAINAGE_STATUS_DRAINING,
}, enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_DRAINING, testVersionDrainageVisibilityGracePeriod)
s.Equal(changed2, changed1)
s.Greater(checked2, checked1)
// tv1 should still be "draining" after a refresh intervals
changed3, checked3 := s.checkVersionDrainageAndVersionStatus(env, tv1, &deploymentpb.VersionDrainageInfo{
Status: enumspb.VERSION_DRAINAGE_STATUS_DRAINING,
}, enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_DRAINING, testVersionDrainageRefreshInterval)
s.Equal(changed3, changed1)
s.Greater(checked3, checked2)
// terminate workflow
_, err = env.FrontendClient().TerminateWorkflowExecution(s.Context(), &workflowservice.TerminateWorkflowExecutionRequest{
Namespace: env.Namespace().String(),
WorkflowExecution: &commonpb.WorkflowExecution{
WorkflowId: run.GetID(),
RunId: run.GetRunID(),
},
Reason: "test",
Identity: tv1.ClientIdentity(),
})
s.NoError(err)
// tv1 should now be "drained"
changed4, checked4 := s.checkVersionDrainageAndVersionStatus(env, tv1, &deploymentpb.VersionDrainageInfo{
Status: enumspb.VERSION_DRAINAGE_STATUS_DRAINED,
}, enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_DRAINED, 0)
s.Greater(changed4, changed3)
s.GreaterOrEqual(checked4, changed4)
}
func (s *DeploymentVersionSuite) startVersionedWorkflow(ctx context.Context, env *testcore.TestEnv, tv *testvars.TestVars, behavior workflow.VersioningBehavior) sdkclient.WorkflowRun {
started := make(chan struct{}, 1)
wf := func(ctx workflow.Context) (string, error) {
started <- struct{}{}
workflow.GetSignalChannel(ctx, "wait").Receive(ctx, nil)
if workflow.GetInfo(ctx).Attempt == 1 {
return "", errors.New("try again") //nolint:err113
}
panic("oops")
}
wId := testcore.RandomizeStr("id")
w := worker.New(env.SdkClient(), tv.TaskQueue().String(), worker.Options{
DeploymentOptions: worker.DeploymentOptions{
Version: tv.SDKDeploymentVersion(),
UseVersioning: true,
},
Identity: wId,
})
w.RegisterWorkflowWithOptions(wf, workflow.RegisterOptions{VersioningBehavior: behavior})
s.NoError(w.Start())
defer w.Stop()
run, err := env.SdkClient().ExecuteWorkflow(ctx, sdkclient.StartWorkflowOptions{TaskQueue: tv.TaskQueue().String()}, wf)
s.NoError(err)
await.Rcv(s.T(), started)
return run
}
func (s *DeploymentVersionSuite) startPinnedWorkflow(ctx context.Context, env *testcore.TestEnv, tv *testvars.TestVars) sdkclient.WorkflowRun {
return s.startVersionedWorkflow(ctx, env, tv, workflow.VersioningBehaviorPinned)
}
func (s *DeploymentVersionSuite) startUnpinnedWorkflow(ctx context.Context, env *testcore.TestEnv, tv *testvars.TestVars) sdkclient.WorkflowRun {
return s.startVersionedWorkflow(ctx, env, tv, workflow.VersioningBehaviorAutoUpgrade)
}
func (s *DeploymentVersionSuite) TestWorkerDeploymentLatencyMetricTags() {
for _, tc := range []struct {
name string
versioned bool
breakdownMetricsByBuildID bool
}{
{
name: "versioned with breakdown enabled",
versioned: true,
breakdownMetricsByBuildID: true,
},
{
name: "versioned with breakdown disabled",
versioned: true,
breakdownMetricsByBuildID: false,
},
{
name: "unversioned with breakdown enabled",
breakdownMetricsByBuildID: true,
},
} {
s.Run(tc.name, func(s *DeploymentVersionSuite) {
env := s.newTestEnv(
testcore.WithDynamicConfig(
dynamicconfig.MetricsBreakdownByBuildID,
tc.breakdownMetricsByBuildID,
),
)
tv := env.Tv()
if tc.versioned {
s.startVersionWorkflowAndStopPoll(s.Context(), env, tv)
s.NoError(s.setCurrent(env, tv, false))
}
capture := env.StartNamespaceMetricCapture()
activity := func(context.Context) error {
return nil
}
wf := func(ctx workflow.Context) error {
ctx = workflow.WithActivityOptions(ctx, workflow.ActivityOptions{
StartToCloseTimeout: time.Minute,
})
return workflow.ExecuteActivity(ctx, activity).Get(ctx, nil)
}
workerOptions := worker.Options{
Identity: tv.WorkerIdentity(),
DisableEagerActivities: true,
}
if tc.versioned {
workerOptions.DeploymentOptions = worker.DeploymentOptions{
Version: tv.SDKDeploymentVersion(),
UseVersioning: true,
}
}
w := worker.New(env.SdkClient(), tv.TaskQueue().GetName(), workerOptions)
if tc.versioned {
w.RegisterWorkflowWithOptions(wf, workflow.RegisterOptions{
VersioningBehavior: workflow.VersioningBehaviorAutoUpgrade,
})
} else {
w.RegisterWorkflow(wf)
}
w.RegisterActivity(activity)
s.NoError(w.Start())
defer w.Stop()
run, err := env.SdkClient().ExecuteWorkflow(
s.Context(),
sdkclient.StartWorkflowOptions{
ID: tv.WorkflowID(),
TaskQueue: tv.TaskQueue().GetName(),
},
wf,
)
s.NoError(err)
s.NoError(run.Get(s.Context(), nil))
expectedDeploymentName := ""
expectedBuildID := ""
if tc.versioned && tc.breakdownMetricsByBuildID {
expectedDeploymentName = tv.DeploymentSeries()
expectedBuildID = tv.BuildID()
}
for _, metricName := range []string{
metrics.TaskScheduleToStartLatency.Name(),
metrics.TaskDispatchLatencyPerTaskQueue.Name(),
} {
seenTaskTypes := make(map[string]bool)
for _, recording := range capture.Metric(metricName) {
if recording.Tags["taskqueue"] != tv.TaskQueue().GetName() {
continue
}
taskType := recording.Tags["task_type"]
if taskType != enumspb.TASK_QUEUE_TYPE_WORKFLOW.String() &&
taskType != enumspb.TASK_QUEUE_TYPE_ACTIVITY.String() {
continue
}
s.Equal(expectedDeploymentName, recording.Tags["worker_deployment_name"])
s.Equal(expectedBuildID, recording.Tags["worker_build_id"])
seenTaskTypes[taskType] = true
}
s.True(seenTaskTypes[enumspb.TASK_QUEUE_TYPE_WORKFLOW.String()])
s.True(seenTaskTypes[enumspb.TASK_QUEUE_TYPE_ACTIVITY.String()])
}
})
}
}
func (s *DeploymentVersionSuite) TestWorkerDeploymentActivityOutcomeMetricTags() {
for _, tc := range []struct {
name string
outcome string
expectError bool
expectedMetrics []string
}{
{
name: "success",
outcome: "success",
expectedMetrics: []string{
metrics.WorkflowTasksCompleted.Name(),
metrics.ActivitySuccess.Name(),
metrics.ActivityE2ELatency.Name(),
metrics.ActivityStartToCloseLatency.Name(),
metrics.ActivityScheduleToCloseLatency.Name(),
},
},
{
name: "terminal failure",
outcome: "fail",
expectError: true,
expectedMetrics: []string{
metrics.ActivityTaskFail.Name(),
metrics.ActivityFail.Name(),
metrics.ActivityE2ELatency.Name(),
metrics.ActivityStartToCloseLatency.Name(),
metrics.ActivityScheduleToCloseLatency.Name(),
},
},
{
name: "cancellation",
outcome: "cancel",
expectError: true,
expectedMetrics: []string{metrics.ActivityCancel.Name()},
},
{
name: "terminal timeout",
outcome: "timeout",
expectError: true,
expectedMetrics: []string{
metrics.ActivityTaskTimeout.Name(),
metrics.ActivityTimeout.Name(),
},
},
} {
s.Run(tc.name, func(s *DeploymentVersionSuite) {
env, tv, capture := s.newWorkerDeploymentMetricTestEnv()
releaseTimeoutActivity := make(chan struct{})
w := worker.New(env.SdkClient(), tv.TaskQueue().GetName(), worker.Options{
DeploymentOptions: worker.DeploymentOptions{
Version: tv.SDKDeploymentVersion(),
UseVersioning: true,
},
Identity: tv.WorkerIdentity(),
DisableEagerActivities: true,
})
activityStarted := make(chan struct{}, 1)
activityFn := func(ctx context.Context) error {
activityStarted <- struct{}{}
switch tc.outcome {
case "success":
return nil
case "fail":
return errors.New("intentional activity failure") //nolint:err113
case "cancel":
heartbeatTicker := time.NewTicker(10 * time.Millisecond)
defer heartbeatTicker.Stop()
for {
activity.RecordHeartbeat(ctx)
select {
case <-ctx.Done():
return ctx.Err()
case <-heartbeatTicker.C:
}
}
case "timeout":
<-releaseTimeoutActivity
return nil
default:
return fmt.Errorf("unknown activity outcome %q", tc.outcome)
}
}
workflowFn := func(ctx workflow.Context) error {
startToCloseTimeout := time.Minute
if tc.outcome == "timeout" {
startToCloseTimeout = 200 * time.Millisecond
}
ctx = workflow.WithActivityOptions(ctx, workflow.ActivityOptions{
StartToCloseTimeout: startToCloseTimeout,
HeartbeatTimeout: time.Second,
WaitForCancellation: true,
RetryPolicy: &temporal.RetryPolicy{
MaximumAttempts: 1,
},
})
return workflow.ExecuteActivity(ctx, "worker-deployment-metric-activity").Get(ctx, nil)
}
w.RegisterWorkflowWithOptions(workflowFn, workflow.RegisterOptions{
Name: "worker-deployment-activity-outcome-metric-workflow",
VersioningBehavior: workflow.VersioningBehaviorAutoUpgrade,
})
w.RegisterActivityWithOptions(activityFn, activity.RegisterOptions{
Name: "worker-deployment-metric-activity",
})
s.NoError(w.Start())
defer w.Stop()
defer close(releaseTimeoutActivity)
run, err := env.SdkClient().ExecuteWorkflow(
s.Context(),
sdkclient.StartWorkflowOptions{
ID: tv.WorkflowID(),
TaskQueue: tv.TaskQueue().GetName(),
},
"worker-deployment-activity-outcome-metric-workflow",
)
s.NoError(err)
if tc.outcome == "cancel" {
await.Rcv(s.T(), activityStarted)
s.NoError(env.SdkClient().CancelWorkflow(s.Context(), run.GetID(), run.GetRunID()))
}
if tc.expectError {
s.Error(run.Get(s.Context(), nil))
} else {
s.NoError(run.Get(s.Context(), nil))
}
requireWorkerDeploymentMetricTags(s, capture, tv, tc.expectedMetrics...)
})
}
}
func (s *DeploymentVersionSuite) TestWorkerDeploymentFailedWorkflowTaskMetricTags() {
env, tv, capture := s.newWorkerDeploymentMetricTestEnv()
run, err := env.SdkClient().ExecuteWorkflow(
s.Context(),
sdkclient.StartWorkflowOptions{
ID: tv.WorkflowID(),
TaskQueue: tv.TaskQueue().GetName(),
VersioningOverride: &sdkclient.PinnedVersioningOverride{
Version: tv.SDKDeploymentVersion(),
},
},
"worker-deployment-failed-workflow-task-metric-workflow",
)
s.NoError(err)
defer func() {
_ = env.SdkClient().TerminateWorkflow(s.Context(), run.GetID(), run.GetRunID(), "test cleanup")
}()
pollCtx, cancelPoll := context.WithTimeout(s.Context(), 10*time.Second)
defer cancelPoll()
task, err := env.FrontendClient().PollWorkflowTaskQueue(pollCtx, &workflowservice.PollWorkflowTaskQueueRequest{
Namespace: env.Namespace().String(),
TaskQueue: tv.TaskQueue(),
Identity: tv.WorkerIdentity(),
DeploymentOptions: tv.WorkerDeploymentOptions(true),
})
s.NoError(err)
s.NotEmpty(task.GetTaskToken())
_, err = env.FrontendClient().RespondWorkflowTaskFailed(s.Context(), &workflowservice.RespondWorkflowTaskFailedRequest{
Namespace: env.Namespace().String(),
TaskToken: task.GetTaskToken(),
Cause: enumspb.WORKFLOW_TASK_FAILED_CAUSE_WORKFLOW_WORKER_UNHANDLED_FAILURE,
Identity: tv.WorkerIdentity(),
DeploymentOptions: tv.WorkerDeploymentOptions(true),
})
s.NoError(err)
requireWorkerDeploymentMetricTags(s, capture, tv, metrics.FailedWorkflowTasksCounter.Name())
}
func (s *DeploymentVersionSuite) newWorkerDeploymentMetricTestEnv() (
*testcore.TestEnv,
*testvars.TestVars,
*testcore.NamespaceMetricCapture,
) {
env := s.newTestEnv(
testcore.WithDynamicConfig(dynamicconfig.MetricsBreakdownByBuildID, true),
testcore.WithDynamicConfig(dynamicconfig.MetricsBreakdownByTaskQueue, true),
)
tv := env.Tv()
s.startVersionWorkflowAndStopPoll(s.Context(), env, tv)
s.NoError(s.setCurrent(env, tv, false))
capture := env.StartNamespaceMetricCapture()
return env, tv, capture
}
func requireWorkerDeploymentMetricTags(
s parallelsuite.Scope,
capture *testcore.NamespaceMetricCapture,
tv *testvars.TestVars,
metricNames ...string,
) {
for _, metricName := range metricNames {
await.Require(s.Context(), s.TB(), func(t *await.T) {
r := t.Require()
recordings := capture.CollectMetric(metricName, func(recording *metricstest.CapturedRecording) bool {
taskQueue, hasTaskQueue := recording.Tags["taskqueue"]
return !hasTaskQueue || taskQueue == tv.TaskQueue().GetName()
})
r.NotEmpty(recordings, "expected %s in namespace %s", metricName, tv.NamespaceName())
hasExpectedTags := false
for _, recording := range recordings {
if recording.Tags["worker_deployment_name"] == tv.DeploymentSeries() &&
recording.Tags["worker_build_id"] == tv.BuildID() {
hasExpectedTags = true
break
}
}
r.True(hasExpectedTags, "expected %s with deployment %q and build ID %q", metricName, tv.DeploymentSeries(), tv.BuildID())
}, 5*time.Second, 50*time.Millisecond)
}
}
func (s *DeploymentVersionSuite) TestVersionIgnoresDrainageSignalWhenCurrentOrRamping() {
env := s.newTestEnv()
tv1 := env.Tv().WithBuildIDNumber(1)
// Start deployment workflow 1 and wait for the deployment version to exist
s.startVersionWorkflow(s.Context(), env, tv1)
// Make it current
err := s.setCurrent(env, tv1, false)
s.NoError(err)
// Signal it to be drained. Only do this in tests.
versionWorkflowID := workerdeployment.GenerateVersionWorkflowID(tv1.DeploymentSeries(), tv1.BuildID())
workflowExecution := &commonpb.WorkflowExecution{
WorkflowId: versionWorkflowID,
}
input := &deploymentpb.VersionDrainageInfo{
Status: enumspb.VERSION_DRAINAGE_STATUS_DRAINED,
LastChangedTime: timestamppb.New(time.Now()),
LastCheckedTime: timestamppb.New(time.Now()),
}
marshaledData, err := input.Marshal()
s.NoError(err)
signalPayload := &commonpb.Payloads{
Payloads: []*commonpb.Payload{
{
Metadata: map[string][]byte{
"encoding": []byte("binary/protobuf"),
},
Data: marshaledData,
},
},
}
err = env.SendSignal(env.Namespace().String(), workflowExecution, workerdeployment.SyncDrainageSignalName, signalPayload, tv1.ClientIdentity())
s.NoError(err)
// describe version and confirm that it is not drained
// add a 3s time requirement so that it does not succeed immediately
sentSignal := time.Now()
s.EventuallyWithT(func(t *assert.CollectT) {
a := require.New(t)
a.Greater(time.Since(sentSignal), 2*time.Second)
resp, err := s.describeVersion(env, tv1)
a.NoError(err)
a.NotEqual(enumspb.VERSION_DRAINAGE_STATUS_DRAINED, resp.GetWorkerDeploymentVersionInfo().GetDrainageInfo().GetStatus())
}, time.Second*10, time.Millisecond*1000)
}
// Testing DeleteVersion
func (s *DeploymentVersionSuite) TestDeleteVersion_DeleteCurrentVersion() {
// Override the dynamic config so that we can verify we don't get any unexpected masked errors.
env := s.newTestEnv(testcore.WithDynamicConfig(dynamicconfig.FrontendMaskInternalErrorDetails, true))
tv1 := env.Tv().WithBuildIDNumber(1)
// Create a deployment version
s.startVersionWorkflow(s.Context(), env, tv1)
// Set version as current
err := s.setCurrent(env, tv1, false)
s.NoError(err)
// Deleting this version should fail since the version is current
s.tryDeleteVersion(env, tv1, fmt.Sprintf(workerdeployment.ErrVersionIsCurrentOrRamping, tv1.DeploymentVersionStringV32()), false)
// Verifying workflow is not in a locked state after an invalid delete request such as the one above. If the workflow were in a locked
// state, the passed context would have timed out making the following operation fail.
s.EventuallyWithT(func(t *assert.CollectT) {
a := require.New(t)
resp, err := env.FrontendClient().DescribeWorkerDeployment(s.Context(), &workflowservice.DescribeWorkerDeploymentRequest{
Namespace: env.Namespace().String(),
DeploymentName: tv1.DeploymentSeries(),
})
a.NoError(err)
a.Equal(tv1.DeploymentVersionString(), resp.GetWorkerDeploymentInfo().GetRoutingConfig().GetCurrentVersion()) //nolint:staticcheck // SA1019: worker versioning v0.31
a.Equal(tv1.ExternalDeploymentVersion(), resp.GetWorkerDeploymentInfo().GetRoutingConfig().GetCurrentDeploymentVersion())
}, time.Second*5, time.Millisecond*200)
}
func (s *DeploymentVersionSuite) TestDeleteVersion_DeleteRampedVersion() {
env := s.newTestEnv()
tv1 := env.Tv().WithBuildIDNumber(1)
// Create a deployment version
s.startVersionWorkflow(s.Context(), env, tv1)
// Set version as ramping
err := s.setRamping(env, tv1, 0)
s.NoError(err)
// Deleting this version should fail since the version is ramping
s.tryDeleteVersion(env, tv1, fmt.Sprintf(workerdeployment.ErrVersionIsCurrentOrRamping, tv1.DeploymentVersionStringV32()), false)
// Verifying workflow is not in a locked state after an invalid delete request such as the one above. If the workflow were in a locked
// state, the passed context would have timed out making the following operation fail.
s.EventuallyWithT(func(t *assert.CollectT) {
a := require.New(t)
resp, err := env.FrontendClient().DescribeWorkerDeployment(s.Context(), &workflowservice.DescribeWorkerDeploymentRequest{
Namespace: env.Namespace().String(),
DeploymentName: tv1.DeploymentSeries(),
})
a.NoError(err)
a.Equal(tv1.DeploymentVersionString(), resp.GetWorkerDeploymentInfo().GetRoutingConfig().GetRampingVersion()) //nolint:staticcheck // SA1019: worker versioning v0.31
a.Equal(tv1.ExternalDeploymentVersion(), resp.GetWorkerDeploymentInfo().GetRoutingConfig().GetRampingDeploymentVersion())
}, time.Second*5, time.Millisecond*200)
}
func (s *DeploymentVersionSuite) TestDeleteVersion_NoWfs() {
env := s.newTestEnv(testcore.WithDynamicConfig(dynamicconfig.PollerHistoryTTL, 500*time.Millisecond))
tv1 := env.Tv().WithBuildIDNumber(1)
// Create a deployment version
s.startVersionWorkflow(s.Context(), env, tv1)
//nolint:forbidigo
time.Sleep(2 * time.Second) // todo (Shivam): remove this after the above skip is removed
// delete should succeed
s.tryDeleteVersion(env, tv1, "", false)
// deployment version does not exist in the deployment list
s.EventuallyWithT(func(t *assert.CollectT) {
a := require.New(t)
resp, err := env.FrontendClient().DescribeWorkerDeployment(s.Context(), &workflowservice.DescribeWorkerDeploymentRequest{
Namespace: env.Namespace().String(),
DeploymentName: tv1.DeploymentSeries(),
})
a.NoError(err)
if resp != nil {
for _, vs := range resp.GetWorkerDeploymentInfo().GetVersionSummaries() {
a.NotEqual(tv1.DeploymentVersionString(), vs.Version) //nolint:staticcheck // SA1019: worker versioning v0.31
a.NotEqual(tv1.ExternalDeploymentVersion().GetDeploymentName(), vs.GetDeploymentVersion().GetDeploymentName())
a.NotEqual(tv1.ExternalDeploymentVersion().GetBuildId(), vs.GetDeploymentVersion().GetBuildId())
}
}
}, time.Second*5, time.Millisecond*200)
}
func (s *DeploymentVersionSuite) TestDeleteVersion_DrainingVersion() {
env := s.newTestEnv()
tv1 := env.Tv().WithBuildIDNumber(1)
// Start deployment workflow 1 and wait for the deployment version to exist
s.startVersionWorkflow(s.Context(), env, tv1)
// Make the version current
err := s.setCurrent(env, tv1, false)
s.NoError(err)
// Start another version workflow
tv2 := env.Tv().WithBuildIDNumber(2)
s.startVersionWorkflow(s.Context(), env, tv2)
// Setting this version to current should start the drainage workflow for version1 and make it draining
err = s.setCurrent(env, tv2, true)
s.NoError(err)
// Version should be draining
s.checkVersionDrainageAndVersionStatus(env, tv1, &deploymentpb.VersionDrainageInfo{
Status: enumspb.VERSION_DRAINAGE_STATUS_DRAINING,
LastChangedTime: nil, // don't test this now
LastCheckedTime: nil, // don't test this now
}, enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_DRAINING, 0)
// delete should fail
s.tryDeleteVersion(env, tv1, fmt.Sprintf(workerdeployment.ErrVersionIsDraining, tv1.DeploymentVersionStringV32()), false)
}
func (s *DeploymentVersionSuite) TestDeleteVersion_Drained_But_Pollers_Exist() {
env := s.newTestEnv()
tv1 := env.Tv().WithBuildIDNumber(1)
// Start deployment workflow 1 and wait for the deployment version to exist
s.startVersionWorkflow(s.Context(), env, tv1)
// Make the version current
err := s.setCurrent(env, tv1, false)
s.NoError(err)
// Start another version workflow
tv2 := env.Tv().WithBuildIDNumber(2)
s.startVersionWorkflow(s.Context(), env, tv2)
// Setting this version to current should start the drainage workflow for version1
err = s.setCurrent(env, tv2, true)
s.NoError(err)
// Signal the first version to be drained. Only do this in tests.
s.signalAndWaitForDrained(env, tv1)
// Version will bypass "drained" check but delete should still fail since we have active pollers.
s.tryDeleteVersion(env, tv1, fmt.Sprintf(workerdeployment.ErrVersionHasPollers, tv1.DeploymentVersionStringV32()), false)
}
func (s *DeploymentVersionSuite) signalAndWaitForDrained(env *testcore.TestEnv, tv *testvars.TestVars) {
versionWorkflowID := workerdeployment.GenerateVersionWorkflowID(tv.DeploymentSeries(), tv.BuildID())
workflowExecution := &commonpb.WorkflowExecution{
WorkflowId: versionWorkflowID,
}
input := &deploymentpb.VersionDrainageInfo{
Status: enumspb.VERSION_DRAINAGE_STATUS_DRAINED,
LastChangedTime: timestamppb.New(time.Now()),
LastCheckedTime: timestamppb.New(time.Now()),
}
marshaledData, err := input.Marshal()
s.NoError(err)
signalPayload := &commonpb.Payloads{
Payloads: []*commonpb.Payload{
{
Metadata: map[string][]byte{
"encoding": []byte("binary/protobuf"),
},
Data: marshaledData,
},
},
}
err = env.SendSignal(env.Namespace().String(), workflowExecution, workerdeployment.SyncDrainageSignalName, signalPayload, tv.ClientIdentity())
s.NoError(err)
// wait for drained
s.EventuallyWithT(func(t *assert.CollectT) {
resp, err := s.describeVersion(env, tv)
assert.NoError(t, err)
assert.Equal(t, enumspb.VERSION_DRAINAGE_STATUS_DRAINED, resp.GetWorkerDeploymentVersionInfo().GetDrainageInfo().GetStatus())
}, 10*time.Second, time.Second)
}
func (s *DeploymentVersionSuite) waitForPollers(env *testcore.TestEnv, tv *testvars.TestVars, moreExpectedVersions ...*testvars.TestVars) {
expectedVersionsStr := []string{tv.DeploymentVersionStringV32()}
for _, tv2 := range moreExpectedVersions {
if !tv2.ExternalDeploymentVersion().Equal(tv.ExternalDeploymentVersion()) {
expectedVersionsStr = append(expectedVersionsStr, tv2.DeploymentVersionStringV32())
}
}
s.EventuallyWithT(func(t *assert.CollectT) {
resp, err := env.FrontendClient().DescribeTaskQueue(s.Context(), &workflowservice.DescribeTaskQueueRequest{
Namespace: env.Namespace().String(),
TaskQueue: tv.TaskQueue(),
TaskQueueType: enumspb.TASK_QUEUE_TYPE_WORKFLOW,
})
require.NoError(t, err)
versionsSeen := 0
for _, poller := range resp.Pollers {
pollerV := worker_versioning.WorkerDeploymentVersionToStringV32(worker_versioning.DeploymentVersionFromOptions(poller.GetDeploymentOptions()))
for _, v := range expectedVersionsStr {
if pollerV == v {
versionsSeen++
}
}
}
require.Equal(t, len(expectedVersionsStr), versionsSeen)
}, 10*time.Second, 100*time.Millisecond)
}
func (s *DeploymentVersionSuite) waitForNoPollers(env *testcore.TestEnv, tv *testvars.TestVars, moreUnexpectedVersions ...*testvars.TestVars) {
unexpectedVersionsStr := []string{tv.DeploymentVersionStringV32()}
for _, tv2 := range moreUnexpectedVersions {
if !tv2.ExternalDeploymentVersion().Equal(tv.ExternalDeploymentVersion()) {
unexpectedVersionsStr = append(unexpectedVersionsStr, tv2.DeploymentVersionStringV32())
}
}
s.EventuallyWithT(func(t *assert.CollectT) {
resp, err := env.FrontendClient().DescribeTaskQueue(s.Context(), &workflowservice.DescribeTaskQueueRequest{
Namespace: env.Namespace().String(),
TaskQueue: tv.TaskQueue(),
TaskQueueType: enumspb.TASK_QUEUE_TYPE_WORKFLOW,
})
require.NoError(t, err)
versionsSeen := 0
fmt.Printf("Pollers: %+v\n", resp.Pollers)
for _, poller := range resp.Pollers {
pollerV := worker_versioning.WorkerDeploymentVersionToStringV32(worker_versioning.DeploymentVersionFromOptions(poller.GetDeploymentOptions()))
for _, v := range unexpectedVersionsStr {
if pollerV == v {
versionsSeen++
}
}
}
require.Equal(t, 0, versionsSeen)
}, 10*time.Second, 100*time.Millisecond)
}
func (s *DeploymentVersionSuite) TestVersionScavenger_DeleteOnAdd() {
env := s.newTestEnv(
testcore.WithDynamicConfig(dynamicconfig.PollerHistoryTTL, 3*time.Second),
testcore.WithDynamicConfig(dynamicconfig.MatchingMaxVersionsInDeployment, testMaxVersionsInDeployment),
// we don't want the version to drain in this test
testcore.WithDynamicConfig(dynamicconfig.VersionDrainageStatusVisibilityGracePeriod, 60*time.Second),
testcore.WithDynamicConfig(dynamicconfig.TaskQueueInfoByBuildIdTTL, 0),
)
// Set deployment register error backoff to zero so to speed up the test.
env.InjectHook(testhooks.NewHook(testhooks.MatchingDeploymentRegisterErrorBackoff, 0*time.Second))
tvs := make([]*testvars.TestVars, testMaxVersionsInDeployment)
// max out the versions
for i := range testMaxVersionsInDeployment {
tvs[i] = env.Tv().WithBuildIDNumber(i)
s.startVersionWorkflow(s.Context(), env, tvs[i])
}
// Make tvs[0] current
err := s.setCurrent(env, tvs[0], false)
s.NoError(err)
// Make tvs[1] current, hence tvs[0] should go to draining
err = s.setCurrent(env, tvs[1], false)
s.NoError(err)
// CI can be slow, keep sending fresh polls to ensure that auto deletion logic sees them when we want to add tvMax so it can't add.
pollContext, cancelPolls := context.WithTimeout(context.Background(), 3*time.Second)
go func() {
for i := range testMaxVersionsInDeployment {
go s.pollFromDeployment(pollContext, env, tvs[i])
}
t := time.NewTicker(time.Second)
for {
select {
case <-pollContext.Done():
return
case <-t.C:
for i := range testMaxVersionsInDeployment {
go s.pollFromDeployment(pollContext, env, tvs[i])
}
}
}
}()
s.waitForPollers(env, tvs[0], tvs...)
tvMax := env.Tv().WithBuildIDNumber(9999)
cancelPolls()
// try to add a version and it fails because none of the versions can be deleted
s.startVersionWorkflowExpectFailAddVersion(s.Context(), env, tvMax)
// this waits for no pollers from any of original versions (tvMax pollers should be fine)
s.waitForNoPollers(env, tvs[0], tvs...)
// try to add the version again, and it succeeds, after deleting tvs[2] version but not tvs[3] (both are eligible)
s.startVersionWorkflow(s.Context(), env, tvMax)
// tvs[0] is draining so can't be deleted. tvs[1] is current, so tvs[2] should be deleted.
s.EventuallyWithT(func(t *assert.CollectT) {
a := require.New(t)
resp, err := env.FrontendClient().DescribeWorkerDeployment(s.Context(), &workflowservice.DescribeWorkerDeploymentRequest{
Namespace: env.Namespace().String(),
DeploymentName: tvMax.DeploymentSeries(),
})
a.NoError(err)
var versions []string
for _, vs := range resp.GetWorkerDeploymentInfo().GetVersionSummaries() {
versions = append(versions, vs.Version) //nolint:staticcheck // SA1019: worker versioning v0.31
}
a.NotContains(versions, tvs[2].DeploymentVersionString())
a.Contains(versions, tvs[0].DeploymentVersionString())
a.Contains(versions, tvs[1].DeploymentVersionString())
a.Contains(versions, tvs[3].DeploymentVersionString())
}, time.Second*5, time.Millisecond*200)
}
func (s *DeploymentVersionSuite) TestDeleteVersion_ValidDelete() {
env := s.newTestEnv(testcore.WithDynamicConfig(dynamicconfig.PollerHistoryTTL, 500*time.Millisecond))
tv1 := env.Tv().WithBuildIDNumber(1)
// Start deployment workflow 1 and wait for the deployment version to exist
s.startVersionWorkflow(s.Context(), env, tv1)
// Signal the first version to be drained. Only do this in tests.
s.signalAndWaitForDrained(env, tv1)
// Wait for pollers going away
s.waitForNoPollers(env, tv1, tv1)
// delete succeeds
s.tryDeleteVersion(env, tv1, "", false)
// deployment version does not exist in the deployment list
s.EventuallyWithT(func(t *assert.CollectT) {
a := require.New(t)
resp, err := env.FrontendClient().DescribeWorkerDeployment(s.Context(), &workflowservice.DescribeWorkerDeploymentRequest{
Namespace: env.Namespace().String(),
DeploymentName: tv1.DeploymentSeries(),
})
a.NoError(err)
if resp != nil {
for _, vs := range resp.GetWorkerDeploymentInfo().GetVersionSummaries() {
a.NotEqual(tv1.DeploymentVersionString(), vs.Version) //nolint:staticcheck // SA1019: worker versioning v0.31
a.NotEqual(tv1.ExternalDeploymentVersion().GetDeploymentName(), vs.GetDeploymentVersion().GetDeploymentName())
a.NotEqual(tv1.ExternalDeploymentVersion().GetBuildId(), vs.GetDeploymentVersion().GetBuildId())
}
}
}, time.Second*5, time.Millisecond*200)
// idempotency check: deleting the same version again should succeed
s.tryDeleteVersion(env, tv1, "", false)
}
func (s *DeploymentVersionSuite) skipBeforeVersion(version workerdeployment.DeploymentWorkflowVersion) {
if workerdeployment.VersionDataRevisionNumber < version {
s.T().Skipf("test supports version %v and newer", version)
}
}
func (s *DeploymentVersionSuite) TestDeleteVersion_ValidDelete_SkipDrainage() {
env := s.newTestEnv(testcore.WithDynamicConfig(dynamicconfig.PollerHistoryTTL, 500*time.Millisecond))
tv1 := env.Tv().WithBuildIDNumber(1)
// Start deployment workflow 1 and wait for the deployment version to exist
s.startVersionWorkflow(s.Context(), env, tv1)
// Wait for pollers going away
s.EventuallyWithT(func(t *assert.CollectT) {
resp, err := env.FrontendClient().DescribeTaskQueue(s.Context(), &workflowservice.DescribeTaskQueueRequest{
Namespace: env.Namespace().String(),
TaskQueue: tv1.TaskQueue(),
TaskQueueType: enumspb.TASK_QUEUE_TYPE_WORKFLOW,
})
require.NoError(t, err)
require.Empty(t, resp.Pollers)
}, 5*time.Second, time.Second)
// skipDrainage=true will make delete succeed
s.tryDeleteVersion(env, tv1, "", false)
// deployment version does not exist in the deployment list
s.EventuallyWithT(func(t *assert.CollectT) {
a := require.New(t)
resp, err := env.FrontendClient().DescribeWorkerDeployment(s.Context(), &workflowservice.DescribeWorkerDeploymentRequest{
Namespace: env.Namespace().String(),
DeploymentName: tv1.DeploymentSeries(),
})
a.NoError(err)
if resp != nil {
for _, vs := range resp.GetWorkerDeploymentInfo().GetVersionSummaries() {
a.NotEqual(tv1.DeploymentVersionString(), vs.Version) //nolint:staticcheck // SA1019: worker versioning v0.31
a.NotEqual(tv1.ExternalDeploymentVersion().GetDeploymentName(), vs.GetDeploymentVersion().GetDeploymentName())
a.NotEqual(tv1.ExternalDeploymentVersion().GetBuildId(), vs.GetDeploymentVersion().GetBuildId())
}
}
}, time.Second*5, time.Millisecond*200)
// idempotency check: deleting the same version again should succeed
s.tryDeleteVersion(env, tv1, "", false)
// Describe Worker Deployment should give not found
// describe deployment version gives not found error
s.EventuallyWithT(func(t *assert.CollectT) {
a := assert.New(t)
_, err := s.describeVersion(env, tv1)
a.Error(err)
var nfe *serviceerror.NotFound
a.ErrorAs(err, &nfe)
}, time.Second*5, time.Millisecond*200)
}
func (s *DeploymentVersionSuite) TestDeleteVersion_ConcurrentDeleteVersion() {
env := s.newTestEnv(testcore.WithDynamicConfig(dynamicconfig.PollerHistoryTTL, 500*time.Millisecond))
tv1 := env.Tv().WithBuildIDNumber(1)
// Start deployment workflow 1 and wait for the deployment version to exist
s.startVersionWorkflow(s.Context(), env, tv1)
// Wait for pollers going away
s.EventuallyWithT(func(t *assert.CollectT) {
resp, err := env.FrontendClient().DescribeTaskQueue(s.Context(), &workflowservice.DescribeTaskQueueRequest{
Namespace: env.Namespace().String(),
TaskQueue: tv1.TaskQueue(),
TaskQueueType: enumspb.TASK_QUEUE_TYPE_WORKFLOW,
})
require.NoError(t, err)
require.Empty(t, resp.Pollers)
}, 10*time.Second, time.Second)
// concurrent delete version requests should not break the system.
var wg sync.WaitGroup
wg.Add(2)
go func() {
defer wg.Done()
s.tryDeleteVersion(env, tv1, "", false) //nolint:testifylint // concurrent delete requests are expected to all succeed
}()
go func() {
defer wg.Done()
s.tryDeleteVersion(env, tv1, "", false) //nolint:testifylint // concurrent delete requests are expected to all succeed
}()
wg.Wait()
// deployment version does not exist in the deployment list
s.EventuallyWithT(func(t *assert.CollectT) {
a := require.New(t)
resp, err := env.FrontendClient().DescribeWorkerDeployment(s.Context(), &workflowservice.DescribeWorkerDeploymentRequest{
Namespace: env.Namespace().String(),
DeploymentName: tv1.DeploymentSeries(),
})
a.NoError(err)
if resp != nil {
for _, vs := range resp.GetWorkerDeploymentInfo().GetVersionSummaries() {
a.NotEqual(tv1.DeploymentVersionString(), vs.Version) //nolint:staticcheck // SA1019: worker versioning v0.31
a.NotEqual(tv1.ExternalDeploymentVersion().GetDeploymentName(), vs.GetDeploymentVersion().GetDeploymentName())
a.NotEqual(tv1.ExternalDeploymentVersion().GetBuildId(), vs.GetDeploymentVersion().GetBuildId())
}
}
}, time.Second*10, time.Millisecond*200)
}
// VersionMissingTaskQueues
func (s *DeploymentVersionSuite) TestVersionMissingTaskQueues_InvalidSetCurrentVersion() {
// Override the dynamic config to verify we don't get any unexpected masked errors.
env := s.newTestEnv(testcore.WithDynamicConfig(dynamicconfig.FrontendMaskInternalErrorDetails, true))
tv1 := env.Tv().WithBuildIDNumber(1).WithTaskQueue(env.Tv().Any().String())
// Start deployment workflow 1 and wait for the deployment version to exist
pollerCtx1, pollerCancel1 := context.WithCancel(s.Context())
s.startVersionWorkflow(pollerCtx1, env, tv1)
// SetCurrent so that the task queue puts the version in its versions info
err := s.setCurrent(env, tv1, false)
s.NoError(err)
// new version with a different registered task-queue
tv2 := env.Tv().WithBuildIDNumber(2).WithTaskQueue(env.Tv().Any().String())
s.startVersionWorkflow(s.Context(), env, tv2)
// Cancel pollers on task_queue_1 to increase the backlog of tasks
pollerCancel1()
// Start a workflow on task_queue_1 to increase the add rate
s.startWorkflow(env, tv1, tv1.VersioningOverridePinned())
// SetCurrent tv2
err = s.setCurrent(env, tv2, false)
// SetCurrent should fail since task_queue_1 does not have a current version than the deployment's existing current version
// and it either has a backlog of tasks being present or an add rate > 0.
s.EqualError(err, fmt.Sprintf(workerdeployment.ErrCurrentVersionDoesNotHaveAllTaskQueues, tv2.DeploymentVersionStringV32()))
}
func (s *DeploymentVersionSuite) TestVersionMissingTaskQueues_ValidSetCurrentVersion() {
env := s.newTestEnv()
tv1 := env.Tv().WithBuildIDNumber(1).WithTaskQueue(env.Tv().Any().String())
s.startVersionWorkflow(s.Context(), env, tv1)
// SetCurrent so that the task queue puts the version in its versions info
err := s.setCurrent(env, tv1, false)
s.NoError(err)
// new version with a different registered task-queue
tv2 := env.Tv().WithBuildIDNumber(2).WithTaskQueue(env.Tv().Any().String())
s.startVersionWorkflow(s.Context(), env, tv2)
// SetCurrent tv2
err = s.setCurrent(env, tv2, false)
// SetCurrent tv2 should succeed as task_queue_1, despite missing from the new current version, has no backlogged tasks/add-rate > 0
s.NoError(err)
}
func (s *DeploymentVersionSuite) TestVersionMissingTaskQueues_InvalidSetRampingVersion() {
// Override the dynamic config to verify we don't get any unexpected masked errors.
env := s.newTestEnv(testcore.WithDynamicConfig(dynamicconfig.FrontendMaskInternalErrorDetails, true))
tv1 := env.Tv().WithBuildIDNumber(1).WithTaskQueue(env.Tv().Any().String())
// Start deployment workflow 1 and wait for the deployment version to exist
pollerCtx1, pollerCancel1 := context.WithCancel(s.Context())
s.startVersionWorkflow(pollerCtx1, env, tv1)
// SetCurrent so that the task queue puts the version in its versions info
err := s.setCurrent(env, tv1, false)
s.NoError(err)
// new version with a different registered task-queue
tv2 := env.Tv().WithBuildIDNumber(2).WithTaskQueue(env.Tv().Any().String())
s.startVersionWorkflow(s.Context(), env, tv2)
// Cancel pollers on task_queue_1 to increase the backlog of tasks
pollerCancel1()
// Start a workflow on task_queue_1 to increase the add rate
s.startWorkflow(env, tv1, tv1.VersioningOverridePinned())
// SetRampingVersion to tv2
err = s.setRamping(env, tv2, 0)
// SetRampingVersion should fail since task_queue_1 does not have a current version than the deployment's existing current version
// and it either has a backlog of tasks being present or an add rate > 0.
s.EqualError(err, fmt.Sprintf(workerdeployment.ErrRampingVersionDoesNotHaveAllTaskQueues, tv2.DeploymentVersionStringV32()))
}
func (s *DeploymentVersionSuite) TestVersionMissingTaskQueues_ValidSetRampingVersion() {
env := s.newTestEnv()
tv1 := env.Tv().WithBuildIDNumber(1).WithTaskQueue(env.Tv().Any().String())
// Start deployment workflow 1 and wait for the deployment version to exist
s.startVersionWorkflow(s.Context(), env, tv1)
// SetCurrent so that the task queue puts the version in its versions info
err := s.setCurrent(env, tv1, false)
s.NoError(err)
// new version with a different registered task-queue
tv2 := env.Tv().WithBuildIDNumber(2).WithTaskQueue(env.Tv().Any().String())
s.startVersionWorkflow(s.Context(), env, tv2)
// SetRampingVersion to tv2
err = s.setRamping(env, tv2, 0)
// SetRampingVersion to tv2 should succeed as task_queue_1, despite missing from the new current version, has no backlogged tasks/add-rate > 0
s.NoError(err)
}
func (s *DeploymentVersionSuite) TestUpdateVersionMetadata() {
env := s.newTestEnv()
tv1 := env.Tv().WithBuildIDNumber(1)
// Start deployment workflow 1 and wait for the deployment version to exist
s.startVersionWorkflow(s.Context(), env, tv1)
metadata := map[string]*commonpb.Payload{
"key1": {Data: testRandomMetadataValue},
"key2": {Data: testRandomMetadataValue},
}
_, err := s.updateMetadata(env, tv1, metadata, nil)
s.NoError(err)
resp, err := s.describeVersion(env, tv1)
s.NoError(err)
// validating the metadata
entries := resp.GetWorkerDeploymentVersionInfo().GetMetadata().GetEntries()
s.Len(entries, 2)
s.Equal(testRandomMetadataValue, entries["key1"].Data)
s.Equal(testRandomMetadataValue, entries["key2"].Data)
// Remove all the entries
_, err = s.updateMetadata(env, tv1, nil, []string{"key1", "key2"})
s.NoError(err)
resp, err = s.describeVersion(env, tv1)
s.NoError(err)
entries = resp.GetWorkerDeploymentVersionInfo().GetMetadata().GetEntries()
s.Empty(entries)
// update metadata for the second time with an explicit identity
metadataIdentity := tv1.Any().String()
metadataReq := &workflowservice.UpdateWorkerDeploymentVersionMetadataRequest{
Namespace: env.Namespace().String(),
UpsertEntries: metadata,
Identity: metadataIdentity,
}
if useV32 {
metadataReq.DeploymentVersion = tv1.ExternalDeploymentVersion()
} else {
metadataReq.Version = tv1.DeploymentVersionString() //nolint:staticcheck // SA1019: worker versioning v0.31
}
_, err = env.FrontendClient().UpdateWorkerDeploymentVersionMetadata(s.Context(), metadataReq)
s.NoError(err)
resp, err = s.describeVersion(env, tv1)
s.NoError(err)
// validating the metadata
entries = resp.GetWorkerDeploymentVersionInfo().GetMetadata().GetEntries()
s.Len(entries, 2)
s.Equal(testRandomMetadataValue, entries["key1"].Data)
s.Equal(testRandomMetadataValue, entries["key2"].Data)
// LastModifierIdentity should match the identity provided in the metadata update
s.Equal(metadataIdentity, resp.GetWorkerDeploymentVersionInfo().GetLastModifierIdentity())
}
// testInvokeScaler returns the scaler that scaling groups backed by the test-invoke compute
// provider must declare.
func testInvokeScaler() *computepb.ComputeScaler {
return &computepb.ComputeScaler{Type: "no-sync"}
}
func (s *DeploymentVersionSuite) createDeploymentAndVersion(
env *testcore.TestEnv,
tv *testvars.TestVars,
identity string,
computeConfig *computepb.ComputeConfig,
) {
s.T().Helper()
_, err := env.FrontendClient().CreateWorkerDeployment(s.Context(), &workflowservice.CreateWorkerDeploymentRequest{
Namespace: env.Namespace().String(),
DeploymentName: tv.DeploymentSeries(),
RequestId: tv.Any().String(),
})
s.NoError(err)
_, err = env.FrontendClient().CreateWorkerDeploymentVersion(s.Context(), &workflowservice.CreateWorkerDeploymentVersionRequest{
Namespace: env.Namespace().String(),
DeploymentVersion: tv.ExternalDeploymentVersion(),
Identity: identity,
RequestId: tv.Any().String(),
ComputeConfig: computeConfig,
})
s.NoError(err)
// Wait for version to be created.
s.EventuallyWithT(func(t *assert.CollectT) {
a := require.New(t)
descResp, err := s.describeVersion(env, tv)
a.NoError(err)
a.NotNil(descResp.GetWorkerDeploymentVersionInfo())
}, 10*time.Second, 500*time.Millisecond)
}
func (s *DeploymentVersionSuite) TestUpdateComputeConfig_Success() {
env := s.newTestEnv()
createIdentity := env.Tv().Any().String()
validProvider := computeprovider.TestInvokeComputeProviderValidComputeProvider()
s.createDeploymentAndVersion(env, env.Tv(), createIdentity, &computepb.ComputeConfig{
ScalingGroups: map[string]*computepb.ComputeConfigScalingGroup{
"sg1": {Provider: validProvider, Scaler: testInvokeScaler()},
},
})
// Update compute config with a different identity and a new scaling group.
updateIdentity := env.Tv().Any().String()
_, err := env.FrontendClient().UpdateWorkerDeploymentVersionComputeConfig(s.Context(), &workflowservice.UpdateWorkerDeploymentVersionComputeConfigRequest{
Namespace: env.Namespace().String(),
DeploymentVersion: env.Tv().ExternalDeploymentVersion(),
ComputeConfigScalingGroups: map[string]*computepb.ComputeConfigScalingGroupUpdate{
"sg2": {
ScalingGroup: &computepb.ComputeConfigScalingGroup{
TaskQueueTypes: []enumspb.TaskQueueType{enumspb.TASK_QUEUE_TYPE_ACTIVITY},
Provider: validProvider,
Scaler: testInvokeScaler(),
},
},
},
Identity: updateIdentity,
RequestId: env.Tv().Any().String(),
})
s.NoError(err)
// Verify both scaling groups exist and LastModifierIdentity is updated.
s.EventuallyWithT(func(t *assert.CollectT) {
a := require.New(t)
descResp, err := s.describeVersion(env, env.Tv())
a.NoError(err)
info := descResp.GetWorkerDeploymentVersionInfo()
a.Equal(updateIdentity, info.GetLastModifierIdentity())
a.True(proto.Equal(&computepb.ComputeConfig{
ScalingGroups: map[string]*computepb.ComputeConfigScalingGroup{
"sg1": {Provider: validProvider, Scaler: testInvokeScaler()},
"sg2": {
TaskQueueTypes: []enumspb.TaskQueueType{enumspb.TASK_QUEUE_TYPE_ACTIVITY},
Provider: validProvider,
Scaler: testInvokeScaler(),
},
},
}, info.GetComputeConfig()))
}, 10*time.Second, 500*time.Millisecond)
// Verify the compute config summary is reflected in DescribeWorkerDeployment version summaries.
s.EventuallyWithT(func(t *assert.CollectT) {
a := require.New(t)
descDeployResp, err := env.FrontendClient().DescribeWorkerDeployment(s.Context(), &workflowservice.DescribeWorkerDeploymentRequest{
Namespace: env.Namespace().String(),
DeploymentName: env.Tv().DeploymentSeries(),
})
a.NoError(err)
var versionSummary *deploymentpb.WorkerDeploymentInfo_WorkerDeploymentVersionSummary
for _, vs := range descDeployResp.GetWorkerDeploymentInfo().GetVersionSummaries() {
if vs.GetVersion() == env.Tv().DeploymentVersionString() { //nolint:staticcheck // SA1019: worker versioning v0.31
versionSummary = vs
break
}
}
a.NotNil(versionSummary, "version %s not found in DescribeWorkerDeployment", env.Tv().DeploymentVersionString())
a.True(proto.Equal(&computepb.ComputeConfigSummary{
ScalingGroups: map[string]*computepb.ComputeConfigScalingGroupSummary{
"sg1": {
ProviderType: validProvider.GetType(),
},
"sg2": {
TaskQueueTypes: []enumspb.TaskQueueType{enumspb.TASK_QUEUE_TYPE_ACTIVITY},
ProviderType: validProvider.GetType(),
},
},
}, versionSummary.GetComputeConfig()))
}, 10*time.Second, 500*time.Millisecond)
// Verify the compute config summary is reflected in ListWorkerDeployments latest version summary.
s.EventuallyWithT(func(t *assert.CollectT) {
a := require.New(t)
listResp, err := env.FrontendClient().ListWorkerDeployments(s.Context(), &workflowservice.ListWorkerDeploymentsRequest{
Namespace: env.Namespace().String(),
})
a.NoError(err)
var found *workflowservice.ListWorkerDeploymentsResponse_WorkerDeploymentSummary
for _, d := range listResp.GetWorkerDeployments() {
if d.GetName() == env.Tv().DeploymentSeries() {
found = d
break
}
}
a.NotNil(found, "deployment %s not found in ListWorkerDeployments", env.Tv().DeploymentSeries())
a.True(proto.Equal(&computepb.ComputeConfigSummary{
ScalingGroups: map[string]*computepb.ComputeConfigScalingGroupSummary{
"sg1": {
ProviderType: validProvider.GetType(),
},
"sg2": {
TaskQueueTypes: []enumspb.TaskQueueType{enumspb.TASK_QUEUE_TYPE_ACTIVITY},
ProviderType: validProvider.GetType(),
},
},
}, found.GetLatestVersionSummary().GetComputeConfig()))
}, 60*time.Second, 500*time.Millisecond)
}
func (s *DeploymentVersionSuite) TestUpdateComputeConfig_UpdateExistingGroup() {
env := s.newTestEnv()
validProvider := computeprovider.TestInvokeComputeProviderValidComputeProvider()
s.createDeploymentAndVersion(env, env.Tv(), env.Tv().Any().String(), &computepb.ComputeConfig{
ScalingGroups: map[string]*computepb.ComputeConfigScalingGroup{
"sg1": {
TaskQueueTypes: []enumspb.TaskQueueType{enumspb.TASK_QUEUE_TYPE_WORKFLOW},
Provider: validProvider,
Scaler: testInvokeScaler(),
},
},
})
// Partially update sg1's task queue types via field mask.
_, err := env.FrontendClient().UpdateWorkerDeploymentVersionComputeConfig(s.Context(), &workflowservice.UpdateWorkerDeploymentVersionComputeConfigRequest{
Namespace: env.Namespace().String(),
DeploymentVersion: env.Tv().ExternalDeploymentVersion(),
ComputeConfigScalingGroups: map[string]*computepb.ComputeConfigScalingGroupUpdate{
"sg1": {
ScalingGroup: &computepb.ComputeConfigScalingGroup{
TaskQueueTypes: []enumspb.TaskQueueType{enumspb.TASK_QUEUE_TYPE_ACTIVITY},
},
UpdateMask: &fieldmaskpb.FieldMask{Paths: []string{"task_queue_types"}},
},
},
Identity: env.Tv().Any().String(),
RequestId: env.Tv().Any().String(),
})
s.NoError(err)
// Verify task queue types changed but provider is preserved.
s.EventuallyWithT(func(t *assert.CollectT) {
a := require.New(t)
descResp, err := s.describeVersion(env, env.Tv())
a.NoError(err)
a.True(proto.Equal(&computepb.ComputeConfig{
ScalingGroups: map[string]*computepb.ComputeConfigScalingGroup{
"sg1": {
TaskQueueTypes: []enumspb.TaskQueueType{enumspb.TASK_QUEUE_TYPE_ACTIVITY},
Provider: validProvider,
Scaler: testInvokeScaler(),
},
},
}, descResp.GetWorkerDeploymentVersionInfo().GetComputeConfig()))
}, 10*time.Second, 500*time.Millisecond)
}
func (s *DeploymentVersionSuite) TestUpdateComputeConfig_RemoveScalingGroup() {
env := s.newTestEnv()
validProvider := computeprovider.TestInvokeComputeProviderValidComputeProvider()
s.createDeploymentAndVersion(env, env.Tv(), env.Tv().Any().String(), &computepb.ComputeConfig{
ScalingGroups: map[string]*computepb.ComputeConfigScalingGroup{
"sg1": {Provider: validProvider, Scaler: testInvokeScaler()},
"sg2": {
TaskQueueTypes: []enumspb.TaskQueueType{enumspb.TASK_QUEUE_TYPE_ACTIVITY},
Provider: validProvider,
Scaler: testInvokeScaler(),
},
},
})
// Remove sg1.
_, err := env.FrontendClient().UpdateWorkerDeploymentVersionComputeConfig(s.Context(), &workflowservice.UpdateWorkerDeploymentVersionComputeConfigRequest{
Namespace: env.Namespace().String(),
DeploymentVersion: env.Tv().ExternalDeploymentVersion(),
RemoveComputeConfigScalingGroups: []string{"sg1"},
Identity: env.Tv().Any().String(),
RequestId: env.Tv().Any().String(),
})
s.NoError(err)
// Verify only sg2 remains.
s.EventuallyWithT(func(t *assert.CollectT) {
a := require.New(t)
descResp, err := s.describeVersion(env, env.Tv())
a.NoError(err)
a.True(proto.Equal(&computepb.ComputeConfig{
ScalingGroups: map[string]*computepb.ComputeConfigScalingGroup{
"sg2": {
TaskQueueTypes: []enumspb.TaskQueueType{enumspb.TASK_QUEUE_TYPE_ACTIVITY},
Provider: validProvider,
Scaler: testInvokeScaler(),
},
},
}, descResp.GetWorkerDeploymentVersionInfo().GetComputeConfig()))
}, 10*time.Second, 500*time.Millisecond)
}
func (s *DeploymentVersionSuite) TestUpdateComputeConfig_VersionNotFound() {
env := s.newTestEnv()
// Create deployment but no version.
_, err := env.FrontendClient().CreateWorkerDeployment(s.Context(), &workflowservice.CreateWorkerDeploymentRequest{
Namespace: env.Namespace().String(),
DeploymentName: env.Tv().DeploymentSeries(),
RequestId: env.Tv().Any().String(),
})
s.NoError(err)
_, err = env.FrontendClient().UpdateWorkerDeploymentVersionComputeConfig(s.Context(), &workflowservice.UpdateWorkerDeploymentVersionComputeConfigRequest{
Namespace: env.Namespace().String(),
DeploymentVersion: env.Tv().ExternalDeploymentVersion(),
ComputeConfigScalingGroups: map[string]*computepb.ComputeConfigScalingGroupUpdate{
"sg1": {
ScalingGroup: &computepb.ComputeConfigScalingGroup{
Provider: computeprovider.TestInvokeComputeProviderValidComputeProvider(),
},
},
},
Identity: env.Tv().Any().String(),
RequestId: env.Tv().Any().String(),
})
s.Error(err)
var notFound *serviceerror.NotFound
s.ErrorAs(err, &notFound)
}
func (s *DeploymentVersionSuite) TestUpdateComputeConfig_InvalidProvider() {
env := s.newTestEnv()
s.createDeploymentAndVersion(env, env.Tv(), env.Tv().Any().String(), &computepb.ComputeConfig{
ScalingGroups: map[string]*computepb.ComputeConfigScalingGroup{
"sg1": {Provider: computeprovider.TestInvokeComputeProviderValidComputeProvider(), Scaler: testInvokeScaler()},
},
})
_, err := env.FrontendClient().UpdateWorkerDeploymentVersionComputeConfig(s.Context(), &workflowservice.UpdateWorkerDeploymentVersionComputeConfigRequest{
Namespace: env.Namespace().String(),
DeploymentVersion: env.Tv().ExternalDeploymentVersion(),
ComputeConfigScalingGroups: map[string]*computepb.ComputeConfigScalingGroupUpdate{
"sg2": {
ScalingGroup: &computepb.ComputeConfigScalingGroup{
TaskQueueTypes: []enumspb.TaskQueueType{enumspb.TASK_QUEUE_TYPE_ACTIVITY},
Provider: &computepb.ComputeProvider{Type: "invalid-provider"},
},
},
},
Identity: env.Tv().Any().String(),
RequestId: env.Tv().Any().String(),
})
s.Error(err)
var invalidArg *serviceerror.InvalidArgument
s.ErrorAs(err, &invalidArg)
s.Contains(invalidArg.Message, "invalid compute provider type")
// Verify compute config summary is unchanged — the failed update should not have added sg2.
s.EventuallyWithT(func(t *assert.CollectT) {
a := require.New(t)
descDeployResp, err := env.FrontendClient().DescribeWorkerDeployment(s.Context(), &workflowservice.DescribeWorkerDeploymentRequest{
Namespace: env.Namespace().String(),
DeploymentName: env.Tv().DeploymentSeries(),
})
a.NoError(err)
var versionSummary *deploymentpb.WorkerDeploymentInfo_WorkerDeploymentVersionSummary
for _, vs := range descDeployResp.GetWorkerDeploymentInfo().GetVersionSummaries() {
if vs.GetVersion() == env.Tv().DeploymentVersionString() { //nolint:staticcheck // SA1019: worker versioning v0.31
versionSummary = vs
break
}
}
a.NotNil(versionSummary, "version %s not found in deployment summaries", env.Tv().DeploymentVersionString())
a.True(proto.Equal(&computepb.ComputeConfigSummary{
ScalingGroups: map[string]*computepb.ComputeConfigScalingGroupSummary{
"sg1": {
ProviderType: computeprovider.TestInvokeComputeProviderValidComputeProvider().GetType(),
},
},
}, versionSummary.GetComputeConfig()))
}, 10*time.Second, 500*time.Millisecond)
}
func (s *DeploymentVersionSuite) TestUpdateComputeConfig_DeletedVersion() {
env := s.newTestEnv(testcore.WithDynamicConfig(dynamicconfig.PollerHistoryTTL, 500*time.Millisecond))
s.createDeploymentAndVersion(env, env.Tv(), env.Tv().Any().String(), nil)
// Delete the version (skip drainage, no pollers since we created via CreateWorkerDeploymentVersion).
s.tryDeleteVersion(env, env.Tv(), "", true)
// Try to update compute config on the deleted version.
_, err := env.FrontendClient().UpdateWorkerDeploymentVersionComputeConfig(s.Context(), &workflowservice.UpdateWorkerDeploymentVersionComputeConfigRequest{
Namespace: env.Namespace().String(),
DeploymentVersion: env.Tv().ExternalDeploymentVersion(),
ComputeConfigScalingGroups: map[string]*computepb.ComputeConfigScalingGroupUpdate{
"sg1": {
ScalingGroup: &computepb.ComputeConfigScalingGroup{
Provider: computeprovider.TestInvokeComputeProviderValidComputeProvider(),
},
},
},
Identity: env.Tv().Any().String(),
RequestId: env.Tv().Any().String(),
})
s.Error(err)
}
func (s *DeploymentVersionSuite) TestValidateComputeConfig_Valid() {
env := s.newTestEnv()
s.createDeploymentAndVersion(env, env.Tv(), env.Tv().Any().String(), nil)
_, err := env.FrontendClient().ValidateWorkerDeploymentVersionComputeConfig(s.Context(), &workflowservice.ValidateWorkerDeploymentVersionComputeConfigRequest{
Namespace: env.Namespace().String(),
DeploymentVersion: env.Tv().ExternalDeploymentVersion(),
ComputeConfigScalingGroups: map[string]*computepb.ComputeConfigScalingGroupUpdate{
"sg1": {
ScalingGroup: &computepb.ComputeConfigScalingGroup{
Provider: computeprovider.TestInvokeComputeProviderValidComputeProvider(),
Scaler: testInvokeScaler(),
},
},
},
Identity: env.Tv().Any().String(),
})
s.NoError(err)
// Verify the validation had no side effects — no compute config should exist.
descResp, err := s.describeVersion(env, env.Tv())
s.NoError(err)
s.Nil(descResp.GetWorkerDeploymentVersionInfo().GetComputeConfig())
}
func (s *DeploymentVersionSuite) TestValidateComputeConfig_InvalidProvider() {
env := s.newTestEnv()
s.createDeploymentAndVersion(env, env.Tv(), env.Tv().Any().String(), nil)
_, err := env.FrontendClient().ValidateWorkerDeploymentVersionComputeConfig(s.Context(), &workflowservice.ValidateWorkerDeploymentVersionComputeConfigRequest{
Namespace: env.Namespace().String(),
DeploymentVersion: env.Tv().ExternalDeploymentVersion(),
ComputeConfigScalingGroups: map[string]*computepb.ComputeConfigScalingGroupUpdate{
"sg1": {
ScalingGroup: &computepb.ComputeConfigScalingGroup{
Provider: &computepb.ComputeProvider{Type: "invalid-provider"},
},
},
},
Identity: env.Tv().Any().String(),
})
s.Error(err)
var invalidArg *serviceerror.InvalidArgument
s.ErrorAs(err, &invalidArg)
s.Contains(invalidArg.Message, "invalid compute provider type")
}
func (s *DeploymentVersionSuite) TestValidateComputeConfig_VersionNotFound() {
env := s.newTestEnv()
// No deployment or version created — validate should still work since
// the WCI validates independently.
_, err := env.FrontendClient().CreateWorkerDeployment(s.Context(), &workflowservice.CreateWorkerDeploymentRequest{
Namespace: env.Namespace().String(),
DeploymentName: env.Tv().DeploymentSeries(),
RequestId: env.Tv().Any().String(),
})
s.NoError(err)
_, err = env.FrontendClient().ValidateWorkerDeploymentVersionComputeConfig(s.Context(), &workflowservice.ValidateWorkerDeploymentVersionComputeConfigRequest{
Namespace: env.Namespace().String(),
DeploymentVersion: env.Tv().ExternalDeploymentVersion(),
ComputeConfigScalingGroups: map[string]*computepb.ComputeConfigScalingGroupUpdate{
"sg1": {
ScalingGroup: &computepb.ComputeConfigScalingGroup{
Provider: computeprovider.TestInvokeComputeProviderValidComputeProvider(),
Scaler: testInvokeScaler(),
},
},
},
Identity: env.Tv().Any().String(),
})
s.NoError(err)
}
func (s *DeploymentVersionSuite) checkVersionDrainageAndVersionStatus(
env *testcore.TestEnv,
tv *testvars.TestVars,
expectedDrainageInfo *deploymentpb.VersionDrainageInfo,
expectedStatus enumspb.WorkerDeploymentVersionStatus,
waitFor time.Duration,
) (changedTime, checkedTime time.Time) {
if waitFor > 0 {
// wait for the requested duration before looking at the result ( +1 sec for system latency)
time.Sleep(waitFor + 1*time.Second) //nolint:forbidigo
}
s.EventuallyWithT(func(t *assert.CollectT) {
a := assert.New(t)
resp, err := s.describeVersion(env, tv)
a.NoError(err)
dInfo := resp.GetWorkerDeploymentVersionInfo().GetDrainageInfo()
a.Equal(expectedDrainageInfo.Status, dInfo.GetStatus())
if expectedDrainageInfo.LastCheckedTime != nil {
a.Equal(expectedDrainageInfo.LastCheckedTime, dInfo.GetLastCheckedTime())
}
if expectedDrainageInfo.LastChangedTime != nil {
a.Equal(expectedDrainageInfo.LastChangedTime, dInfo.GetLastChangedTime())
}
a.Equal(expectedStatus, resp.GetWorkerDeploymentVersionInfo().GetStatus())
changedTime = dInfo.GetLastChangedTime().AsTime()
checkedTime = dInfo.GetLastCheckedTime().AsTime()
}, 15*time.Second, time.Second)
return changedTime, checkedTime
}
func (s *DeploymentVersionSuite) checkVersionStatusInDeployment(
env *testcore.TestEnv,
tv *testvars.TestVars,
expectedStatus enumspb.WorkerDeploymentVersionStatus,
) {
s.EventuallyWithT(func(t *assert.CollectT) {
a := assert.New(t)
resp, err := env.FrontendClient().DescribeWorkerDeployment(s.Context(), &workflowservice.DescribeWorkerDeploymentRequest{
Namespace: env.Namespace().String(),
DeploymentName: tv.DeploymentSeries(),
})
a.NoError(err)
found := false
for _, versionSummary := range resp.GetWorkerDeploymentInfo().GetVersionSummaries() {
if versionSummary.GetVersion() == tv.DeploymentVersionString() { //nolint:staticcheck // SA1019: worker versioning v0.31
a.Equal(expectedStatus, versionSummary.GetStatus(),
"DescribeWorkerDeployment should show version %s as %s", tv.DeploymentVersionString(), expectedStatus)
found = true
break
}
}
a.True(found, "Version %s should be found in DescribeWorkerDeployment response", tv.DeploymentVersionString())
}, 10*time.Second, 500*time.Millisecond)
}
func (s *DeploymentVersionSuite) checkDescribeWorkflowAfterOverride(
env *testcore.TestEnv,
wf *commonpb.WorkflowExecution,
expectedOverride *workflowpb.VersioningOverride,
) {
s.EventuallyWithT(func(t *assert.CollectT) {
a := require.New(t)
resp, err := env.FrontendClient().DescribeWorkflowExecution(s.Context(), &workflowservice.DescribeWorkflowExecutionRequest{
Namespace: env.Namespace().String(),
Execution: wf,
})
a.NoError(err)
a.NotNil(resp)
a.NotNil(resp.GetWorkflowExecutionInfo())
actualOverride := resp.GetWorkflowExecutionInfo().GetVersioningInfo().GetVersioningOverride()
if useV32 {
// v0.32 override
a.Equal(expectedOverride.GetAutoUpgrade(), actualOverride.GetAutoUpgrade())
a.Equalf(expectedOverride.GetPinned().GetVersion().GetBuildId(), actualOverride.GetPinned().GetVersion().GetBuildId(),
"expected pinned version build id %v, got %v", expectedOverride.GetPinned().GetVersion().GetBuildId(), actualOverride.GetPinned().GetVersion().GetBuildId())
a.Equalf(expectedOverride.GetPinned().GetVersion().GetDeploymentName(), actualOverride.GetPinned().GetVersion().GetDeploymentName(),
"expected pinned version deployment name %v, got %v", expectedOverride.GetPinned().GetVersion().GetDeploymentName(), actualOverride.GetPinned().GetVersion().GetDeploymentName())
a.Equalf(expectedOverride.GetPinned().GetBehavior(), actualOverride.GetPinned().GetBehavior(),
"expected pinned override behavior %v, got %v", expectedOverride.GetPinned().GetBehavior(), actualOverride.GetPinned().GetBehavior())
if worker_versioning.OverrideIsPinned(expectedOverride) {
a.Equal(expectedOverride.GetPinned().GetVersion().GetDeploymentName(), resp.GetWorkflowExecutionInfo().GetWorkerDeploymentName())
}
} else {
// v0.31 override
a.Equal(expectedOverride.GetBehavior().String(), actualOverride.GetBehavior().String()) //nolint:staticcheck // SA1019: worker versioning v0.31
if actualOverrideDeployment := actualOverride.GetPinnedVersion(); expectedOverride.GetPinnedVersion() != actualOverrideDeployment { //nolint:staticcheck // SA1019: worker versioning v0.31
a.Fail(fmt.Sprintf("pinned override mismatch. expected: {%s}, actual: {%s}",
expectedOverride.GetPinnedVersion(), //nolint:staticcheck // SA1019: worker versioning v0.31
actualOverrideDeployment,
))
}
if worker_versioning.OverrideIsPinned(expectedOverride) {
d, _ := worker_versioning.WorkerDeploymentVersionFromStringV31(expectedOverride.GetPinnedVersion()) //nolint:staticcheck // SA1019: worker versioning v0.31
a.Equal(d.GetDeploymentName(), resp.GetWorkflowExecutionInfo().GetWorkerDeploymentName())
}
}
}, 10*time.Second, 50*time.Millisecond)
}
func (s *DeploymentVersionSuite) checkWorkflowUpdateOptionsEventIdentity(
ctx context.Context,
env *testcore.TestEnv,
wf *commonpb.WorkflowExecution,
expectedIdentity string,
) {
s.EventuallyWithT(func(t *assert.CollectT) {
a := require.New(t)
resp, err := env.FrontendClient().GetWorkflowExecutionHistory(ctx, &workflowservice.GetWorkflowExecutionHistoryRequest{
Namespace: env.Namespace().String(),
Execution: wf,
})
a.NoError(err)
a.NotNil(resp)
events := resp.GetHistory().GetEvents()
for resp.NextPageToken != nil { // probably there won't ever be more than one page of events in these tests
resp, err = env.FrontendClient().GetWorkflowExecutionHistory(ctx, &workflowservice.GetWorkflowExecutionHistoryRequest{
Namespace: env.Namespace().String(),
Execution: wf,
NextPageToken: resp.NextPageToken,
})
a.NoError(err)
a.NotNil(resp)
events = append(events, resp.GetHistory().GetEvents()...)
}
for _, event := range events {
if event.GetEventType() == enumspb.EVENT_TYPE_WORKFLOW_EXECUTION_OPTIONS_UPDATED {
a.Equal(expectedIdentity, event.GetWorkflowExecutionOptionsUpdatedEventAttributes().GetIdentity())
}
}
}, 10*time.Second, 50*time.Millisecond)
}
func (s *DeploymentVersionSuite) checkVersionIsCurrent(ctx context.Context, env *testcore.TestEnv, tv *testvars.TestVars) {
// Querying the Deployment Version
s.EventuallyWithT(func(t *assert.CollectT) {
a := assert.New(t)
resp, err := s.describeVersion(env, tv)
if !a.NoError(err) {
return
}
a.Equal(tv.DeploymentVersionString(), resp.GetWorkerDeploymentVersionInfo().GetVersion()) //nolint:staticcheck // SA1019: worker versioning v0.31
a.Equal(tv.ExternalDeploymentVersion().GetDeploymentName(), resp.GetWorkerDeploymentVersionInfo().GetDeploymentVersion().GetDeploymentName())
a.Equal(tv.ExternalDeploymentVersion().GetBuildId(), resp.GetWorkerDeploymentVersionInfo().GetDeploymentVersion().GetBuildId())
a.NotNil(resp.GetWorkerDeploymentVersionInfo().GetCurrentSinceTime())
a.Equal(enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_CURRENT, resp.GetWorkerDeploymentVersionInfo().GetStatus())
}, time.Second*10, time.Millisecond*1000)
}
func (s *DeploymentVersionSuite) checkVersionIsRamping(ctx context.Context, env *testcore.TestEnv, tv *testvars.TestVars) {
// Querying the Deployment Version
s.EventuallyWithT(func(t *assert.CollectT) {
a := assert.New(t)
resp, err := s.describeVersion(env, tv)
if !a.NoError(err) {
return
}
a.Equal(tv.DeploymentVersionString(), resp.GetWorkerDeploymentVersionInfo().GetVersion()) //nolint:staticcheck // SA1019: worker versioning v0.31
a.Equal(tv.ExternalDeploymentVersion().GetDeploymentName(), resp.GetWorkerDeploymentVersionInfo().GetDeploymentVersion().GetDeploymentName())
a.Equal(tv.ExternalDeploymentVersion().GetBuildId(), resp.GetWorkerDeploymentVersionInfo().GetDeploymentVersion().GetBuildId())
a.NotNil(resp.GetWorkerDeploymentVersionInfo().GetRampingSinceTime())
a.Equal(enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_RAMPING, resp.GetWorkerDeploymentVersionInfo().GetStatus())
}, time.Second*10, time.Millisecond*1000)
}
func (s *DeploymentVersionSuite) setCurrent(env *testcore.TestEnv, tv *testvars.TestVars, ignoreMissingTQs bool) error {
ctx, cancel := context.WithTimeout(s.Context(), 10*time.Second)
defer cancel()
req := &workflowservice.SetWorkerDeploymentCurrentVersionRequest{
Namespace: env.Namespace().String(),
DeploymentName: tv.DeploymentSeries(),
IgnoreMissingTaskQueues: ignoreMissingTQs,
Identity: tv.ClientIdentity(),
}
if useV32 {
req.BuildId = tv.BuildID()
} else {
req.Version = tv.DeploymentVersionString() //nolint:staticcheck // SA1019: worker versioning v0.31
}
_, err := env.FrontendClient().SetWorkerDeploymentCurrentVersion(ctx, req)
if err == nil {
s.checkVersionIsCurrent(ctx, env, tv)
}
return err
}
func (s *DeploymentVersionSuite) setRamping(
env *testcore.TestEnv,
tv *testvars.TestVars,
percentage float32,
) error {
ctx, cancel := context.WithTimeout(s.Context(), 20*time.Second)
defer cancel()
v := tv.DeploymentVersionString()
bid := tv.BuildID()
req := &workflowservice.SetWorkerDeploymentRampingVersionRequest{
Namespace: env.Namespace().String(),
DeploymentName: tv.DeploymentSeries(),
Percentage: percentage,
Identity: tv.ClientIdentity(),
}
if useV32 {
req.BuildId = bid
} else {
req.Version = v //nolint:staticcheck // SA1019: worker versioning v0.31
}
_, err := env.FrontendClient().SetWorkerDeploymentRampingVersion(ctx, req)
if err == nil {
s.checkVersionIsRamping(ctx, env, tv)
}
return err
}
func (s *DeploymentVersionSuite) startWorkflow(
env *testcore.TestEnv,
tv *testvars.TestVars,
override *workflowpb.VersioningOverride,
) string {
request := &workflowservice.StartWorkflowExecutionRequest{
RequestId: tv.Any().String(),
Namespace: env.Namespace().String(),
WorkflowId: tv.WorkflowID(),
WorkflowType: tv.WorkflowType(),
TaskQueue: tv.TaskQueue(),
Identity: tv.WorkerIdentity(),
VersioningOverride: override,
}
we, err0 := env.FrontendClient().StartWorkflowExecution(s.Context(), request)
s.NoError(err0)
return we.GetRunId()
}
func (s *DeploymentVersionSuite) tryDeleteVersion(
env *testcore.TestEnv,
tv *testvars.TestVars,
expectedError string,
skipDrainage bool,
) {
req := &workflowservice.DeleteWorkerDeploymentVersionRequest{
Namespace: env.Namespace().String(),
SkipDrainage: skipDrainage,
}
if useV32 {
req.DeploymentVersion = tv.ExternalDeploymentVersion()
} else {
req.Version = tv.DeploymentVersionString() //nolint:staticcheck // SA1019: worker versioning v0.31
}
_, err := env.FrontendClient().DeleteWorkerDeploymentVersion(s.Context(), req)
if expectedError == "" {
s.NoError(err)
} else {
s.EqualErrorf(err, expectedError, err.Error())
}
}
func (s *DeploymentVersionSuite) setAndCheckOverride(env *testcore.TestEnv, tv *testvars.TestVars, override *workflowpb.VersioningOverride) {
s.setAndCheckOverrideWithExpectedOutput(env, tv, override, override)
}
func (s *DeploymentVersionSuite) setAndCheckOverrideWithExpectedOutput(env *testcore.TestEnv, tv *testvars.TestVars, inputOverride, expectedOutputOverride *workflowpb.VersioningOverride) {
optsIn := &workflowpb.WorkflowExecutionOptions{VersioningOverride: inputOverride}
optsOut := &workflowpb.WorkflowExecutionOptions{VersioningOverride: expectedOutputOverride}
// Set input override --> describe workflow shows the expected output override
updateResp, err := env.FrontendClient().UpdateWorkflowExecutionOptions(s.Context(), &workflowservice.UpdateWorkflowExecutionOptionsRequest{
Namespace: env.Namespace().String(),
WorkflowExecution: tv.WorkflowExecution(),
WorkflowExecutionOptions: optsIn,
UpdateMask: &fieldmaskpb.FieldMask{Paths: []string{"versioning_override"}},
Identity: tv.ClientIdentity(),
})
s.NoError(err)
s.True(proto.Equal(updateResp.GetWorkflowExecutionOptions(), optsOut))
s.checkDescribeWorkflowAfterOverride(env, tv.WorkflowExecution(), expectedOutputOverride)
s.checkWorkflowUpdateOptionsEventIdentity(s.Context(), env, tv.WorkflowExecution(), tv.ClientIdentity())
}
// The following tests test the VersioningOverride functionality when passed via the UpdateWorkflowExecutionOptions API.
func (s *DeploymentVersionSuite) TestUpdateWorkflowExecutionOptions_SetPinned_CacheMissAndHits() {
env := s.newTestEnv(
// TODO: remove WithWorkerService once legacy suite-scoped cluster behavior is removed.
testcore.WithWorkerService("worker-deployment version membership cache test"),
testcore.WithDynamicConfig(dynamicconfig.VersionMembershipCacheTTL, 5*time.Second),
)
// start an unversioned workflow
s.startWorkflow(env, env.Tv(), nil)
opts := &workflowpb.WorkflowExecutionOptions{VersioningOverride: s.makePinnedOverride(env.Tv())}
// Setting a pinned override should fail since the version does not exist
resp, err := env.FrontendClient().UpdateWorkflowExecutionOptions(s.Context(), &workflowservice.UpdateWorkflowExecutionOptionsRequest{
Namespace: env.Namespace().String(),
WorkflowExecution: env.Tv().WorkflowExecution(),
WorkflowExecutionOptions: opts,
UpdateMask: &fieldmaskpb.FieldMask{Paths: []string{"versioning_override"}},
Identity: env.Tv().ClientIdentity(),
})
s.Error(err)
s.Nil(resp)
// Start a versioned poller which shall create a version; however, the cache TTL is not expired yet. This would result in a cache hit which would return
// a stale value for the version presence in the task queue.
s.startVersionWorkflow(s.Context(), env, env.Tv())
// Setting a pinned override should fail since the stale cache entry is returned.
resp, err = env.FrontendClient().UpdateWorkflowExecutionOptions(s.Context(), &workflowservice.UpdateWorkflowExecutionOptionsRequest{
Namespace: env.Namespace().String(),
WorkflowExecution: env.Tv().WorkflowExecution(),
WorkflowExecutionOptions: opts,
UpdateMask: &fieldmaskpb.FieldMask{Paths: []string{"versioning_override"}},
Identity: env.Tv().ClientIdentity(),
})
s.Error(err)
s.Nil(resp)
// Wait for the cache TTL to expire
s.Eventually(func() bool {
_, err := env.FrontendClient().UpdateWorkflowExecutionOptions(s.Context(), &workflowservice.UpdateWorkflowExecutionOptionsRequest{
Namespace: env.Namespace().String(),
WorkflowExecution: env.Tv().WorkflowExecution(),
WorkflowExecutionOptions: opts,
UpdateMask: &fieldmaskpb.FieldMask{Paths: []string{"versioning_override"}},
Identity: env.Tv().ClientIdentity(),
})
return err == nil
}, 10*time.Second, 500*time.Millisecond)
// The Pinned Override should have now succeeded with no error. Verify that the
// the workflow shows the override.
s.checkDescribeWorkflowAfterOverride(env, env.Tv().WorkflowExecution(), opts.VersioningOverride)
}
func (s *DeploymentVersionSuite) TestUpdateWorkflowExecutionOptions_SetUnpinnedThenUnset() {
env := s.newTestEnv()
// start an unversioned workflow
s.startWorkflow(env, env.Tv(), nil)
// 1. Set unpinned override --> describe workflow shows the override
s.setAndCheckOverride(env, env.Tv(), s.makeAutoUpgradeOverride())
// 2. Unset using empty update opts with mutation mask --> describe workflow shows no more override
s.setAndCheckOverride(env, env.Tv(), nil)
}
func (s *DeploymentVersionSuite) TestUpdateWorkflowExecutionOptions_SetPinnedThenUnset() {
env := s.newTestEnv()
// Start a versioned poller which shall create a version; the version must be present before it can be set as an override.
s.startVersionWorkflow(s.Context(), env, env.Tv())
// start an unversioned workflow
s.startWorkflow(env, env.Tv(), nil)
// 1. Set pinned override on our new unversioned workflow --> describe workflow shows the override
s.setAndCheckOverride(env, env.Tv(), s.makePinnedOverride(env.Tv()))
// 2. Unset using empty update opts with mutation mask --> describe workflow shows no more override
s.setAndCheckOverride(env, env.Tv(), nil)
}
func (s *DeploymentVersionSuite) TestUpdateWorkflowExecutionOptions_EmptyFields() {
env := s.newTestEnv()
// Start a versioned poller which shall create a version; the version must be present before it can be set as an override.
s.startVersionWorkflow(s.Context(), env, env.Tv())
// start an unversioned workflow
s.startWorkflow(env, env.Tv(), nil)
// 1. Pinned update with empty mask --> describe workflow shows no change
updateResp, err := env.FrontendClient().UpdateWorkflowExecutionOptions(s.Context(), &workflowservice.UpdateWorkflowExecutionOptionsRequest{
Namespace: env.Namespace().String(),
WorkflowExecution: env.Tv().WorkflowExecution(),
WorkflowExecutionOptions: &workflowpb.WorkflowExecutionOptions{
VersioningOverride: s.makePinnedOverride(env.Tv()),
},
UpdateMask: &fieldmaskpb.FieldMask{Paths: []string{}},
})
s.NoError(err)
s.True(proto.Equal(updateResp.GetWorkflowExecutionOptions(), &workflowpb.WorkflowExecutionOptions{}))
s.checkDescribeWorkflowAfterOverride(env, env.Tv().WorkflowExecution(), nil)
}
func (s *DeploymentVersionSuite) TestUpdateWorkflowExecutionOptions_SetPinnedSetPinned() {
env := s.newTestEnv()
tv1 := env.Tv().WithBuildIDNumber(1)
tv2 := env.Tv().WithBuildIDNumber(2)
// Start a versioned poller which shall create the two versions; the versions must be present before they can be set as overrides.
s.startVersionWorkflow(s.Context(), env, tv1)
s.startVersionWorkflow(s.Context(), env, tv2)
// start an unversioned workflow
s.startWorkflow(env, env.Tv(), nil)
// 1. Set pinned override 1 --> describe workflow shows the override
s.setAndCheckOverride(env, env.Tv(), s.makePinnedOverride(tv1))
// 3. Set pinned override 2 --> describe workflow shows the override
s.setAndCheckOverride(env, env.Tv(), s.makePinnedOverride(tv2))
}
func (s *DeploymentVersionSuite) TestUpdateWorkflowExecutionOptions_SetImpliedPinnedSuccess() {
env := s.newTestEnv()
if !useV32 {
s.T().Skip("Implied pinned overrides are only supported in v3.2+")
}
tv1 := env.Tv().WithBuildIDNumber(1)
// Start a versioned poller which shall create the two versions; the versions must be present before they can be set as overrides.
s.startVersionWorkflow(s.Context(), env, tv1)
// Set tv1 to current, so that the test workflow will be naturally pinned to tv1.
err := s.setCurrent(env, tv1, true)
s.NoError(err)
// Start a workflow pinned to tv1.
run := s.startPinnedWorkflow(s.Context(), env, tv1)
noVersionPinnedOverride := &workflowpb.VersioningOverride{Override: &workflowpb.VersioningOverride_Pinned{
Pinned: &workflowpb.VersioningOverride_PinnedOverride{
Behavior: workflowpb.VersioningOverride_PINNED_OVERRIDE_BEHAVIOR_PINNED,
},
}}
yesVersionPinnedOverride := &workflowpb.VersioningOverride{Override: &workflowpb.VersioningOverride_Pinned{
Pinned: &workflowpb.VersioningOverride_PinnedOverride{
Behavior: workflowpb.VersioningOverride_PINNED_OVERRIDE_BEHAVIOR_PINNED,
Version: tv1.ExternalDeploymentVersion(),
},
}}
// 1. Set pinned override without a version --> describe workflow shows the override with pinned override version set to tv1.
s.setAndCheckOverrideWithExpectedOutput(env, tv1.WithWorkflowID(run.GetID()), noVersionPinnedOverride, yesVersionPinnedOverride)
}
func (s *DeploymentVersionSuite) TestUpdateWorkflowExecutionOptions_SetImpliedPinnedError() {
env := s.newTestEnv()
if !useV32 {
s.T().Skip("Implied pinned overrides are only supported in v3.2+")
}
tv1 := env.Tv().WithBuildIDNumber(1)
// Start a versioned poller which shall create the two versions; the versions must be present before they can be set as overrides.
s.startVersionWorkflow(s.Context(), env, tv1)
// Set tv1 to current, so that the test workflow will run on v1.
err := s.setCurrent(env, tv1, true)
s.NoError(err)
// Start an auto-upgrade workflow.
run := s.startUnpinnedWorkflow(s.Context(), env, tv1)
noVersionPinnedOverride := &workflowpb.VersioningOverride{Override: &workflowpb.VersioningOverride_Pinned{
Pinned: &workflowpb.VersioningOverride_PinnedOverride{
Behavior: workflowpb.VersioningOverride_PINNED_OVERRIDE_BEHAVIOR_PINNED,
},
}}
// 1. Set pinned override without a version --> errors because workflow is not already pinned to a version
optsIn := &workflowpb.WorkflowExecutionOptions{VersioningOverride: noVersionPinnedOverride}
// Set input override --> describe workflow shows the expected output override
_, err = env.FrontendClient().UpdateWorkflowExecutionOptions(s.Context(), &workflowservice.UpdateWorkflowExecutionOptionsRequest{
Namespace: env.Namespace().String(),
WorkflowExecution: tv1.WithWorkflowID(run.GetID()).WorkflowExecution(),
WorkflowExecutionOptions: optsIn,
UpdateMask: &fieldmaskpb.FieldMask{Paths: []string{"versioning_override"}},
Identity: tv1.ClientIdentity(),
})
s.Error(err)
s.Contains(err.Error(),
fmt.Sprintf("must specify a specific pinned override version because workflow with id %v has behavior %s and is not yet pinned to any version",
run.GetID(), enumspb.VERSIONING_BEHAVIOR_AUTO_UPGRADE.String()),
)
}
func (s *DeploymentVersionSuite) TestUpdateWorkflowExecutionOptions_SetUnpinnedSetUnpinned() {
env := s.newTestEnv()
// start an unversioned workflow
s.startWorkflow(env, env.Tv(), nil)
// 1. Set unpinned override --> describe workflow shows the override
s.setAndCheckOverride(env, env.Tv(), s.makeAutoUpgradeOverride())
// 2. Set unpinned override --> describe workflow shows the override
s.setAndCheckOverride(env, env.Tv(), s.makeAutoUpgradeOverride())
}
func (s *DeploymentVersionSuite) TestUpdateWorkflowExecutionOptions_SetUnpinnedSetPinned() {
env := s.newTestEnv()
// Start a versioned poller which shall create a version; the version must be present before it can be set as an override.
s.startVersionWorkflow(s.Context(), env, env.Tv())
// start an unversioned workflow
s.startWorkflow(env, env.Tv(), nil)
// 1. Set unpinned override --> describe workflow shows the override
s.setAndCheckOverride(env, env.Tv(), s.makeAutoUpgradeOverride())
// 2. Set pinned override 1 --> describe workflow shows the override
s.setAndCheckOverride(env, env.Tv(), s.makePinnedOverride(env.Tv()))
}
func (s *DeploymentVersionSuite) TestUpdateWorkflowExecutionOptions_SetPinnedSetUnpinned() {
env := s.newTestEnv()
// Start a versioned poller which shall create a version; the version must be present before it can be set as an override.
s.startVersionWorkflow(s.Context(), env, env.Tv())
// start an unversioned workflow
s.startWorkflow(env, env.Tv(), nil)
// 1. Set pinned override 1 --> describe workflow shows the override
s.setAndCheckOverride(env, env.Tv(), s.makePinnedOverride(env.Tv()))
// 2. Set unpinned override --> describe workflow shows the override
s.setAndCheckOverride(env, env.Tv(), s.makeAutoUpgradeOverride())
}
func (s *DeploymentVersionSuite) TestUpdateWorkflowExecutionOptions_ReactivateVersionOnPinned() {
env := s.newTestEnv()
tv1 := env.Tv().WithBuildIDNumber(1) // Pinned target (INACTIVE)
tv2 := env.Tv().WithBuildIDNumber(2) // Current version
// v1 starts INACTIVE (never set as current). The reactivation signal handler in
// version_workflow.go treats DRAINED and INACTIVE identically — both flip to DRAINING.
s.startVersionWorkflow(s.Context(), env, tv1)
// v2 becomes the current version so the initial (non-pinned) workflow has a target.
s.startVersionWorkflow(s.Context(), env, tv2)
err := s.setCurrent(env, tv2, true)
s.NoError(err)
wf := func(version string) func(ctx workflow.Context) (string, error) {
return func(ctx workflow.Context) (string, error) {
workflow.GetSignalChannel(ctx, "complete").Receive(ctx, nil)
return "done from " + version, nil
}
}
// Register a worker for version 1 (INACTIVE) so it can accept workflows when
// UpdateWorkflowExecutionOptions is called to pin the workflow to version 1.
w1 := worker.New(env.SdkClient(), tv1.TaskQueue().String(), worker.Options{
DeploymentOptions: worker.DeploymentOptions{
Version: tv1.SDKDeploymentVersion(),
UseVersioning: true,
},
})
w1.RegisterWorkflowWithOptions(wf("v1"), workflow.RegisterOptions{
Name: "waitingWorkflow",
VersioningBehavior: workflow.VersioningBehaviorAutoUpgrade,
})
s.NoError(w1.Start())
defer w1.Stop()
// Register and start worker for version 2 on THE SAME task queue as version 1
w2 := worker.New(env.SdkClient(), tv1.TaskQueue().String(), worker.Options{
DeploymentOptions: worker.DeploymentOptions{
Version: tv2.SDKDeploymentVersion(),
UseVersioning: true,
},
})
w2.RegisterWorkflowWithOptions(wf("v2"), workflow.RegisterOptions{
Name: "waitingWorkflow",
VersioningBehavior: workflow.VersioningBehaviorAutoUpgrade,
})
s.NoError(w2.Start())
defer w2.Stop()
s.waitForPollers(env, tv1, tv2)
// Start the workflow. The workflow shall start on version 2, by default, since it is the current version.
wfTV := env.Tv()
run, err := env.SdkClient().ExecuteWorkflow(s.Context(), sdkclient.StartWorkflowOptions{
TaskQueue: tv1.TaskQueue().String(),
ID: wfTV.WorkflowID(),
}, "waitingWorkflow")
s.NoError(err)
// Pin the running workflow to version 1 (INACTIVE) using UpdateWorkflowExecutionOptions.
pinnedOverride := &workflowpb.VersioningOverride{
Override: &workflowpb.VersioningOverride_Pinned{
Pinned: &workflowpb.VersioningOverride_PinnedOverride{
Behavior: workflowpb.VersioningOverride_PINNED_OVERRIDE_BEHAVIOR_PINNED,
Version: tv1.ExternalDeploymentVersion(),
},
},
}
// Pin the workflow to version 1 (both versions are on the same task queue).
// Use Eventually to bypass version membership cache checks.
s.Eventually(func() bool {
_, err = env.FrontendClient().UpdateWorkflowExecutionOptions(s.Context(),
&workflowservice.UpdateWorkflowExecutionOptionsRequest{
Namespace: env.Namespace().String(),
WorkflowExecution: &commonpb.WorkflowExecution{
WorkflowId: wfTV.WorkflowID(),
RunId: run.GetRunID(),
},
WorkflowExecutionOptions: &workflowpb.WorkflowExecutionOptions{
VersioningOverride: pinnedOverride,
},
UpdateMask: &fieldmaskpb.FieldMask{Paths: []string{"versioning_override"}},
})
return err == nil
}, 10*time.Second, 500*time.Millisecond)
// Verify workflow has the pinned override
s.checkDescribeWorkflowAfterOverride(env,
&commonpb.WorkflowExecution{
WorkflowId: wfTV.WorkflowID(),
RunId: run.GetRunID(),
},
pinnedOverride)
// Wait for version 1 to show up as DRAINING
s.checkVersionDrainageAndVersionStatus(env, tv1,
&deploymentpb.VersionDrainageInfo{
Status: enumspb.VERSION_DRAINAGE_STATUS_DRAINING,
},
enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_DRAINING,
0)
// Verify via DescribeWorkerDeployment that the version status is updated
s.checkVersionStatusInDeployment(env, tv1, enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_DRAINING)
// Signal workflow to complete
s.NoError(env.SdkClient().SignalWorkflow(s.Context(),
wfTV.WorkflowID(), run.GetRunID(), "complete", nil))
// Wait for workflow to complete and verify it ran on version 1
var result string
s.NoError(run.Get(s.Context(), &result))
s.Equal("done from v1", result, "Workflow should have completed on version 1")
}
func (s *DeploymentVersionSuite) TestStartWorkflowExecution_ReactivateVersionOnPinned() {
env := s.newTestEnv()
tv1 := env.Tv()
// v1 starts INACTIVE (never set as current). The reactivation signal handler in
// version_workflow.go treats DRAINED and INACTIVE identically — both flip to DRAINING.
s.startVersionWorkflow(s.Context(), env, tv1)
wf := func(version string) func(ctx workflow.Context) (string, error) {
return func(ctx workflow.Context) (string, error) {
workflow.GetSignalChannel(ctx, "complete").Receive(ctx, nil)
return "done from " + version, nil
}
}
// Register a worker for version 1 (INACTIVE) so it can accept workflows when
// StartWorkflowExecution is called with a pinned override to version 1.
w1 := worker.New(env.SdkClient(), tv1.TaskQueue().String(), worker.Options{
DeploymentOptions: worker.DeploymentOptions{
Version: tv1.SDKDeploymentVersion(),
UseVersioning: true,
},
})
w1.RegisterWorkflowWithOptions(wf("v1"), workflow.RegisterOptions{
Name: "waitingWorkflow",
VersioningBehavior: workflow.VersioningBehaviorAutoUpgrade,
})
s.NoError(w1.Start())
defer w1.Stop()
// Start a new workflow with the pinned override pointing to version 1 (INACTIVE).
wfTV := env.Tv()
var run sdkclient.WorkflowRun
s.Eventually(func() bool {
var startErr error
run, startErr = env.SdkClient().ExecuteWorkflow(s.Context(), sdkclient.StartWorkflowOptions{
TaskQueue: tv1.TaskQueue().String(),
ID: wfTV.WorkflowID(),
VersioningOverride: &sdkclient.PinnedVersioningOverride{
Version: tv1.SDKDeploymentVersion(),
},
}, "waitingWorkflow")
return startErr == nil
}, 10*time.Second, 500*time.Millisecond)
// Verify workflow has the pinned override
s.checkDescribeWorkflowAfterOverride(env,
&commonpb.WorkflowExecution{
WorkflowId: wfTV.WorkflowID(),
RunId: run.GetRunID(),
},
&workflowpb.VersioningOverride{
Override: &workflowpb.VersioningOverride_Pinned{
Pinned: &workflowpb.VersioningOverride_PinnedOverride{
Behavior: workflowpb.VersioningOverride_PINNED_OVERRIDE_BEHAVIOR_PINNED,
Version: tv1.ExternalDeploymentVersion(),
},
},
})
// Wait for version 1 to show up as DRAINING (reactivated from INACTIVE)
s.checkVersionDrainageAndVersionStatus(env, tv1,
&deploymentpb.VersionDrainageInfo{
Status: enumspb.VERSION_DRAINAGE_STATUS_DRAINING,
},
enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_DRAINING,
0)
// Verify via DescribeWorkerDeployment that the version status is updated
s.checkVersionStatusInDeployment(env, tv1, enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_DRAINING)
// Signal workflow to complete
s.NoError(env.SdkClient().SignalWorkflow(s.Context(),
wfTV.WorkflowID(), run.GetRunID(), "complete", nil))
// Wait for workflow to complete and verify it ran on version 1
var result string
s.NoError(run.Get(s.Context(), &result))
s.Equal("done from v1", result, "Workflow should have completed on version 1")
}
func (s *DeploymentVersionSuite) TestStartWorkflowExecution_ReactivateVersionOnPinned_WithConflictPolicy() {
env := s.newTestEnv()
tv1 := env.Tv().WithBuildIDNumber(1)
tv2 := env.Tv().WithBuildIDNumber(2)
// v1 starts INACTIVE (never set as current). The reactivation signal handler in
// version_workflow.go treats DRAINED and INACTIVE identically — both flip to DRAINING.
s.startVersionWorkflow(s.Context(), env, tv1)
// v2 becomes the current version so the initial (non-pinned) workflow has a target.
s.startVersionWorkflow(s.Context(), env, tv2)
err := s.setCurrent(env, tv2, true)
s.NoError(err)
wf := func(version string) func(ctx workflow.Context) (string, error) {
return func(ctx workflow.Context) (string, error) {
workflow.GetSignalChannel(ctx, "complete").Receive(ctx, nil)
return "done from " + version, nil
}
}
// Register a worker for version 1 (INACTIVE) so it can accept workflows
w1 := worker.New(env.SdkClient(), tv1.TaskQueue().String(), worker.Options{
DeploymentOptions: worker.DeploymentOptions{
Version: tv1.SDKDeploymentVersion(),
UseVersioning: true,
},
})
w1.RegisterWorkflowWithOptions(wf("v1"), workflow.RegisterOptions{
Name: "waitingWorkflow",
VersioningBehavior: workflow.VersioningBehaviorAutoUpgrade,
})
s.NoError(w1.Start())
defer w1.Stop()
// Start a first workflow (no pinning, uses current version v2) to create a running execution
// with a specific workflow ID that we will terminate via conflict policy.
wfTV := env.Tv()
w2 := worker.New(env.SdkClient(), tv2.TaskQueue().String(), worker.Options{
DeploymentOptions: worker.DeploymentOptions{
Version: tv2.SDKDeploymentVersion(),
UseVersioning: true,
},
})
w2.RegisterWorkflowWithOptions(wf("v2"), workflow.RegisterOptions{
Name: "waitingWorkflow",
VersioningBehavior: workflow.VersioningBehaviorAutoUpgrade,
})
s.NoError(w2.Start())
defer w2.Stop()
s.Eventually(func() bool {
_, startErr := env.SdkClient().ExecuteWorkflow(s.Context(), sdkclient.StartWorkflowOptions{
TaskQueue: tv2.TaskQueue().String(),
ID: wfTV.WorkflowID(),
}, "waitingWorkflow")
return startErr == nil
}, 10*time.Second, 500*time.Millisecond)
// Now start a second workflow with the SAME workflow ID, pinned to v1 (INACTIVE),
// using TERMINATE_EXISTING conflict policy. This goes through the handleConflict method in api.go.
var run sdkclient.WorkflowRun
s.Eventually(func() bool {
var startErr error
run, startErr = env.SdkClient().ExecuteWorkflow(s.Context(), sdkclient.StartWorkflowOptions{
TaskQueue: tv1.TaskQueue().String(),
ID: wfTV.WorkflowID(),
VersioningOverride: &sdkclient.PinnedVersioningOverride{
Version: tv1.SDKDeploymentVersion(),
},
WorkflowIDConflictPolicy: enumspb.WORKFLOW_ID_CONFLICT_POLICY_TERMINATE_EXISTING,
}, "waitingWorkflow")
return startErr == nil
}, 10*time.Second, 500*time.Millisecond)
// Verify workflow has the pinned override
s.checkDescribeWorkflowAfterOverride(env,
&commonpb.WorkflowExecution{
WorkflowId: wfTV.WorkflowID(),
RunId: run.GetRunID(),
},
&workflowpb.VersioningOverride{
Override: &workflowpb.VersioningOverride_Pinned{
Pinned: &workflowpb.VersioningOverride_PinnedOverride{
Behavior: workflowpb.VersioningOverride_PINNED_OVERRIDE_BEHAVIOR_PINNED,
Version: tv1.ExternalDeploymentVersion(),
},
},
})
// Wait for version 1 to show up as DRAINING (reactivated from INACTIVE)
s.checkVersionDrainageAndVersionStatus(env, tv1,
&deploymentpb.VersionDrainageInfo{
Status: enumspb.VERSION_DRAINAGE_STATUS_DRAINING,
},
enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_DRAINING,
0)
// Verify via DescribeWorkerDeployment that the version status is updated
s.checkVersionStatusInDeployment(env, tv1, enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_DRAINING)
// Signal workflow to complete
s.NoError(env.SdkClient().SignalWorkflow(s.Context(),
wfTV.WorkflowID(), run.GetRunID(), "complete", nil))
// Wait for workflow to complete and verify it ran on version 1
var result string
s.NoError(run.Get(s.Context(), &result))
s.Equal("done from v1", result, "Workflow should have completed on version 1")
}
func (s *DeploymentVersionSuite) TestSignalWithStartWorkflowExecution_ReactivateVersionOnPinned() {
env := s.newTestEnv()
tv1 := env.Tv()
// v1 starts INACTIVE (never set as current). The reactivation signal handler in
// version_workflow.go treats DRAINED and INACTIVE identically — both flip to DRAINING.
s.startVersionWorkflow(s.Context(), env, tv1)
wf := func(version string) func(ctx workflow.Context) (string, error) {
return func(ctx workflow.Context) (string, error) {
workflow.GetSignalChannel(ctx, "complete").Receive(ctx, nil)
return "done from " + version, nil
}
}
// Register a worker for version 1 (INACTIVE) so it can accept workflows when
// SignalWithStartWorkflowExecution is called with a pinned override to version 1.
w1 := worker.New(env.SdkClient(), tv1.TaskQueue().String(), worker.Options{
DeploymentOptions: worker.DeploymentOptions{
Version: tv1.SDKDeploymentVersion(),
UseVersioning: true,
},
})
w1.RegisterWorkflowWithOptions(wf("v1"), workflow.RegisterOptions{
Name: "waitingWorkflow",
VersioningBehavior: workflow.VersioningBehaviorAutoUpgrade,
})
s.NoError(w1.Start())
defer w1.Stop()
// Use SignalWithStart with the pinned override pointing to version 1 (INACTIVE).
// This should START a new workflow (not signal an existing one) since no workflow exists yet.
wfTV := env.Tv()
var run sdkclient.WorkflowRun
s.Eventually(func() bool {
var startErr error
run, startErr = env.SdkClient().SignalWithStartWorkflow(s.Context(),
wfTV.WorkflowID(),
"start-signal", // signal name
nil, // signal arg
sdkclient.StartWorkflowOptions{
TaskQueue: tv1.TaskQueue().String(),
VersioningOverride: &sdkclient.PinnedVersioningOverride{
Version: tv1.SDKDeploymentVersion(),
},
},
"waitingWorkflow",
)
return startErr == nil
}, 10*time.Second, 500*time.Millisecond)
// Verify workflow has the pinned override
s.checkDescribeWorkflowAfterOverride(env,
&commonpb.WorkflowExecution{
WorkflowId: wfTV.WorkflowID(),
RunId: run.GetRunID(),
},
&workflowpb.VersioningOverride{
Override: &workflowpb.VersioningOverride_Pinned{
Pinned: &workflowpb.VersioningOverride_PinnedOverride{
Behavior: workflowpb.VersioningOverride_PINNED_OVERRIDE_BEHAVIOR_PINNED,
Version: tv1.ExternalDeploymentVersion(),
},
},
})
// Wait for version 1 to show up as DRAINING (reactivated from INACTIVE)
s.checkVersionDrainageAndVersionStatus(env, tv1,
&deploymentpb.VersionDrainageInfo{
Status: enumspb.VERSION_DRAINAGE_STATUS_DRAINING,
},
enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_DRAINING,
0)
// Verify via DescribeWorkerDeployment that the version status is updated
s.checkVersionStatusInDeployment(env, tv1, enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_DRAINING)
// Signal workflow to complete
s.NoError(env.SdkClient().SignalWorkflow(s.Context(),
wfTV.WorkflowID(), run.GetRunID(), "complete", nil))
// Wait for workflow to complete and verify it ran on version 1
var result string
s.NoError(run.Get(s.Context(), &result))
s.Equal("done from v1", result, "Workflow should have completed on version 1")
}
func (s *DeploymentVersionSuite) TestResetWorkflowExecution_ReactivateVersionOnPinned() {
env := s.newTestEnv()
tv1 := env.Tv().WithBuildIDNumber(1) // Pinned target (INACTIVE)
tv2 := env.Tv().WithBuildIDNumber(2) // Current version
// v1 starts INACTIVE (never set as current). The reactivation signal handler in
// version_workflow.go treats DRAINED and INACTIVE identically — both flip to DRAINING.
s.startVersionWorkflow(s.Context(), env, tv1)
// v2 becomes the current version so the initial (non-pinned) workflow has a target.
s.startVersionWorkflow(s.Context(), env, tv2)
err := s.setCurrent(env, tv2, true)
s.NoError(err)
// Workflow that waits for a signal, used for both versions.
// Returns a string indicating which version completed it.
wf := func(version string) func(ctx workflow.Context) (string, error) {
return func(ctx workflow.Context) (string, error) {
workflow.GetSignalChannel(ctx, "complete").Receive(ctx, nil)
return "done from " + version, nil
}
}
// Register a worker for version 1 (INACTIVE) so it can accept workflows when
// ResetWorkflowExecution is called with a pinned override to version 1.
w1 := worker.New(env.SdkClient(), tv1.TaskQueue().String(), worker.Options{
DeploymentOptions: worker.DeploymentOptions{
Version: tv1.SDKDeploymentVersion(),
UseVersioning: true,
},
})
w1.RegisterWorkflowWithOptions(wf("v1"), workflow.RegisterOptions{
Name: "waitingWorkflow",
VersioningBehavior: workflow.VersioningBehaviorAutoUpgrade,
})
s.NoError(w1.Start())
defer w1.Stop()
// Register and start worker for version 2 on THE SAME task queue as version 1
w2 := worker.New(env.SdkClient(), tv1.TaskQueue().String(), worker.Options{
DeploymentOptions: worker.DeploymentOptions{
Version: tv2.SDKDeploymentVersion(),
UseVersioning: true,
},
})
w2.RegisterWorkflowWithOptions(wf("v2"), workflow.RegisterOptions{
Name: "waitingWorkflow",
VersioningBehavior: workflow.VersioningBehaviorAutoUpgrade,
})
s.NoError(w2.Start())
defer w2.Stop()
s.waitForPollers(env, tv1, tv2)
// Start a workflow on the current version (v2)
wfTV := env.Tv()
run, err := env.SdkClient().ExecuteWorkflow(s.Context(), sdkclient.StartWorkflowOptions{
TaskQueue: tv1.TaskQueue().String(),
ID: wfTV.WorkflowID(),
}, "waitingWorkflow")
s.NoError(err)
// Wait for the workflow to start and complete its first workflow task (creates a reset point)
s.Eventually(func() bool {
hist := env.SdkClient().GetWorkflowHistory(s.Context(), wfTV.WorkflowID(), run.GetRunID(), false, enumspb.HISTORY_EVENT_FILTER_TYPE_ALL_EVENT)
for hist.HasNext() {
event, err := hist.Next()
if err != nil {
return false
}
if event.EventType == enumspb.EVENT_TYPE_WORKFLOW_TASK_COMPLETED {
return true
}
}
return false
}, 10*time.Second, 200*time.Millisecond, "Workflow should have completed its first workflow task")
// Find the first workflow task complete event ID for the reset point
var resetEventID int64
hist := env.SdkClient().GetWorkflowHistory(s.Context(), wfTV.WorkflowID(), run.GetRunID(), false, enumspb.HISTORY_EVENT_FILTER_TYPE_ALL_EVENT)
for hist.HasNext() {
event, err := hist.Next()
s.NoError(err)
if event.EventType == enumspb.EVENT_TYPE_WORKFLOW_TASK_COMPLETED {
resetEventID = event.EventId
break
}
}
s.Positive(resetEventID, "Should have found a workflow task complete event")
// Reset the workflow with PostResetOperations containing a versioning override pinned to v1 (which is currently INACTIVE)
var resetResp *workflowservice.ResetWorkflowExecutionResponse
s.Eventually(func() bool {
var resetErr error
resetResp, resetErr = env.FrontendClient().ResetWorkflowExecution(s.Context(), &workflowservice.ResetWorkflowExecutionRequest{
Namespace: env.Namespace().String(),
WorkflowExecution: &commonpb.WorkflowExecution{
WorkflowId: wfTV.WorkflowID(),
RunId: run.GetRunID(),
},
Reason: "testing-reset-reactivation",
RequestId: uuid.NewString(),
WorkflowTaskFinishEventId: resetEventID,
PostResetOperations: []*workflowpb.PostResetOperation{
{
Variant: &workflowpb.PostResetOperation_UpdateWorkflowOptions_{
UpdateWorkflowOptions: &workflowpb.PostResetOperation_UpdateWorkflowOptions{
WorkflowExecutionOptions: &workflowpb.WorkflowExecutionOptions{
VersioningOverride: &workflowpb.VersioningOverride{
Override: &workflowpb.VersioningOverride_Pinned{
Pinned: &workflowpb.VersioningOverride_PinnedOverride{
Behavior: workflowpb.VersioningOverride_PINNED_OVERRIDE_BEHAVIOR_PINNED,
Version: tv1.ExternalDeploymentVersion(),
},
},
},
},
UpdateMask: &fieldmaskpb.FieldMask{
Paths: []string{"versioning_override"},
},
},
},
},
},
})
return resetErr == nil
}, 10*time.Second, 500*time.Millisecond)
newRunID := resetResp.RunId
// Verify the reset workflow has the pinned override
s.checkDescribeWorkflowAfterOverride(env,
&commonpb.WorkflowExecution{
WorkflowId: wfTV.WorkflowID(),
RunId: newRunID,
},
&workflowpb.VersioningOverride{
Override: &workflowpb.VersioningOverride_Pinned{
Pinned: &workflowpb.VersioningOverride_PinnedOverride{
Behavior: workflowpb.VersioningOverride_PINNED_OVERRIDE_BEHAVIOR_PINNED,
Version: tv1.ExternalDeploymentVersion(),
},
},
})
// Wait for version 1 to show up as DRAINING (reactivated from INACTIVE)
s.checkVersionDrainageAndVersionStatus(env, tv1,
&deploymentpb.VersionDrainageInfo{
Status: enumspb.VERSION_DRAINAGE_STATUS_DRAINING,
},
enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_DRAINING,
0)
// Verify via DescribeWorkerDeployment that the version status is DRAINING
s.checkVersionStatusInDeployment(env, tv1, enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_DRAINING)
// Signal the reset workflow to complete
s.NoError(env.SdkClient().SignalWorkflow(s.Context(),
wfTV.WorkflowID(), newRunID, "complete", nil))
// Wait for the reset workflow to complete and verify it ran on version 1
resetRun := env.SdkClient().GetWorkflow(s.Context(), wfTV.WorkflowID(), newRunID)
var result string
s.NoError(resetRun.Get(s.Context(), &result))
s.Equal("done from v1", result, "Reset workflow should have completed on version 1")
}
// The following tests test the VersioningOverride functionality when passed via the BatchUpdateWorkflowExecutionOptions API.
func (s *DeploymentVersionSuite) TestBatchUpdateWorkflowExecutionOptions_SetPinned_VersionDoesNotExist() {
s.runBatchUpdateWorkflowExecutionOptionsTest(false)
}
func (s *DeploymentVersionSuite) TestBatchUpdateWorkflowExecutionOptions_SetPinnedThenUnset() {
s.runBatchUpdateWorkflowExecutionOptionsTest(true)
}
func (s *DeploymentVersionSuite) runBatchUpdateWorkflowExecutionOptionsTest(createVersionFirst bool) {
env := s.newTestEnv()
// start some unversioned workflows
workflowType := "UpdateOptionsBatchTestFunc"
workflows := make([]*commonpb.WorkflowExecution, 0)
for range 5 {
run, err := env.SdkClient().ExecuteWorkflow(s.Context(), sdkclient.StartWorkflowOptions{TaskQueue: env.Tv().TaskQueue().Name}, workflowType)
s.NoError(err)
workflows = append(workflows, &commonpb.WorkflowExecution{
WorkflowId: run.GetID(),
RunId: run.GetRunID(),
})
}
pinnedOverride := s.makePinnedOverride(env.Tv())
batchJobID := uuid.NewString()
if createVersionFirst {
// Start a versioned poller which shall create a version
s.startVersionWorkflow(s.Context(), env, env.Tv())
}
// start batch update-options operation
_, err := env.SdkClient().WorkflowService().StartBatchOperation(context.Background(), &workflowservice.StartBatchOperationRequest{
Namespace: env.Namespace().String(),
Operation: &workflowservice.StartBatchOperationRequest_UpdateWorkflowOptionsOperation{
UpdateWorkflowOptionsOperation: &batchpb.BatchOperationUpdateWorkflowExecutionOptions{
Identity: env.Tv().ClientIdentity(),
WorkflowExecutionOptions: &workflowpb.WorkflowExecutionOptions{VersioningOverride: pinnedOverride},
UpdateMask: &fieldmaskpb.FieldMask{Paths: []string{"versioning_override"}},
},
},
Executions: workflows,
JobId: batchJobID,
Reason: "test",
})
s.NoError(err)
if !createVersionFirst {
s.checkBatchOperationFails(env, batchJobID, len(workflows))
for _, wf := range workflows {
s.checkDescribeWorkflowAfterOverride(env, wf, nil)
}
return
}
// wait til batch completes successfully
s.checkListAndWaitForBatchCompletion(env, batchJobID)
// check all the workflows
for _, wf := range workflows {
s.checkDescribeWorkflowAfterOverride(env, wf, pinnedOverride)
s.checkWorkflowUpdateOptionsEventIdentity(s.Context(), env, wf, env.Tv().ClientIdentity())
}
// unset with empty update opts with mutation mask
batchJobID = uuid.NewString()
err = s.startBatchJobWithinConcurrentJobLimit(env, &workflowservice.StartBatchOperationRequest{
Namespace: env.Namespace().String(),
JobId: batchJobID,
Reason: "test",
Executions: workflows,
Operation: &workflowservice.StartBatchOperationRequest_UpdateWorkflowOptionsOperation{
UpdateWorkflowOptionsOperation: &batchpb.BatchOperationUpdateWorkflowExecutionOptions{
Identity: env.Tv().ClientIdentity(),
WorkflowExecutionOptions: &workflowpb.WorkflowExecutionOptions{},
UpdateMask: &fieldmaskpb.FieldMask{Paths: []string{"versioning_override"}},
},
},
})
s.NoError(err)
// wait til batch completes
s.checkListAndWaitForBatchCompletion(env, batchJobID)
// check all the workflows
for _, wf := range workflows {
s.checkDescribeWorkflowAfterOverride(env, wf, nil)
s.checkWorkflowUpdateOptionsEventIdentity(s.Context(), env, wf, env.Tv().ClientIdentity())
}
}
func (s *DeploymentVersionSuite) startBatchJobWithinConcurrentJobLimit(env *testcore.TestEnv, req *workflowservice.StartBatchOperationRequest) error {
var err error
s.Eventually(func() bool {
_, err = env.FrontendClient().StartBatchOperation(s.Context(), req)
if err == nil {
return true
} else if strings.Contains(err.Error(), "Max concurrent batch operations is reached") {
return false // retry
}
return true
}, 5*time.Second, 500*time.Millisecond)
return err
}
func (s *DeploymentVersionSuite) checkListAndWaitForBatchCompletion(env *testcore.TestEnv, jobID string) {
s.EventuallyWithT(func(t *assert.CollectT) {
a := assert.New(t)
listResp, err := env.FrontendClient().ListBatchOperations(s.Context(), &workflowservice.ListBatchOperationsRequest{
Namespace: env.Namespace().String(),
})
a.NoError(err)
a.NotEmpty(listResp.GetOperationInfo())
if len(listResp.GetOperationInfo()) > 0 {
a.Equal(jobID, listResp.GetOperationInfo()[0].GetJobId())
}
}, 10*time.Second, 50*time.Millisecond)
s.EventuallyWithT(func(t *assert.CollectT) {
a := assert.New(t)
descResp, err := env.FrontendClient().DescribeBatchOperation(s.Context(), &workflowservice.DescribeBatchOperationRequest{
Namespace: env.Namespace().String(),
JobId: jobID,
})
a.NoError(err)
a.NotEqual(enumspb.BATCH_OPERATION_STATE_FAILED, descResp.GetState(), "batch operation failed. description: %+v", descResp)
a.Equal(enumspb.BATCH_OPERATION_STATE_COMPLETED, descResp.GetState())
}, 10*time.Second, 50*time.Millisecond)
}
func (s *DeploymentVersionSuite) checkBatchOperationFails(env *testcore.TestEnv, jobID string, numWorkflows int) {
s.EventuallyWithT(func(t *assert.CollectT) {
a := assert.New(t)
descResp, err := env.FrontendClient().DescribeBatchOperation(s.Context(), &workflowservice.DescribeBatchOperationRequest{
Namespace: env.Namespace().String(),
JobId: jobID,
})
a.NoError(err)
// All workflows should have failed validation
a.Equal(int64(numWorkflows), descResp.GetFailureOperationCount(), "expected all operations to fail")
}, 30*time.Second, 500*time.Millisecond)
}
func (s *DeploymentVersionSuite) makePinnedOverride(tv *testvars.TestVars) *workflowpb.VersioningOverride {
if useV32 {
return &workflowpb.VersioningOverride{Override: &workflowpb.VersioningOverride_Pinned{
Pinned: &workflowpb.VersioningOverride_PinnedOverride{
Behavior: workflowpb.VersioningOverride_PINNED_OVERRIDE_BEHAVIOR_PINNED,
Version: tv.ExternalDeploymentVersion(),
},
}}
}
return &workflowpb.VersioningOverride{
PinnedVersion: tv.DeploymentVersionString(), //nolint:staticcheck // SA1019: worker versioning v0.31
Behavior: enumspb.VERSIONING_BEHAVIOR_PINNED, //nolint:staticcheck // SA1019: worker versioning v0.31
}
}
func (s *DeploymentVersionSuite) makeAutoUpgradeOverride() *workflowpb.VersioningOverride {
if useV32 {
return &workflowpb.VersioningOverride{Override: &workflowpb.VersioningOverride_AutoUpgrade{AutoUpgrade: true}}
}
return &workflowpb.VersioningOverride{Behavior: enumspb.VERSIONING_BEHAVIOR_AUTO_UPGRADE} //nolint:staticcheck // SA1019: worker versioning v0.31
}
func (s *DeploymentVersionSuite) TestStartWorkflowExecution_WithPinnedOverride_CacheMissAndHits() {
env := s.newTestEnv(
// TODO: remove WithWorkerService once legacy suite-scoped cluster behavior is removed.
testcore.WithWorkerService("worker-deployment version membership cache test"),
testcore.WithDynamicConfig(dynamicconfig.VersionMembershipCacheTTL, 5*time.Second),
)
override := s.makePinnedOverride(env.Tv())
request := &workflowservice.StartWorkflowExecutionRequest{
RequestId: env.Tv().Any().String(),
Namespace: env.Namespace().String(),
WorkflowId: env.Tv().WorkflowID(),
WorkflowType: env.Tv().WorkflowType(),
TaskQueue: env.Tv().TaskQueue(),
Identity: env.Tv().WorkerIdentity(),
VersioningOverride: override,
}
// First call should fail since the version to override is not present in the task queue.
_, err0 := env.FrontendClient().StartWorkflowExecution(s.Context(), request)
s.Error(err0)
// Start a versioned poller which shall create the version; the version must be present before it can be set as an override.
s.startVersionWorkflow(s.Context(), env, env.Tv())
// Wait for the cache TTL to expire; On expiry of the cache TTL, it would result in a fresh RPC which would verify the version presence,
// eventually leading to the StartWorkflowExecution call succeeding.
var resp *workflowservice.StartWorkflowExecutionResponse
s.Eventually(func() bool {
var err error
resp, err = env.FrontendClient().StartWorkflowExecution(s.Context(), request)
return err == nil
}, 10*time.Second, 500*time.Millisecond)
// The StartWorkflowExecution should now succeed with no error. Verify that the workflow shows the override.
s.checkDescribeWorkflowAfterOverride(env, &commonpb.WorkflowExecution{
WorkflowId: env.Tv().WorkflowID(),
RunId: resp.GetRunId(),
}, override)
}
func (s *DeploymentVersionSuite) TestStartWorkflowExecution_WithUnpinnedOverride() {
env := s.newTestEnv()
override := s.makeAutoUpgradeOverride()
wf := &commonpb.WorkflowExecution{
WorkflowId: env.Tv().WorkflowID(),
RunId: s.startWorkflow(env, env.Tv(), override),
}
s.checkDescribeWorkflowAfterOverride(env, wf, override)
}
func (s *DeploymentVersionSuite) TestSignalWithStartWorkflowExecution_WithPinnedOverride_CacheMissAndHits() {
env := s.newTestEnv(
// TODO: remove WithWorkerService once legacy suite-scoped cluster behavior is removed.
testcore.WithWorkerService("worker-deployment version membership cache test"),
testcore.WithDynamicConfig(dynamicconfig.VersionMembershipCacheTTL, 5*time.Second),
)
override := s.makePinnedOverride(env.Tv())
request := &workflowservice.SignalWithStartWorkflowExecutionRequest{
Namespace: env.Namespace().String(),
WorkflowId: env.Tv().WorkflowID(),
WorkflowType: env.Tv().WorkflowType(),
TaskQueue: env.Tv().TaskQueue(),
Identity: env.Tv().ClientIdentity(),
RequestId: env.Tv().RequestID(),
SignalName: "test-signal",
SignalInput: nil,
VersioningOverride: override,
}
// Since the version to override is not present in the task queue, the call should fail.
_, err := env.FrontendClient().SignalWithStartWorkflowExecution(s.Context(), request)
s.Error(err)
// Start a versioned poller which shall create the version; the version must be present before it can be set as an override.
s.startVersionWorkflow(s.Context(), env, env.Tv())
// Wait for the cache TTL to expire; On expiry of the cache TTL, it would result in a fresh RPC which would verify the version presence,
// eventually leading to the SignalWithStartWorkflowExecution call succeeding.
var resp *workflowservice.SignalWithStartWorkflowExecutionResponse
s.Eventually(func() bool {
var err error
resp, err = env.FrontendClient().SignalWithStartWorkflowExecution(s.Context(), request)
return err == nil && resp.GetStarted()
}, 10*time.Second, 500*time.Millisecond)
wf := &commonpb.WorkflowExecution{
WorkflowId: env.Tv().WorkflowID(),
RunId: resp.GetRunId(),
}
s.checkDescribeWorkflowAfterOverride(env, wf, override)
}
func (s *DeploymentVersionSuite) TestSignalWithStartWorkflowExecution_WithUnpinnedOverride() {
env := s.newTestEnv()
override := s.makeAutoUpgradeOverride()
resp, err := env.FrontendClient().SignalWithStartWorkflowExecution(s.Context(), &workflowservice.SignalWithStartWorkflowExecutionRequest{
Namespace: env.Namespace().String(),
WorkflowId: env.Tv().WorkflowID(),
WorkflowType: env.Tv().WorkflowType(),
TaskQueue: env.Tv().TaskQueue(),
Identity: env.Tv().ClientIdentity(),
RequestId: env.Tv().RequestID(),
SignalName: "test-signal",
SignalInput: nil,
VersioningOverride: override,
})
s.NoError(err)
s.True(resp.GetStarted())
wf := &commonpb.WorkflowExecution{
WorkflowId: env.Tv().WorkflowID(),
RunId: resp.GetRunId(),
}
s.checkDescribeWorkflowAfterOverride(env, wf, override)
}
func (s *DeploymentVersionSuite) TestDeleteVersion_ThenRecreateByPolling() {
s.skipBeforeVersion(workerdeployment.VersionDataRevisionNumber)
env := s.newTestEnv(testcore.WithDynamicConfig(dynamicconfig.PollerHistoryTTL, 500*time.Millisecond))
tv1 := env.Tv().WithBuildIDNumber(1)
s.startVersionWorkflow(s.Context(), env, tv1)
vd := s.getTaskQueueVersionData(env, tv1, enumspb.TASK_QUEUE_TYPE_WORKFLOW, tv1.ExternalDeploymentVersion())
s.Equal(int64(0), vd.GetRevisionNumber())
s.False(vd.GetDeleted())
// Wait for pollers to go away
s.EventuallyWithT(func(t *assert.CollectT) {
resp, err := env.FrontendClient().DescribeTaskQueue(s.Context(), &workflowservice.DescribeTaskQueueRequest{
Namespace: env.Namespace().String(),
TaskQueue: tv1.TaskQueue(),
TaskQueueType: enumspb.TASK_QUEUE_TYPE_WORKFLOW,
})
require.NoError(t, err)
require.Empty(t, resp.Pollers)
}, 5*time.Second, time.Second)
// Delete the version
s.tryDeleteVersion(env, tv1, "", false)
// Verify the version is gone from the task queue
s.EventuallyWithT(func(t *assert.CollectT) {
vd = s.getTaskQueueVersionData(env, tv1, enumspb.TASK_QUEUE_TYPE_WORKFLOW, tv1.ExternalDeploymentVersion())
require.New(t).Nil(vd)
}, time.Second*5, time.Millisecond*200)
// Verify the version is gone from the deployment
s.EventuallyWithT(func(t *assert.CollectT) {
a := require.New(t)
resp, err := env.FrontendClient().DescribeWorkerDeployment(s.Context(), &workflowservice.DescribeWorkerDeploymentRequest{
Namespace: env.Namespace().String(),
DeploymentName: tv1.DeploymentSeries(),
})
a.NoError(err)
for _, vs := range resp.GetWorkerDeploymentInfo().GetVersionSummaries() {
//nolint:staticcheck // SA1019 deprecated Version will clean up later
a.NotEqual(tv1.DeploymentVersionString(), vs.Version)
}
}, time.Second*5, time.Millisecond*200)
// Poll again to recreate the version
s.startVersionWorkflow(s.Context(), env, tv1)
// Verify the version is back (undeleted) in the deployment
s.EventuallyWithT(func(t *assert.CollectT) {
a := require.New(t)
resp, err := env.FrontendClient().DescribeWorkerDeployment(s.Context(), &workflowservice.DescribeWorkerDeploymentRequest{
Namespace: env.Namespace().String(),
DeploymentName: tv1.DeploymentSeries(),
})
a.NoError(err)
found := false
for _, vs := range resp.GetWorkerDeploymentInfo().GetVersionSummaries() {
//nolint:staticcheck // SA1019 deprecated Version will clean up later
if vs.Version == tv1.DeploymentVersionString() {
found = true
}
}
a.True(found, "version should be recreated after polling")
}, time.Second*5, time.Millisecond*200)
// Ensure the version data revived properly in the task queue
vd = s.getTaskQueueVersionData(env, tv1, enumspb.TASK_QUEUE_TYPE_WORKFLOW, tv1.ExternalDeploymentVersion())
s.Equal(int64(0), vd.GetRevisionNumber())
s.False(vd.GetDeleted())
}
// getTaskQueueDeploymentData gets the deployment data for a given TQ type. The data is always
// returned from the WF type root partition, so no need to wait for propagation before calling this
// function.
func (s *DeploymentVersionSuite) getTaskQueueDeploymentData(
env *testcore.TestEnv,
tv *testvars.TestVars,
tqType enumspb.TaskQueueType,
) *persistencespb.DeploymentData {
ctx, cancel := context.WithTimeout(s.Context(), time.Second*5)
defer cancel()
resp, err := env.GetTestCluster().MatchingClient().GetTaskQueueUserData(
ctx,
&matchingservice.GetTaskQueueUserDataRequest{
NamespaceId: env.NamespaceID().String(),
TaskQueue: tv.TaskQueue().GetName(),
TaskQueueType: tqTypeWf,
})
s.NoError(err)
return resp.GetUserData().GetData().GetPerType()[int32(tqType)].GetDeploymentData()
}
func (s *DeploymentVersionSuite) getTaskQueueVersionData(
env *testcore.TestEnv,
tv *testvars.TestVars,
tqType enumspb.TaskQueueType,
version *deploymentpb.WorkerDeploymentVersion,
) *deploymentspb.WorkerDeploymentVersionData {
data := s.getTaskQueueDeploymentData(env, tv, tqType)
return data.GetDeploymentsData()[version.GetDeploymentName()].GetVersions()[version.GetBuildId()]
}
func (s *DeploymentVersionSuite) TestCreateWorkerDeploymentVersion_Success() {
env := s.newTestEnv()
buildID := env.Tv().BuildID()
requestID := env.Tv().Any().String()
identity := env.Tv().Any().String()
// First create the deployment
_, err := env.FrontendClient().CreateWorkerDeployment(s.Context(), &workflowservice.CreateWorkerDeploymentRequest{
Namespace: env.Namespace().String(),
DeploymentName: env.Tv().DeploymentSeries(),
RequestId: env.Tv().Any().String(),
})
s.NoError(err)
computeConfig := &computepb.ComputeConfig{
ScalingGroups: map[string]*computepb.ComputeConfigScalingGroup{
"sg1": {
Provider: computeprovider.TestInvokeComputeProviderValidComputeProvider(),
Scaler: testInvokeScaler(),
},
},
}
// Create a version in the deployment
resp, err := env.FrontendClient().CreateWorkerDeploymentVersion(s.Context(), &workflowservice.CreateWorkerDeploymentVersionRequest{
Namespace: env.Namespace().String(),
DeploymentVersion: &deploymentpb.WorkerDeploymentVersion{
DeploymentName: env.Tv().DeploymentSeries(),
BuildId: buildID,
},
Identity: identity,
RequestId: requestID,
ComputeConfig: computeConfig,
})
s.NoError(err)
s.NotNil(resp)
// Verify the version exists via DescribeWorkerDeploymentVersion
s.EventuallyWithT(func(t *assert.CollectT) {
a := require.New(t)
descResp, err := env.FrontendClient().DescribeWorkerDeploymentVersion(s.Context(), &workflowservice.DescribeWorkerDeploymentVersionRequest{
Namespace: env.Namespace().String(),
Version: env.Tv().DeploymentVersionString(),
})
a.NoError(err)
a.NotNil(descResp.GetWorkerDeploymentVersionInfo())
a.Equal(env.Tv().DeploymentVersionStringV32(), worker_versioning.ExternalWorkerDeploymentVersionToString(descResp.GetWorkerDeploymentVersionInfo().GetDeploymentVersion()))
a.NotNil(descResp.GetWorkerDeploymentVersionInfo().GetCreateTime())
a.True(proto.Equal(computeConfig, descResp.GetWorkerDeploymentVersionInfo().GetComputeConfig()))
a.Equal(enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_CREATED, descResp.GetWorkerDeploymentVersionInfo().GetStatus())
a.Equal(identity, descResp.GetWorkerDeploymentVersionInfo().GetLastModifierIdentity())
}, 10*time.Second, 500*time.Millisecond)
// Verify the version shows up in deployment's version summaries with CREATED status and correct compute config summary.
s.EventuallyWithT(func(t *assert.CollectT) {
a := require.New(t)
descDeployResp, err := env.FrontendClient().DescribeWorkerDeployment(s.Context(), &workflowservice.DescribeWorkerDeploymentRequest{
Namespace: env.Namespace().String(),
DeploymentName: env.Tv().DeploymentSeries(),
})
a.NoError(err)
a.Len(descDeployResp.GetWorkerDeploymentInfo().GetVersionSummaries(), 1)
versionSummary := descDeployResp.GetWorkerDeploymentInfo().GetVersionSummaries()[0]
a.Equal(env.Tv().DeploymentVersionStringV32(), worker_versioning.ExternalWorkerDeploymentVersionToString(versionSummary.GetDeploymentVersion()))
a.Equal(enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_CREATED, versionSummary.GetStatus())
a.True(proto.Equal(&computepb.ComputeConfigSummary{
ScalingGroups: map[string]*computepb.ComputeConfigScalingGroupSummary{
"sg1": {
ProviderType: computeConfig.GetScalingGroups()["sg1"].GetProvider().GetType(),
},
},
}, versionSummary.GetComputeConfig()))
}, 10*time.Second, 500*time.Millisecond)
// Verify the compute config summary is reflected in ListWorkerDeployments latest version summary.
s.EventuallyWithT(func(t *assert.CollectT) {
a := require.New(t)
listResp, err := env.FrontendClient().ListWorkerDeployments(s.Context(), &workflowservice.ListWorkerDeploymentsRequest{
Namespace: env.Namespace().String(),
})
a.NoError(err)
var found *workflowservice.ListWorkerDeploymentsResponse_WorkerDeploymentSummary
for _, d := range listResp.GetWorkerDeployments() {
if d.GetName() == env.Tv().DeploymentSeries() {
found = d
break
}
}
a.NotNil(found, "deployment %s not found in ListWorkerDeployments", env.Tv().DeploymentSeries())
a.True(proto.Equal(&computepb.ComputeConfigSummary{
ScalingGroups: map[string]*computepb.ComputeConfigScalingGroupSummary{
"sg1": {
ProviderType: computeConfig.GetScalingGroups()["sg1"].GetProvider().GetType(),
},
},
}, found.GetLatestVersionSummary().GetComputeConfig()))
}, 10*time.Second, 500*time.Millisecond)
}
func (s *DeploymentVersionSuite) TestCreateWorkerDeploymentVersion_ThenPoll_TaskQueueInVersionInfo() {
env := s.newTestEnv()
buildID := env.Tv().BuildID()
// Create the deployment
_, err := env.FrontendClient().CreateWorkerDeployment(s.Context(), &workflowservice.CreateWorkerDeploymentRequest{
Namespace: env.Namespace().String(),
DeploymentName: env.Tv().DeploymentSeries(),
RequestId: env.Tv().Any().String(),
})
s.NoError(err)
// Create the version explicitly
_, err = env.FrontendClient().CreateWorkerDeploymentVersion(s.Context(), &workflowservice.CreateWorkerDeploymentVersionRequest{
Namespace: env.Namespace().String(),
DeploymentVersion: &deploymentpb.WorkerDeploymentVersion{
DeploymentName: env.Tv().DeploymentSeries(),
BuildId: buildID,
},
RequestId: env.Tv().Any().String(),
})
s.NoError(err)
// Verify the version starts with CREATED status
s.EventuallyWithT(func(t *assert.CollectT) {
a := require.New(t)
descResp, err := env.FrontendClient().DescribeWorkerDeploymentVersion(s.Context(), &workflowservice.DescribeWorkerDeploymentVersionRequest{
Namespace: env.Namespace().String(),
Version: env.Tv().DeploymentVersionString(),
})
a.NoError(err)
a.Equal(enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_CREATED, descResp.GetWorkerDeploymentVersionInfo().GetStatus())
}, 10*time.Second, 500*time.Millisecond)
// Poll from the version to register a task queue
go s.pollFromDeployment(s.Context(), env, env.Tv())
// Verify the task queue shows up and status transitions to INACTIVE
s.EventuallyWithT(func(t *assert.CollectT) {
a := require.New(t)
descResp, err := env.FrontendClient().DescribeWorkerDeploymentVersion(s.Context(), &workflowservice.DescribeWorkerDeploymentVersionRequest{
Namespace: env.Namespace().String(),
Version: env.Tv().DeploymentVersionString(),
})
a.NoError(err)
tqInfos := descResp.GetWorkerDeploymentVersionInfo().GetTaskQueueInfos()
a.GreaterOrEqual(len(tqInfos), 1)
found := false
for _, tqInfo := range tqInfos {
if tqInfo.GetName() == env.Tv().TaskQueue().GetName() {
found = true
break
}
}
a.True(found, "expected task queue %q in version info, got %v", env.Tv().TaskQueue().GetName(), tqInfos)
a.Equal(enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_INACTIVE, descResp.GetWorkerDeploymentVersionInfo().GetStatus())
}, 30*time.Second, 500*time.Millisecond)
// Verify the version shows up in deployment's version summaries with INACTIVE status
s.EventuallyWithT(func(t *assert.CollectT) {
a := require.New(t)
descDeployResp, err := env.FrontendClient().DescribeWorkerDeployment(s.Context(), &workflowservice.DescribeWorkerDeploymentRequest{
Namespace: env.Namespace().String(),
DeploymentName: env.Tv().DeploymentSeries(),
})
a.NoError(err)
a.Len(descDeployResp.GetWorkerDeploymentInfo().GetVersionSummaries(), 1)
a.Equal(env.Tv().DeploymentVersionStringV32(), worker_versioning.ExternalWorkerDeploymentVersionToString(descDeployResp.GetWorkerDeploymentInfo().GetVersionSummaries()[0].GetDeploymentVersion()))
a.Equal(enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_INACTIVE, descDeployResp.GetWorkerDeploymentInfo().GetVersionSummaries()[0].GetStatus())
}, 10*time.Second, 500*time.Millisecond)
}
func (s *DeploymentVersionSuite) TestCreateWorkerDeploymentVersion_Idempotent() {
env := s.newTestEnv()
buildID := env.Tv().BuildID()
requestID := env.Tv().Any().String()
// First create the deployment
_, err := env.FrontendClient().CreateWorkerDeployment(s.Context(), &workflowservice.CreateWorkerDeploymentRequest{
Namespace: env.Namespace().String(),
DeploymentName: env.Tv().DeploymentSeries(),
RequestId: env.Tv().Any().String(),
})
s.NoError(err)
// Create a version
_, err = env.FrontendClient().CreateWorkerDeploymentVersion(s.Context(), &workflowservice.CreateWorkerDeploymentVersionRequest{
Namespace: env.Namespace().String(),
DeploymentVersion: &deploymentpb.WorkerDeploymentVersion{
DeploymentName: env.Tv().DeploymentSeries(),
BuildId: buildID,
},
RequestId: requestID,
})
s.NoError(err)
// Create the same version again with same request ID - should be idempotent
_, err = env.FrontendClient().CreateWorkerDeploymentVersion(s.Context(), &workflowservice.CreateWorkerDeploymentVersionRequest{
Namespace: env.Namespace().String(),
DeploymentVersion: &deploymentpb.WorkerDeploymentVersion{
DeploymentName: env.Tv().DeploymentSeries(),
BuildId: buildID,
},
RequestId: requestID,
})
s.NoError(err)
}
func (s *DeploymentVersionSuite) TestCreateWorkerDeploymentVersion_AlreadyExists_DifferentRequestID() {
env := s.newTestEnv()
buildID := env.Tv().BuildID()
requestID1 := env.Tv().Any().String()
requestID2 := env.Tv().Any().String()
// First create the deployment
_, err := env.FrontendClient().CreateWorkerDeployment(s.Context(), &workflowservice.CreateWorkerDeploymentRequest{
Namespace: env.Namespace().String(),
DeploymentName: env.Tv().DeploymentSeries(),
RequestId: env.Tv().Any().String(),
})
s.NoError(err)
// Create a version
_, err = env.FrontendClient().CreateWorkerDeploymentVersion(s.Context(), &workflowservice.CreateWorkerDeploymentVersionRequest{
Namespace: env.Namespace().String(),
DeploymentVersion: &deploymentpb.WorkerDeploymentVersion{
DeploymentName: env.Tv().DeploymentSeries(),
BuildId: buildID,
},
RequestId: requestID1,
})
s.NoError(err)
// Try to create the same version with different request ID - should fail
_, err = env.FrontendClient().CreateWorkerDeploymentVersion(s.Context(), &workflowservice.CreateWorkerDeploymentVersionRequest{
Namespace: env.Namespace().String(),
DeploymentVersion: &deploymentpb.WorkerDeploymentVersion{
DeploymentName: env.Tv().DeploymentSeries(),
BuildId: buildID,
},
RequestId: requestID2,
})
s.Error(err)
var alreadyExists *serviceerror.AlreadyExists
s.ErrorAs(err, &alreadyExists)
}
func (s *DeploymentVersionSuite) TestCreateWorkerDeploymentVersion_DeploymentNotFound() {
env := s.newTestEnv()
// Try to create a version for a deployment that doesn't exist
_, err := env.FrontendClient().CreateWorkerDeploymentVersion(s.Context(), &workflowservice.CreateWorkerDeploymentVersionRequest{
Namespace: env.Namespace().String(),
DeploymentVersion: &deploymentpb.WorkerDeploymentVersion{
DeploymentName: env.Tv().DeploymentSeries(),
BuildId: env.Tv().BuildID(),
},
RequestId: env.Tv().Any().String(),
})
s.Error(err)
var notFound *serviceerror.NotFound
s.ErrorAs(err, &notFound)
}
func (s *DeploymentVersionSuite) TestCreateWorkerDeploymentVersion_InvalidArgs() {
testCases := []struct {
name string
deploymentName string
buildID string
expectedError string
}{
{
name: "empty deployment name",
deploymentName: "",
buildID: "build-1",
expectedError: "deployment name cannot be empty",
},
{
name: "empty build ID",
deploymentName: "my-deployment",
buildID: "",
expectedError: "build ID cannot be empty",
},
}
for _, tc := range testCases {
s.Run(tc.name, func(s *DeploymentVersionSuite) {
env := s.newTestEnv()
_, err := env.FrontendClient().CreateWorkerDeploymentVersion(s.Context(), &workflowservice.CreateWorkerDeploymentVersionRequest{
Namespace: env.Namespace().String(),
DeploymentVersion: &deploymentpb.WorkerDeploymentVersion{
DeploymentName: tc.deploymentName,
BuildId: tc.buildID,
},
RequestId: env.Tv().Any().String(),
})
s.Error(err)
var invalidArg *serviceerror.InvalidArgument
s.ErrorAs(err, &invalidArg)
s.Contains(invalidArg.Message, tc.expectedError)
})
}
}
func (s *DeploymentVersionSuite) TestCreateWorkerDeploymentVersion_AutoCreatedByPoller_ConflictWithExplicitCreate() {
env := s.newTestEnv()
// Create version via polling (auto-creates deployment and version)
s.startVersionWorkflow(s.Context(), env, env.Tv())
// Try to explicitly create the same version with a different request ID
_, err := env.FrontendClient().CreateWorkerDeploymentVersion(s.Context(), &workflowservice.CreateWorkerDeploymentVersionRequest{
Namespace: env.Namespace().String(),
DeploymentVersion: &deploymentpb.WorkerDeploymentVersion{
DeploymentName: env.Tv().DeploymentSeries(),
BuildId: env.Tv().BuildID(),
},
RequestId: env.Tv().Any().String(),
})
s.Error(err)
var alreadyExists *serviceerror.AlreadyExists
s.ErrorAs(err, &alreadyExists)
}
func (s *DeploymentVersionSuite) TestCreateWorkerDeploymentVersion_MultipleVersions() {
env := s.newTestEnv()
// First create the deployment
_, err := env.FrontendClient().CreateWorkerDeployment(s.Context(), &workflowservice.CreateWorkerDeploymentRequest{
Namespace: env.Namespace().String(),
DeploymentName: env.Tv().DeploymentSeries(),
RequestId: env.Tv().Any().String(),
})
s.NoError(err)
computeConfig1 := &computepb.ComputeConfig{
ScalingGroups: map[string]*computepb.ComputeConfigScalingGroup{
"sg1": {
Provider: computeprovider.TestInvokeComputeProviderValidComputeProvider(),
Scaler: testInvokeScaler(),
},
},
}
computeConfig2 := &computepb.ComputeConfig{
ScalingGroups: map[string]*computepb.ComputeConfigScalingGroup{
"sg2": {
Provider: computeprovider.TestInvokeComputeProviderValidComputeProvider(),
Scaler: testInvokeScaler(),
},
},
}
// Create first version
tv1 := env.Tv().WithBuildIDNumber(1)
_, err = env.FrontendClient().CreateWorkerDeploymentVersion(s.Context(), &workflowservice.CreateWorkerDeploymentVersionRequest{
Namespace: env.Namespace().String(),
DeploymentVersion: &deploymentpb.WorkerDeploymentVersion{
DeploymentName: env.Tv().DeploymentSeries(),
BuildId: tv1.BuildID(),
},
RequestId: env.Tv().Any().String(),
ComputeConfig: computeConfig1,
})
s.NoError(err)
// Create second version
tv2 := env.Tv().WithBuildIDNumber(2)
_, err = env.FrontendClient().CreateWorkerDeploymentVersion(s.Context(), &workflowservice.CreateWorkerDeploymentVersionRequest{
Namespace: env.Namespace().String(),
DeploymentVersion: &deploymentpb.WorkerDeploymentVersion{
DeploymentName: env.Tv().DeploymentSeries(),
BuildId: tv2.BuildID(),
},
RequestId: env.Tv().Any().String(),
ComputeConfig: computeConfig2,
})
s.NoError(err)
// Verify both versions show up in deployment's version summaries
s.EventuallyWithT(func(t *assert.CollectT) {
a := require.New(t)
descResp, err := env.FrontendClient().DescribeWorkerDeployment(s.Context(), &workflowservice.DescribeWorkerDeploymentRequest{
Namespace: env.Namespace().String(),
DeploymentName: env.Tv().DeploymentSeries(),
})
a.NoError(err)
a.Len(descResp.GetWorkerDeploymentInfo().GetVersionSummaries(), 2)
}, 10*time.Second, 500*time.Millisecond)
// Verify compute configs via DescribeWorkerDeploymentVersion
s.EventuallyWithT(func(t *assert.CollectT) {
a := require.New(t)
descV1, err := env.FrontendClient().DescribeWorkerDeploymentVersion(s.Context(), &workflowservice.DescribeWorkerDeploymentVersionRequest{
Namespace: env.Namespace().String(),
Version: tv1.DeploymentVersionString(),
})
a.NoError(err)
a.True(proto.Equal(computeConfig1, descV1.GetWorkerDeploymentVersionInfo().GetComputeConfig()))
}, 10*time.Second, 500*time.Millisecond)
s.EventuallyWithT(func(t *assert.CollectT) {
a := require.New(t)
descV2, err := env.FrontendClient().DescribeWorkerDeploymentVersion(s.Context(), &workflowservice.DescribeWorkerDeploymentVersionRequest{
Namespace: env.Namespace().String(),
Version: tv2.DeploymentVersionString(),
})
a.NoError(err)
a.True(proto.Equal(computeConfig2, descV2.GetWorkerDeploymentVersionInfo().GetComputeConfig()))
}, 10*time.Second, 500*time.Millisecond)
}
// TestCreateWorkerDeploymentVersion_InvalidComputeConfig_ReturnsPromptly verifies that
// creating a version with invalid compute config returns an InvalidArgument error promptly,
// even when the deployment workflow is busy (e.g., processing a recently polled version).
// This is a regression test for a bug where the activity error was retryable, causing the
// deployment workflow to stay busy and the client to time out with "too many requests".
func (s *DeploymentVersionSuite) TestCreateWorkerDeploymentVersion_InvalidComputeConfig_ReturnsPromptly() {
env := s.newTestEnv()
// Create a version via polling, which auto-creates the deployment and keeps
// the deployment workflow active with lazy-creation processing.
s.startVersionWorkflow(s.Context(), env, env.Tv())
// Immediately attempt to create another version on the same deployment with
// invalid compute config. This should return an InvalidArgument error
// promptly (not time out with "too many requests").
tv2 := env.Tv().WithBuildIDNumber(2)
invalidProvider := computeprovider.TestInvokeComputeProviderInvalidComputeProvider()
start := time.Now()
_, err := env.FrontendClient().CreateWorkerDeploymentVersion(s.Context(), &workflowservice.CreateWorkerDeploymentVersionRequest{
Namespace: env.Namespace().String(),
DeploymentVersion: &deploymentpb.WorkerDeploymentVersion{
DeploymentName: env.Tv().DeploymentSeries(),
BuildId: tv2.BuildID(),
},
RequestId: tv2.Any().String(),
ComputeConfig: &computepb.ComputeConfig{
ScalingGroups: map[string]*computepb.ComputeConfigScalingGroup{
"sg1": {Provider: invalidProvider},
},
},
})
elapsed := time.Since(start)
s.Error(err)
var invalidArg *serviceerror.InvalidArgument
s.ErrorAs(err, &invalidArg)
s.Contains(invalidArg.Message, "illegal_field found in config")
// The error should return well before the deployment client's 1-minute
// retry timeout. If it takes longer than 30 seconds, the activity error
// is likely still retryable (the bug we're testing for).
s.Less(elapsed, 30*time.Second, "create-version with invalid compute config should fail promptly, not time out")
}
func (s *DeploymentVersionSuite) TestCreateWorkerDeploymentVersion_InvalidScalingGroups() {
validProvider := computeprovider.TestInvokeComputeProviderValidComputeProvider()
testCases := []struct {
name string
computeConfig *computepb.ComputeConfig
expectedError string
}{
{
name: "invalid compute provider type",
computeConfig: &computepb.ComputeConfig{
ScalingGroups: map[string]*computepb.ComputeConfigScalingGroup{
"sg1": {TaskQueueTypes: nil, Provider: &computepb.ComputeProvider{Type: "invalid-provider"}},
},
},
expectedError: "invalid compute provider type",
},
{
name: "invalid compute provider details",
computeConfig: &computepb.ComputeConfig{
ScalingGroups: map[string]*computepb.ComputeConfigScalingGroup{
"sg1": {TaskQueueTypes: nil, Provider: computeprovider.TestInvokeComputeProviderInvalidComputeProvider()},
},
},
expectedError: "illegal_field found in config",
},
{
name: "two catch-all scaling groups",
computeConfig: &computepb.ComputeConfig{
ScalingGroups: map[string]*computepb.ComputeConfigScalingGroup{
"sg1": {TaskQueueTypes: nil, Provider: validProvider},
"sg2": {TaskQueueTypes: nil, Provider: validProvider},
},
},
expectedError: "only one scaling group can have no task types defined",
},
{
name: "overlapping workflow task queue type",
computeConfig: &computepb.ComputeConfig{
ScalingGroups: map[string]*computepb.ComputeConfigScalingGroup{
"sg1": {TaskQueueTypes: []enumspb.TaskQueueType{enumspb.TASK_QUEUE_TYPE_WORKFLOW}, Provider: validProvider},
"sg2": {TaskQueueTypes: []enumspb.TaskQueueType{enumspb.TASK_QUEUE_TYPE_WORKFLOW, enumspb.TASK_QUEUE_TYPE_ACTIVITY}, Provider: validProvider},
},
},
expectedError: "task type Workflow appears in more than one entry",
},
}
for _, tc := range testCases {
s.Run(tc.name, func(s *DeploymentVersionSuite) {
env := s.newTestEnv()
// Create the deployment first
_, err := env.FrontendClient().CreateWorkerDeployment(s.Context(), &workflowservice.CreateWorkerDeploymentRequest{
Namespace: env.Namespace().String(),
DeploymentName: env.Tv().DeploymentSeries(),
RequestId: env.Tv().Any().String(),
})
s.NoError(err)
_, err = env.FrontendClient().CreateWorkerDeploymentVersion(s.Context(), &workflowservice.CreateWorkerDeploymentVersionRequest{
Namespace: env.Namespace().String(),
DeploymentVersion: &deploymentpb.WorkerDeploymentVersion{
DeploymentName: env.Tv().DeploymentSeries(),
BuildId: env.Tv().BuildID(),
},
RequestId: env.Tv().Any().String(),
ComputeConfig: tc.computeConfig,
})
s.Error(err)
var invalidArg *serviceerror.InvalidArgument
s.ErrorAs(err, &invalidArg)
s.Contains(invalidArg.Message, tc.expectedError)
})
}
}