mirror of
https://github.com/temporalio/temporal.git
synced 2026-08-30 18:41:49 -07:00
## 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 -->
3970 lines
156 KiB
Go
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, ¬Found)
|
|
}
|
|
|
|
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, ¬Found)
|
|
}
|
|
|
|
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)
|
|
})
|
|
}
|
|
}
|