worker_versioning

package
v1.31.2 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Jul 7, 2026 License: MIT Imports: 28 Imported by: 1

Documentation

Index

Constants

View Source
const (
	BuildIdSearchAttributePrefixPinned = "pinned"

	BuildIdSearchAttributeDelimiter = ":"
	BuildIdSearchAttributeEscape    = "|"
	// UnversionedSearchAttribute is the sentinel value used to mark all unversioned workflows
	UnversionedSearchAttribute = buildIdSearchAttributePrefixUnversioned
	UnversionedVersionId       = "__unversioned__"

	// ErrPinnedVersionNotInTaskQueueSubstring is the key substring used to identify
	// when a pinned version is not present in a task queue. This is used for error
	// classification in batch operations.
	ErrPinnedVersionNotInTaskQueueSubstring = "is not present in task queue"

	// WorkerDeploymentVersionIdDelimiterV31 will be deleted once we stop supporting v31 version string fields
	// in external and internal APIs. Until then, both delimiters are banned in deployment name. All
	// deprecated version string fields in APIs keep using the old delimiter. Workflow SA uses new delimiter.
	WorkerDeploymentVersionIDDelimiterV31   = "."
	WorkerDeploymentVersionDelimiter        = ":"
	WorkerDeploymentVersionWorkflowIDEscape = "|"

	// Prefixes, Delimeters and Keys that are used in the internal entity workflows backing worker-versioning
	WorkerDeploymentWorkflowIDPrefix             = "temporal-sys-worker-deployment"
	WorkerDeploymentVersionWorkflowIDPrefix      = "temporal-sys-worker-deployment-version"
	WorkerDeploymentVersionWorkflowIDInitialSize = len(WorkerDeploymentVersionWorkflowIDPrefix) + len(WorkerDeploymentVersionDelimiter) // 39
	WorkerDeploymentNameFieldName                = "WorkerDeploymentName"
	WorkerDeploymentBuildIDFieldName             = "BuildID"
)

Variables

This section is empty.

Functions

func AddV31VersioningInfoToV32 added in v1.28.0

We store versioning info in the modern v0.32 format, so call this before returning the object to readers to mutatively populate the missing fields.

func AssignedBuildIdSearchAttribute added in v1.24.0

func AssignedBuildIdSearchAttribute(buildId string) string

AssignedBuildIdSearchAttribute returns the search attribute value for the currently assigned build ID

func BuildIDToStringV32 added in v1.30.0

func BuildIDToStringV32(deploymentName, buildID string) string

func BuildIdFromCapabilities added in v1.28.0

func BuildIdFromCapabilities(capabilities *commonpb.WorkerVersionCapabilities, options *deploymentpb.WorkerDeploymentOptions) string

func BuildIdIfUsingVersioning added in v1.24.0

func BuildIdIfUsingVersioning(stamp *commonpb.WorkerVersionStamp) string

BuildIdIfUsingVersioning returns the given WorkerVersionStamp if it is using versioning, otherwise returns nil.

func CalculateTaskQueueVersioningInfo added in v1.27.0

CalculateTaskQueueVersioningInfo calculates the current and ramping versioning info for a task queue.

func CleanupOldDeletedVersions added in v1.30.0

func CleanupOldDeletedVersions(deploymentData *persistencespb.WorkerDeploymentData, maxVersions int) bool

CleanupOldDeletedVersions removes versions deleted more than 7 days ago. Also removes more deleted versions if the limit is being exceeded. Never removes undeleted versions. Deprecated. Versions now are deleted serially without using the deleted flag in versionData. TODO: remove this cleanup logic after next major release.

func ConvertOverrideToV32 added in v1.28.0

func ConvertOverrideToV32(override *workflowpb.VersioningOverride) *workflowpb.VersioningOverride

ConvertOverrideToV32 reads from deprecated fields and returns a new object with ONLY the equivalent non-deprecated v0.32 fields. Should be used to replace any passed in override that is stored in persistence.

func CountDeploymentVersions added in v1.30.0

func CountDeploymentVersions(deployments *persistencespb.DeploymentData) int

func DeploymentFromCapabilities added in v1.26.2

func DeploymentFromCapabilities(capabilities *commonpb.WorkerVersionCapabilities, options *deploymentpb.WorkerDeploymentOptions) (*deploymentpb.Deployment, error)

DeploymentFromCapabilities returns the deployment if it is using versioning V3, otherwise nil. It returns the deployment from the `options` if present, otherwise, from `capabilities`,

func DeploymentFromDeploymentVersion added in v1.27.0

func DeploymentFromDeploymentVersion(dv *deploymentspb.WorkerDeploymentVersion) *deploymentpb.Deployment

DeploymentFromDeploymentVersion Temporary helper function to convert WorkerDeploymentVersion to Deployment proto until we update code to use the new proto in all places.

func DeploymentFromExternalDeploymentVersion added in v1.28.0

func DeploymentFromExternalDeploymentVersion(dv *deploymentpb.WorkerDeploymentVersion) *deploymentpb.Deployment

DeploymentFromExternalDeploymentVersion Temporary helper function to convert WorkerDeploymentVersion to Deployment proto until we update code to use the new proto in all places.

func DeploymentIfValid added in v1.27.0

func DeploymentIfValid(d *deploymentpb.Deployment) *deploymentpb.Deployment

DeploymentIfValid returns the deployment back if is both of its fields have value.

func DeploymentNameFromCapabilities added in v1.28.0

func DeploymentNameFromCapabilities(capabilities *commonpb.WorkerVersionCapabilities, options *deploymentpb.WorkerDeploymentOptions) string

func DeploymentOrVersion added in v1.27.0

DeploymentOrVersion Temporary helper function to return a Deployment based on passed Deployment or WorkerDeploymentVersion objects, if `v` is not nil, it'll take precedence.

func DeploymentVersionFromDeployment added in v1.27.0

func DeploymentVersionFromDeployment(deployment *deploymentpb.Deployment) *deploymentspb.WorkerDeploymentVersion

DeploymentVersionFromDeployment Temporary helper function to convert Deployment to WorkerDeploymentVersion proto until we update code to use the new proto in all places.

func DeploymentVersionFromOptions added in v1.27.0

func DeploymentVersionFromOptions(options *deploymentpb.WorkerDeploymentOptions) *deploymentspb.WorkerDeploymentVersion

func DirectiveDeployment added in v1.27.0

func DirectiveDeployment(directive *taskqueuespb.TaskVersionDirective) *deploymentpb.Deployment

DirectiveDeployment Temporary function until Directive proto is removed.

func ExternalWorkerDeploymentVersionFromDeployment added in v1.28.0

func ExternalWorkerDeploymentVersionFromDeployment(deployment *deploymentpb.Deployment) *deploymentpb.WorkerDeploymentVersion

ExternalWorkerDeploymentVersionFromDeployment Temporary helper function to convert Deployment to WorkerDeploymentVersion proto until we update code to use the new proto in all places.

func ExternalWorkerDeploymentVersionFromStringV31 added in v1.28.0

func ExternalWorkerDeploymentVersionFromStringV31(s string) *deploymentpb.WorkerDeploymentVersion

func ExternalWorkerDeploymentVersionFromVersion added in v1.28.0

func ExternalWorkerDeploymentVersionFromVersion(version *deploymentspb.WorkerDeploymentVersion) *deploymentpb.WorkerDeploymentVersion

ExternalWorkerDeploymentVersionFromVersion Temporary helper function to convert internal Worker Deployment to WorkerDeploymentVersion proto until we update code to use the new proto in all places.

func ExternalWorkerDeploymentVersionToString added in v1.28.0

func ExternalWorkerDeploymentVersionToString(v *deploymentpb.WorkerDeploymentVersion) string

func ExternalWorkerDeploymentVersionToStringV31 added in v1.28.0

func ExternalWorkerDeploymentVersionToStringV31(v *deploymentpb.WorkerDeploymentVersion) string

func ExtractVersioningBehaviorFromOverride added in v1.28.0

func ExtractVersioningBehaviorFromOverride(override *workflowpb.VersioningOverride) enumspb.VersioningBehavior

func FindBuildId

func FindBuildId(versioningData *persistencespb.VersioningData, buildId string) (setIndex, indexInSet int)

FindBuildId finds a build ID in the version data's sets, returning (set index, index within that set). Returns -1, -1 if not found.

func FindOldDeploymentVersion added in v1.31.0

func FindOldDeploymentVersion(deployments *persistencespb.DeploymentData, v *deploymentspb.WorkerDeploymentVersion) int

func FindTargetDeploymentVersionAndRevisionNumberForWorkflowID added in v1.30.0

func FindTargetDeploymentVersionAndRevisionNumberForWorkflowID(
	current *deploymentspb.WorkerDeploymentVersion,
	currentRevisionNumber int64,
	ramping *deploymentspb.WorkerDeploymentVersion,
	rampingPercentage float32,
	rampingRevisionNumber int64,
	workflowId string,
	useRampingVersion bool,
) (*deploymentspb.WorkerDeploymentVersion, int64)

FindTargetDeploymentVersionAndRevisionNumberForWorkflowID returns the deployment version and revision number (if applicable) for the particular workflow ID based on the versioning info of the task queue. Nil means unversioned.

func FormatPinnedVersionNotInTaskQueueError added in v1.30.0

func FormatPinnedVersionNotInTaskQueueError(deploymentName, buildID, taskQueue string, taskQueueType enumspb.TaskQueueType) string

FormatPinnedVersionNotInTaskQueueError formats the error message when a pinned version is not present in a task queue.

func GetOverridePinnedVersion added in v1.28.0

func GetOverridePinnedVersion(override *workflowpb.VersioningOverride) *deploymentpb.WorkerDeploymentVersion

func HasDeploymentVersion added in v1.28.0

func HasDeploymentVersion(deployments *persistencespb.DeploymentData, v *deploymentspb.WorkerDeploymentVersion) bool

func IsUnversionedOrAssignedBuildIdSearchAttribute added in v1.24.0

func IsUnversionedOrAssignedBuildIdSearchAttribute(buildId string) bool

IsUnversionedOrAssignedBuildIdSearchAttribute returns the value is "unversioned" or "assigned:<bld>"

func MakeBuildIdDirective added in v1.24.0

func MakeBuildIdDirective(buildId string) *taskqueuespb.TaskVersionDirective

func MakeDirectiveForWorkflowTask added in v1.22.0

func MakeDirectiveForWorkflowTask(
	inheritedBuildId string,
	assignedBuildId string,
	stamp *commonpb.WorkerVersionStamp,
	hasCompletedWorkflowTask bool,
	behavior enumspb.VersioningBehavior,
	deployment *deploymentpb.Deployment,
	revisionNumber int64,
	useRampingVersion bool,
) *taskqueuespb.TaskVersionDirective

MakeDirectiveForWorkflowTask returns a versioning directive based on the following parameters: - inheritedBuildId: build ID inherited from a past/previous wf execution (for Child WF or CaN) - assignedBuildId: the build ID to which the WF is currently assigned (i.e. mutable state's AssginedBuildId) - stamp: the latest versioning stamp of the execution (only needed for old versioning) - hasCompletedWorkflowTask: if the wf has completed any WFT - behavior: workflow's effective behavior - deployment: workflow's effective deployment

func MakeUseAssignmentRulesDirective added in v1.24.0

func MakeUseAssignmentRulesDirective() *taskqueuespb.TaskVersionDirective

func OverrideIsPinned added in v1.28.0

func OverrideIsPinned(override *workflowpb.VersioningOverride) bool

func PickFinalCurrentAndRamping added in v1.30.0

func PickFinalCurrentAndRamping(
	current *deploymentspb.DeploymentVersionData,
	ramping *deploymentspb.DeploymentVersionData,
	currentVersionRoutingConfig *deploymentpb.RoutingConfig,
	rampingVersionRoutingConfig *deploymentpb.RoutingConfig,
) (
	finalCurrent *deploymentspb.WorkerDeploymentVersion,
	finalCurrentRev int64,
	finalCurrentUpdateTime time.Time,
	finalRamping *deploymentspb.WorkerDeploymentVersion,
	isRamping bool,
	finalRampPercentage float32,
	finalRampingRev int64,
	finalRampingUpdateTime time.Time,
)

PickFinalCurrentAndRamping determines the effective "current" and "ramping" deployment versions by comparing timestamps from the legacy deployment data (old format) and the RoutingConfig (new format). It returns: - final current deployment version and its revision number (0 for old format) - final ramping deployment version, its revision number (0 for old format), and ramp percentage

func PinnedBuildIdSearchAttribute added in v1.26.2

func PinnedBuildIdSearchAttribute(version string) string

PinnedBuildIdSearchAttribute creates the pinned search attribute for the BuildIds list, used as a visibility optimization. For pinned workflows using WorkerDeployment APIs (ms.GetEffectiveVersioningBehavior() == PINNED && ms.executionInfo.VersioningInfo.Version != ""), this will be `pinned:<version>`. The version used will be the override version if set, or the versioningInfo.Version.

If deprecated Deployment-based APIs are in use and the workflow is pinned, `pinned:<deployment_series_name>:<deployment_build_id>` will. The values used will be the override deployment_series and build_id if set, or versioningInfo.Deployment.

If the workflow becomes unpinned or unversioned, this entry will be removed from that list.

func StampFromBuildId added in v1.26.2

func StampFromBuildId(buildId string) *commonpb.WorkerVersionStamp

func StampFromCapabilities added in v1.24.0

func StampIfUsingVersioning added in v1.22.0

func StampIfUsingVersioning(stamp *commonpb.WorkerVersionStamp) *commonpb.WorkerVersionStamp

StampIfUsingVersioning returns the given WorkerVersionStamp if it is using versioning, otherwise returns nil.

func UnversionedBuildIdSearchAttribute

func UnversionedBuildIdSearchAttribute(buildId string) string

UnversionedBuildIdSearchAttribute returns the search attribute value for an unversioned build ID

func ValidateDeployment added in v1.26.2

func ValidateDeployment(deployment *deploymentpb.Deployment) error

ValidateDeployment returns error if the deployment is nil or it has empty build ID or deployment name.

func ValidateDeploymentVersion added in v1.27.0

func ValidateDeploymentVersion(version *deploymentspb.WorkerDeploymentVersion, maxIDLengthLimit int) error

ValidateDeploymentVersion returns error if the deployment version is not a valid entity.

func ValidateDeploymentVersionFields added in v1.30.0

func ValidateDeploymentVersionFields(fieldName string, field string, maxIDLengthLimit int) error

ValidateDeploymentVersionFields is a helper that verifies if the fields within a Worker Deployment Version are valid

func ValidateDeploymentVersionStringV31 added in v1.28.0

func ValidateDeploymentVersionStringV31(version string) (*deploymentspb.WorkerDeploymentVersion, error)

ValidateDeploymentVersionStringV31 returns error if the deployment version is nil or it has empty version or deployment name.

func ValidateTaskVersionDirective added in v1.27.0

func ValidateTaskVersionDirective(
	directive *taskqueuespb.TaskVersionDirective,
	wfBehavior enumspb.VersioningBehavior,
	wfDeployment *deploymentpb.Deployment,
	scheduledDeployment *deploymentpb.Deployment,
) error

func ValidateVersioningOverride added in v1.26.2

func ValidateVersioningOverride(ctx context.Context,
	override *workflowpb.VersioningOverride,
	matchingClient resource.MatchingClient,
	versionMembershipCache VersionMembershipCache,
	tq string,
	tqType enumspb.TaskQueueType,
	namespaceID string) error

func VersionStampToBuildIdSearchAttribute

func VersionStampToBuildIdSearchAttribute(stamp *commonpb.WorkerVersionStamp) string

VersionStampToBuildIdSearchAttribute returns the search attribute value for a version stamp

func VersionedBuildIdSearchAttribute

func VersionedBuildIdSearchAttribute(buildId string) string

VersionedBuildIdSearchAttribute returns the search attribute value for a versioned build ID

func WorkerDeploymentVersionFromStringV31 added in v1.28.0

func WorkerDeploymentVersionFromStringV31(s string) (*deploymentspb.WorkerDeploymentVersion, error)

func WorkerDeploymentVersionFromStringV32 added in v1.28.0

func WorkerDeploymentVersionFromStringV32(s string) (*deploymentspb.WorkerDeploymentVersion, error)

func WorkerDeploymentVersionToStringV31 added in v1.28.0

func WorkerDeploymentVersionToStringV31(v *deploymentspb.WorkerDeploymentVersion) string

func WorkerDeploymentVersionToStringV32 added in v1.28.0

func WorkerDeploymentVersionToStringV32(v *deploymentspb.WorkerDeploymentVersion) string

func WorkflowsExistForBuildId

func WorkflowsExistForBuildId(ctx context.Context, visibilityManager manager.VisibilityManager, ns *namespace.Namespace, taskQueue, buildId string) (bool, error)

Types

type IsWFTaskQueueInVersionDetector added in v1.28.0

type IsWFTaskQueueInVersionDetector = func(ctx context.Context, namespaceID, tq string, version *deploymentpb.WorkerDeploymentVersion) (bool, error)

func GetIsWFTaskQueueInVersionDetector added in v1.28.0

func GetIsWFTaskQueueInVersionDetector(matchingClient resource.MatchingClient, versionMembershipCache VersionMembershipCache) IsWFTaskQueueInVersionDetector

type ReactivationSignalCache added in v1.30.2

type ReactivationSignalCache interface {
	// ShouldSendSignal returns true if signal should be sent (not recently sent)
	// and atomically marks it as sent. Returns false if recently sent.
	ShouldSendSignal(namespaceID, deploymentName, buildID string) bool
}

ReactivationSignalCache deduplicates reactivation signals to version workflows.

Implementations are expected to be safe for concurrent use.

func NewReactivationSignalCache added in v1.30.2

func NewReactivationSignalCache(c cache.Cache, metricsHandler metrics.Handler) ReactivationSignalCache

NewReactivationSignalCache wraps the provided cache with a typed API and metrics.

type ReactivationSignalCacheImpl added in v1.30.2

type ReactivationSignalCacheImpl struct {
	cache.Cache
	// contains filtered or unexported fields
}

ReactivationSignalCache deduplicates reactivation signals to version workflows.

Implementations are expected to be safe for concurrent use.

func (*ReactivationSignalCacheImpl) ShouldSendSignal added in v1.30.2

func (c *ReactivationSignalCacheImpl) ShouldSendSignal(
	namespaceID, deploymentName, buildID string,
) bool

type RoutingInfo added in v1.30.2

type RoutingInfo struct {
	Current               *deploymentspb.WorkerDeploymentVersion
	CurrentRevisionNumber int64
	Ramping               *deploymentspb.WorkerDeploymentVersion
	RampPercentage        float32
	RampingRevisionNumber int64
}

RoutingInfoCache is used to cache results of GetTaskQueueUserData calls followed by CalculateTaskQueueVersioningInfo computation.

Implementations are expected to be safe for concurrent use.

type RoutingInfoCache added in v1.30.2

type RoutingInfoCache interface {
	// Get returns the cached routing info. ok=false means there was no cached value.
	Get(
		namespaceID string,
		taskQueue string,
		taskQueueType enumspb.TaskQueueType,
	) (RoutingInfo, bool)

	Put(
		namespaceID string,
		taskQueue string,
		taskQueueType enumspb.TaskQueueType,
		current *deploymentspb.WorkerDeploymentVersion,
		currentRevisionNumber int64,
		ramping *deploymentspb.WorkerDeploymentVersion,
		rampPercentage float32,
		rampingRevisionNumber int64,
	)
}

RoutingInfoCache is used to cache results of GetTaskQueueUserData calls followed by CalculateTaskQueueVersioningInfo computation.

Implementations are expected to be safe for concurrent use.

func NewRoutingInfoCache added in v1.30.2

func NewRoutingInfoCache(c cache.Cache, metricsHandler metrics.Handler) RoutingInfoCache

NewRoutingInfoCache wraps the provided cache with a typed API and metrics.

type RoutingInfoCacheImpl added in v1.30.2

type RoutingInfoCacheImpl struct {
	cache.Cache
	// contains filtered or unexported fields
}

RoutingInfoCache is used to cache results of GetTaskQueueUserData calls followed by CalculateTaskQueueVersioningInfo computation.

Implementations are expected to be safe for concurrent use.

func (*RoutingInfoCacheImpl) Get added in v1.30.2

func (c *RoutingInfoCacheImpl) Get(
	namespaceID string,
	taskQueue string,
	taskQueueType enumspb.TaskQueueType,
) (RoutingInfo, bool)

func (*RoutingInfoCacheImpl) Put added in v1.30.2

func (c *RoutingInfoCacheImpl) Put(
	namespaceID string,
	taskQueue string,
	taskQueueType enumspb.TaskQueueType,
	current *deploymentspb.WorkerDeploymentVersion,
	currentRevisionNumber int64,
	ramping *deploymentspb.WorkerDeploymentVersion,
	rampPercentage float32,
	rampingRevisionNumber int64,
)

type VersionMembershipCache added in v1.30.0

type VersionMembershipCache interface {
	// Get returns (isMember, ok). ok=false means there was no cached value.
	Get(
		namespaceID string,
		taskQueue string,
		taskQueueType enumspb.TaskQueueType,
		deploymentName string,
		buildID string,
	) (isMember bool, ok bool)

	Put(
		namespaceID string,
		taskQueue string,
		taskQueueType enumspb.TaskQueueType,
		deploymentName string,
		buildID string,
		isMember bool,
	)
}

VersionMembershipCache is used to cache results of Matching's CheckTaskQueueVersionMembership calls (used internally by the worker versioning pinned override validation).

Implementations are expected to be safe for concurrent use.

func NewVersionMembershipCache added in v1.30.0

func NewVersionMembershipCache(c cache.Cache, metricsHandler metrics.Handler) VersionMembershipCache

NewVersionMembershipCache wraps the provided cache with a typed API and metrics.

type VersionMembershipCacheImpl added in v1.30.0

type VersionMembershipCacheImpl struct {
	cache.Cache
	// contains filtered or unexported fields
}

VersionMembershipCache is used to cache results of Matching's CheckTaskQueueVersionMembership calls (used internally by the worker versioning pinned override validation).

Implementations are expected to be safe for concurrent use.

func (*VersionMembershipCacheImpl) Get added in v1.30.0

func (c *VersionMembershipCacheImpl) Get(
	namespaceID string,
	taskQueue string,
	taskQueueType enumspb.TaskQueueType,
	deploymentName string,
	buildID string,
) (isMember bool, ok bool)

func (*VersionMembershipCacheImpl) Put added in v1.30.0

func (c *VersionMembershipCacheImpl) Put(
	namespaceID string,
	taskQueue string,
	taskQueueType enumspb.TaskQueueType,
	deploymentName string,
	buildID string,
	isMember bool,
)

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL