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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
170 changes: 110 additions & 60 deletions internal/temporalcli/commands.worker.deployment.go
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@ import (
"go.temporal.io/sdk/converter"
"go.temporal.io/sdk/worker"
"google.golang.org/protobuf/types/known/fieldmaskpb"
"google.golang.org/protobuf/types/known/timestamppb"
)

type versionSummariesRowType struct {
Expand Down Expand Up @@ -63,9 +64,9 @@ type formattedWorkerDeploymentListEntryType struct {
}

type formattedDrainageInfo struct {
DrainageStatus string `json:"drainageStatus"`
LastChangedTime time.Time `json:"lastChangedTime"`
LastCheckedTime time.Time `json:"lastCheckedTime"`
DrainageStatus string `json:"drainageStatus"`
LastChangedTime *time.Time `json:"lastChangedTime,omitempty"`
LastCheckedTime *time.Time `json:"lastCheckedTime,omitempty"`
}

type formattedTaskQueueInfoRowType struct {
Expand Down Expand Up @@ -109,17 +110,22 @@ type priorityStatsDisplayRow struct {
}

type formattedWorkerDeploymentVersionInfoType struct {
DeploymentName string `json:"deploymentName"`
BuildID string `json:"BuildID"`
CreateTime time.Time `json:"createTime"`
RoutingChangedTime time.Time `json:"routingChangedTime"`
CurrentSinceTime time.Time `json:"currentSinceTime"`
RampingSinceTime time.Time `json:"rampingSinceTime"`
RampPercentage float32 `json:"rampPercentage"`
DrainageInfo formattedDrainageInfo `json:"drainageInfo"`
TaskQueuesInfos []formattedTaskQueueInfoRowType `json:"taskQueuesInfos"`
Metadata map[string]*common.Payload `json:"metadata"`
ComputeConfig *formattedComputeConfig `json:"computeConfig,omitempty"`
DeploymentName string `json:"deploymentName"`
BuildID string `json:"BuildID"`
Status string `json:"status"`
CreateTime time.Time `json:"createTime"`
RoutingChangedTime *time.Time `json:"routingChangedTime,omitempty"`
CurrentSinceTime *time.Time `json:"currentSinceTime,omitempty"`
RampingSinceTime *time.Time `json:"rampingSinceTime,omitempty"`
FirstActivationTime *time.Time `json:"firstActivationTime,omitempty"`
LastCurrentTime *time.Time `json:"lastCurrentTime,omitempty"`
LastDeactivationTime *time.Time `json:"lastDeactivationTime,omitempty"`
RampPercentage float32 `json:"rampPercentage"`
DrainageInfo *formattedDrainageInfo `json:"drainageInfo,omitempty"`
LastModifierIdentity string `json:"lastModifierIdentity,omitempty"`
TaskQueuesInfos []formattedTaskQueueInfoRowType `json:"taskQueuesInfos"`
Metadata map[string]*common.Payload `json:"metadata"`
ComputeConfig *formattedComputeConfig `json:"computeConfig,omitempty"`
}

type formattedComputeConfig struct {
Expand Down Expand Up @@ -308,6 +314,13 @@ func drainageStatusProtoToStr(status enumspb.VersionDrainageStatus) (string, err
}
}

// versionStatusProtoToStr converts a version status to a lowercase string, e.g.
// WORKER_DEPLOYMENT_VERSION_STATUS_CURRENT becomes "current". Values unknown to
// this CLI fall back to their numeric enum string rather than failing.
func versionStatusProtoToStr(status enumspb.WorkerDeploymentVersionStatus) string {
return strings.ToLower(strings.TrimPrefix(status.String(), "WORKER_DEPLOYMENT_VERSION_STATUS_"))
}

func taskQueueTypeProtoToStr(taskQueueType enumspb.TaskQueueType) (string, error) {
switch taskQueueType {
case enumspb.TASK_QUEUE_TYPE_UNSPECIFIED:
Expand Down Expand Up @@ -368,20 +381,42 @@ func formatTaskQueuesInfosProto(tqis []*workflowservice.DescribeWorkerDeployment
return tqiRows, nil
}

func formatDrainageInfoProto(drainageInfo *deploymentpb.VersionDrainageInfo) (formattedDrainageInfo, error) {
// optionalTimestampToTime converts a timestamp that the server may leave unset
// (e.g. RampingSinceTime for a version that was never ramped) to a time.Time.
// Both a nil timestamp and the Unix epoch are treated as unset and return the
// zero time.Time, so that text output omits the field instead of rendering it
// as "a long while ago".
func optionalTimestampToTime(t *timestamppb.Timestamp) time.Time {
if t == nil || (t.GetSeconds() == 0 && t.GetNanos() == 0) {
return time.Time{}
}
return t.AsTime()
}

// optionalTimestampToTimePtr is like optionalTimestampToTime but returns nil for
// an unset timestamp, so that JSON output omits the field.
func optionalTimestampToTimePtr(t *timestamppb.Timestamp) *time.Time {
tm := optionalTimestampToTime(t)
if tm.IsZero() {
return nil
}
return &tm
}

func formatDrainageInfoProto(drainageInfo *deploymentpb.VersionDrainageInfo) (*formattedDrainageInfo, error) {
if drainageInfo == nil {
return formattedDrainageInfo{}, nil
return nil, nil
}

drainageStr, err := drainageStatusProtoToStr(drainageInfo.GetStatus())
if err != nil {
return formattedDrainageInfo{}, err
return nil, err
}

return formattedDrainageInfo{
return &formattedDrainageInfo{
DrainageStatus: drainageStr,
LastChangedTime: drainageInfo.GetLastChangedTime().AsTime(),
LastCheckedTime: drainageInfo.GetLastCheckedTime().AsTime(),
LastChangedTime: optionalTimestampToTimePtr(drainageInfo.GetLastChangedTime()),
LastCheckedTime: optionalTimestampToTimePtr(drainageInfo.GetLastCheckedTime()),
}, nil
}

Expand Down Expand Up @@ -526,17 +561,22 @@ func workerDeploymentVersionInfoProtoToRows(deploymentInfo *deploymentpb.WorkerD
computeConfig := formatComputeConfigProto(deploymentInfo.GetComputeConfig())

return formattedWorkerDeploymentVersionInfoType{
DeploymentName: deploymentInfo.GetDeploymentVersion().GetDeploymentName(),
BuildID: deploymentInfo.GetDeploymentVersion().GetBuildId(),
CreateTime: deploymentInfo.GetCreateTime().AsTime(),
RoutingChangedTime: deploymentInfo.GetRoutingChangedTime().AsTime(),
CurrentSinceTime: deploymentInfo.GetCurrentSinceTime().AsTime(),
RampingSinceTime: deploymentInfo.GetRampingSinceTime().AsTime(),
RampPercentage: deploymentInfo.GetRampPercentage(),
DrainageInfo: drainage,
TaskQueuesInfos: tqi,
Metadata: deploymentInfo.GetMetadata().GetEntries(),
ComputeConfig: computeConfig,
DeploymentName: deploymentInfo.GetDeploymentVersion().GetDeploymentName(),
BuildID: deploymentInfo.GetDeploymentVersion().GetBuildId(),
Status: versionStatusProtoToStr(deploymentInfo.GetStatus()),
CreateTime: deploymentInfo.GetCreateTime().AsTime(),
RoutingChangedTime: optionalTimestampToTimePtr(deploymentInfo.GetRoutingChangedTime()),
CurrentSinceTime: optionalTimestampToTimePtr(deploymentInfo.GetCurrentSinceTime()),
RampingSinceTime: optionalTimestampToTimePtr(deploymentInfo.GetRampingSinceTime()),
FirstActivationTime: optionalTimestampToTimePtr(deploymentInfo.GetFirstActivationTime()),
LastCurrentTime: optionalTimestampToTimePtr(deploymentInfo.GetLastCurrentTime()),
LastDeactivationTime: optionalTimestampToTimePtr(deploymentInfo.GetLastDeactivationTime()),
RampPercentage: deploymentInfo.GetRampPercentage(),
DrainageInfo: drainage,
LastModifierIdentity: deploymentInfo.GetLastModifierIdentity(),
TaskQueuesInfos: tqi,
Metadata: deploymentInfo.GetMetadata().GetEntries(),
ComputeConfig: computeConfig,
}, nil
}

Expand Down Expand Up @@ -598,35 +638,45 @@ func printWorkerDeploymentVersionInfoProto(cctx *CommandContext, deploymentInfo
if err != nil {
return err
}
drainageLastChangedTime = deploymentInfo.GetDrainageInfo().GetLastChangedTime().AsTime()
drainageLastCheckedTime = deploymentInfo.GetDrainageInfo().GetLastCheckedTime().AsTime()
drainageLastChangedTime = optionalTimestampToTime(deploymentInfo.GetDrainageInfo().GetLastChangedTime())
drainageLastCheckedTime = optionalTimestampToTime(deploymentInfo.GetDrainageInfo().GetLastCheckedTime())
}
computeConfigSummary := computeConfigSummaryStr(deploymentInfo.GetComputeConfig())

printMe := struct {
DeploymentName string
BuildID string
Status string
CreateTime time.Time
RoutingChangedTime time.Time `cli:",cardOmitEmpty"`
CurrentSinceTime time.Time `cli:",cardOmitEmpty"`
RampingSinceTime time.Time `cli:",cardOmitEmpty"`
FirstActivationTime time.Time `cli:",cardOmitEmpty"`
LastCurrentTime time.Time `cli:",cardOmitEmpty"`
LastDeactivationTime time.Time `cli:",cardOmitEmpty"`
RampPercentage float32
DrainageStatus string `cli:",cardOmitEmpty"`
DrainageLastChangedTime time.Time `cli:",cardOmitEmpty"`
DrainageLastCheckedTime time.Time `cli:",cardOmitEmpty"`
LastModifierIdentity string `cli:",cardOmitEmpty"`
Metadata map[string]*common.Payload `cli:",cardOmitEmpty"`
ComputeConfigSummary string `cli:",cardOmitEmpty"`
}{
DeploymentName: deploymentInfo.GetDeploymentVersion().GetDeploymentName(),
BuildID: deploymentInfo.GetDeploymentVersion().GetBuildId(),
Status: fDeploymentInfo.Status,
CreateTime: deploymentInfo.GetCreateTime().AsTime(),
RoutingChangedTime: deploymentInfo.GetRoutingChangedTime().AsTime(),
CurrentSinceTime: deploymentInfo.GetCurrentSinceTime().AsTime(),
RampingSinceTime: deploymentInfo.GetRampingSinceTime().AsTime(),
RoutingChangedTime: optionalTimestampToTime(deploymentInfo.GetRoutingChangedTime()),
CurrentSinceTime: optionalTimestampToTime(deploymentInfo.GetCurrentSinceTime()),
RampingSinceTime: optionalTimestampToTime(deploymentInfo.GetRampingSinceTime()),
FirstActivationTime: optionalTimestampToTime(deploymentInfo.GetFirstActivationTime()),
LastCurrentTime: optionalTimestampToTime(deploymentInfo.GetLastCurrentTime()),
LastDeactivationTime: optionalTimestampToTime(deploymentInfo.GetLastDeactivationTime()),
RampPercentage: deploymentInfo.GetRampPercentage(),
DrainageStatus: drainageStr,
DrainageLastChangedTime: drainageLastChangedTime,
DrainageLastCheckedTime: drainageLastCheckedTime,
LastModifierIdentity: deploymentInfo.GetLastModifierIdentity(),
Metadata: deploymentInfo.GetMetadata().GetEntries(),
ComputeConfigSummary: computeConfigSummary,
}
Expand Down Expand Up @@ -1350,18 +1400,18 @@ func (c *TemporalWorkerDeploymentCreateVersionCommand) run(cctx *CommandContext,
requestID := uuid.NewString()

providerType, detailsPayload, err := computeProviderConfig(&ComputeConfigArgs{
awsLambdaFunctionArn: c.AwsLambdaFunctionArn,
awsLambdaAssumeRoleArn: c.AwsLambdaAssumeRoleArn,
awsLambdaAssumeRoleExternalId: c.AwsLambdaAssumeRoleExternalId,
awsLambdaSkipRoleAndExternalId: c.AwsLambdaSkipRoleAndExternalId,
awsAgentcoreEndpointArn: c.AwsAgentcoreEndpointArn,
awsAgentcoreAssumeRoleArn: c.AwsAgentcoreAssumeRoleArn,
awsAgentcoreAssumeRoleExternalId: c.AwsAgentcoreAssumeRoleExternalId,
awsLambdaFunctionArn: c.AwsLambdaFunctionArn,
awsLambdaAssumeRoleArn: c.AwsLambdaAssumeRoleArn,
awsLambdaAssumeRoleExternalId: c.AwsLambdaAssumeRoleExternalId,
awsLambdaSkipRoleAndExternalId: c.AwsLambdaSkipRoleAndExternalId,
awsAgentcoreEndpointArn: c.AwsAgentcoreEndpointArn,
awsAgentcoreAssumeRoleArn: c.AwsAgentcoreAssumeRoleArn,
awsAgentcoreAssumeRoleExternalId: c.AwsAgentcoreAssumeRoleExternalId,
awsAgentcoreSkipRoleAndExternalId: c.AwsAgentcoreSkipRoleAndExternalId,
gcpCloudRunProject: c.GcpCloudRunProject,
gcpCloudRunRegion: c.GcpCloudRunRegion,
gcpCloudRunWorkerPool: c.GcpCloudRunWorkerPool,
gcpCloudRunServiceAccount: c.GcpCloudRunServiceAccount,
gcpCloudRunProject: c.GcpCloudRunProject,
gcpCloudRunRegion: c.GcpCloudRunRegion,
gcpCloudRunWorkerPool: c.GcpCloudRunWorkerPool,
gcpCloudRunServiceAccount: c.GcpCloudRunServiceAccount,
})
if err != nil {
return err
Expand Down Expand Up @@ -1447,19 +1497,19 @@ func (c *TemporalWorkerDeploymentUpdateVersionComputeConfigCommand) run(cctx *Co
}

computeConfigArgs := &ComputeConfigArgs{
awsLambdaFunctionArn: c.AwsLambdaFunctionArn,
awsLambdaAssumeRoleArn: c.AwsLambdaAssumeRoleArn,
awsLambdaAssumeRoleExternalId: c.AwsLambdaAssumeRoleExternalId,
awsLambdaSkipRoleAndExternalId: c.AwsLambdaSkipRoleAndExternalId,
awsAgentcoreEndpointArn: c.AwsAgentcoreEndpointArn,
awsAgentcoreAssumeRoleArn: c.AwsAgentcoreAssumeRoleArn,
awsAgentcoreAssumeRoleExternalId: c.AwsAgentcoreAssumeRoleExternalId,
awsAgentcoreSkipRoleAndExternalId: c.AwsAgentcoreSkipRoleAndExternalId,
gcpCloudRunProject: c.GcpCloudRunProject,
gcpCloudRunRegion: c.GcpCloudRunRegion,
gcpCloudRunWorkerPool: c.GcpCloudRunWorkerPool,
gcpCloudRunServiceAccount: c.GcpCloudRunServiceAccount,
}
awsLambdaFunctionArn: c.AwsLambdaFunctionArn,
awsLambdaAssumeRoleArn: c.AwsLambdaAssumeRoleArn,
awsLambdaAssumeRoleExternalId: c.AwsLambdaAssumeRoleExternalId,
awsLambdaSkipRoleAndExternalId: c.AwsLambdaSkipRoleAndExternalId,
awsAgentcoreEndpointArn: c.AwsAgentcoreEndpointArn,
awsAgentcoreAssumeRoleArn: c.AwsAgentcoreAssumeRoleArn,
awsAgentcoreAssumeRoleExternalId: c.AwsAgentcoreAssumeRoleExternalId,
awsAgentcoreSkipRoleAndExternalId: c.AwsAgentcoreSkipRoleAndExternalId,
gcpCloudRunProject: c.GcpCloudRunProject,
gcpCloudRunRegion: c.GcpCloudRunRegion,
gcpCloudRunWorkerPool: c.GcpCloudRunWorkerPool,
gcpCloudRunServiceAccount: c.GcpCloudRunServiceAccount,
}

if c.Remove {
if computeConfigArgs.hasAwsLambdaArgs() || computeConfigArgs.hasAwsAgentcoreArgs() || computeConfigArgs.hasGcpCloudRunArgs() ||
Expand Down
Original file line number Diff line number Diff line change
@@ -1,12 +1,18 @@
package temporalcli

import (
"bytes"
"encoding/json"
"testing"
"time"

"github.com/stretchr/testify/require"
"github.com/temporalio/cli/internal/printer"
computepb "go.temporal.io/api/compute/v1"
deploymentpb "go.temporal.io/api/deployment/v1"
enumspb "go.temporal.io/api/enums/v1"
"go.temporal.io/sdk/converter"
"google.golang.org/protobuf/types/known/timestamppb"
)

func TestScalerTypeForProvider(t *testing.T) {
Expand Down Expand Up @@ -213,3 +219,85 @@ func TestFormatComputeConfigProto_ScalerBounds(t *testing.T) {
require.Empty(t, sg.Scaler.ScaleDownStabilization)
require.Equal(t, "gcp-cloud-run", computeConfigSummaryStr(ccNoBounds))
}

func TestPrintWorkerDeploymentVersionInfoProto_AllFields(t *testing.T) {
ts := func(sec int64) *timestamppb.Timestamp { return timestamppb.New(time.Unix(sec, 0).UTC()) }
info := &deploymentpb.WorkerDeploymentVersionInfo{
Status: enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_DRAINING,
DeploymentVersion: &deploymentpb.WorkerDeploymentVersion{DeploymentName: "my-deployment", BuildId: "v1"},
CreateTime: ts(1000),
RoutingChangedTime: ts(2000),
FirstActivationTime: ts(3000),
LastCurrentTime: ts(4000),
LastDeactivationTime: ts(5000),
DrainageInfo: &deploymentpb.VersionDrainageInfo{
Status: enumspb.VERSION_DRAINAGE_STATUS_DRAINING,
LastChangedTime: ts(6000),
LastCheckedTime: ts(7000),
},
LastModifierIdentity: "some-identity",
}

t.Run("text", func(t *testing.T) {
var buf bytes.Buffer
cctx := &CommandContext{Printer: &printer.Printer{Output: &buf}}
require.NoError(t, printWorkerDeploymentVersionInfoProto(cctx, info, nil, "Worker Deployment Version:", printVersionInfoOptions{}))
out := buf.String()
for _, field := range []string{
"Status", "FirstActivationTime", "LastCurrentTime", "LastDeactivationTime",
"DrainageStatus", "DrainageLastChangedTime", "DrainageLastCheckedTime", "LastModifierIdentity",
} {
require.Contains(t, out, field)
}
require.Contains(t, out, "draining")
require.Contains(t, out, "some-identity")
require.Contains(t, out, time.Unix(3000, 0).UTC().Format(time.RFC3339))
// Never current or ramping, so these must not be rendered.
require.NotContains(t, out, "CurrentSinceTime")
require.NotContains(t, out, "RampingSinceTime")
})

t.Run("json", func(t *testing.T) {
var buf bytes.Buffer
cctx := &CommandContext{JSONOutput: true, Printer: &printer.Printer{Output: &buf, JSON: true}}
require.NoError(t, printWorkerDeploymentVersionInfoProto(cctx, info, nil, "", printVersionInfoOptions{}))
var out formattedWorkerDeploymentVersionInfoType
require.NoError(t, json.Unmarshal(buf.Bytes(), &out))
require.Equal(t, "draining", out.Status)
require.Equal(t, time.Unix(3000, 0).UTC(), *out.FirstActivationTime)
require.Equal(t, time.Unix(4000, 0).UTC(), *out.LastCurrentTime)
require.Equal(t, time.Unix(5000, 0).UTC(), *out.LastDeactivationTime)
require.Equal(t, "some-identity", out.LastModifierIdentity)
require.NotNil(t, out.DrainageInfo)
require.Equal(t, "draining", out.DrainageInfo.DrainageStatus)
require.Equal(t, time.Unix(6000, 0).UTC(), *out.DrainageInfo.LastChangedTime)
require.Equal(t, time.Unix(7000, 0).UTC(), *out.DrainageInfo.LastCheckedTime)
require.Nil(t, out.CurrentSinceTime)
require.Nil(t, out.RampingSinceTime)
})

t.Run("json omits unset drainage info", func(t *testing.T) {
var buf bytes.Buffer
cctx := &CommandContext{JSONOutput: true, Printer: &printer.Printer{Output: &buf, JSON: true}}
current := &deploymentpb.WorkerDeploymentVersionInfo{
Status: enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_CURRENT,
DeploymentVersion: &deploymentpb.WorkerDeploymentVersion{DeploymentName: "my-deployment", BuildId: "v2"},
CreateTime: ts(1000),
}
require.NoError(t, printWorkerDeploymentVersionInfoProto(cctx, current, nil, "", printVersionInfoOptions{}))
var raw map[string]any
require.NoError(t, json.Unmarshal(buf.Bytes(), &raw))
require.Equal(t, "current", raw["status"])
require.NotContains(t, raw, "drainageInfo")
})
}

func TestVersionStatusProtoToStr(t *testing.T) {
require.Equal(t, "unspecified", versionStatusProtoToStr(enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_UNSPECIFIED))
require.Equal(t, "inactive", versionStatusProtoToStr(enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_INACTIVE))
require.Equal(t, "current", versionStatusProtoToStr(enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_CURRENT))
require.Equal(t, "ramping", versionStatusProtoToStr(enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_RAMPING))
require.Equal(t, "draining", versionStatusProtoToStr(enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_DRAINING))
require.Equal(t, "drained", versionStatusProtoToStr(enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_DRAINED))
require.Equal(t, "created", versionStatusProtoToStr(enumspb.WORKER_DEPLOYMENT_VERSION_STATUS_CREATED))
}
Loading
Loading