Documentation
¶
Index ¶
- Constants
- func AddV31VersioningInfoToV32(info *workflowpb.WorkflowExecutionVersioningInfo) *workflowpb.WorkflowExecutionVersioningInfo
- func AssignedBuildIdSearchAttribute(buildId string) string
- func BuildIDToStringV32(deploymentName, buildID string) string
- func BuildIdFromCapabilities(capabilities *commonpb.WorkerVersionCapabilities, ...) string
- func BuildIdIfUsingVersioning(stamp *commonpb.WorkerVersionStamp) string
- func CalculateTaskQueueVersioningInfo(deployments *persistencespb.DeploymentData) (*deploymentspb.WorkerDeploymentVersion, int64, time.Time, ...)
- func CleanupOldDeletedVersions(deploymentData *persistencespb.WorkerDeploymentData, maxVersions int) bool
- func ConvertOverrideToV32(override *workflowpb.VersioningOverride) *workflowpb.VersioningOverride
- func CountDeploymentVersions(deployments *persistencespb.DeploymentData) int
- func DeploymentFromCapabilities(capabilities *commonpb.WorkerVersionCapabilities, ...) (*deploymentpb.Deployment, error)
- func DeploymentFromDeploymentVersion(dv *deploymentspb.WorkerDeploymentVersion) *deploymentpb.Deployment
- func DeploymentFromExternalDeploymentVersion(dv *deploymentpb.WorkerDeploymentVersion) *deploymentpb.Deployment
- func DeploymentIfValid(d *deploymentpb.Deployment) *deploymentpb.Deployment
- func DeploymentNameFromCapabilities(capabilities *commonpb.WorkerVersionCapabilities, ...) string
- func DeploymentOrVersion(d *deploymentpb.Deployment, v *deploymentspb.WorkerDeploymentVersion) *deploymentpb.Deployment
- func DeploymentVersionFromDeployment(deployment *deploymentpb.Deployment) *deploymentspb.WorkerDeploymentVersion
- func DeploymentVersionFromOptions(options *deploymentpb.WorkerDeploymentOptions) *deploymentspb.WorkerDeploymentVersion
- func DirectiveDeployment(directive *taskqueuespb.TaskVersionDirective) *deploymentpb.Deployment
- func ExternalWorkerDeploymentVersionFromDeployment(deployment *deploymentpb.Deployment) *deploymentpb.WorkerDeploymentVersion
- func ExternalWorkerDeploymentVersionFromStringV31(s string) *deploymentpb.WorkerDeploymentVersion
- func ExternalWorkerDeploymentVersionFromVersion(version *deploymentspb.WorkerDeploymentVersion) *deploymentpb.WorkerDeploymentVersion
- func ExternalWorkerDeploymentVersionToString(v *deploymentpb.WorkerDeploymentVersion) string
- func ExternalWorkerDeploymentVersionToStringV31(v *deploymentpb.WorkerDeploymentVersion) string
- func ExtractVersioningBehaviorFromOverride(override *workflowpb.VersioningOverride) enumspb.VersioningBehavior
- func FindBuildId(versioningData *persistencespb.VersioningData, buildId string) (setIndex, indexInSet int)
- func FindOldDeploymentVersion(deployments *persistencespb.DeploymentData, ...) int
- func FindTargetDeploymentVersionAndRevisionNumberForWorkflowID(current *deploymentspb.WorkerDeploymentVersion, currentRevisionNumber int64, ...) (*deploymentspb.WorkerDeploymentVersion, int64)
- func FormatPinnedVersionNotInTaskQueueError(deploymentName, buildID, taskQueue string, taskQueueType enumspb.TaskQueueType) string
- func GetOverridePinnedVersion(override *workflowpb.VersioningOverride) *deploymentpb.WorkerDeploymentVersion
- func HasDeploymentVersion(deployments *persistencespb.DeploymentData, ...) bool
- func IsUnversionedOrAssignedBuildIdSearchAttribute(buildId string) bool
- func MakeBuildIdDirective(buildId string) *taskqueuespb.TaskVersionDirective
- func MakeDirectiveForWorkflowTask(inheritedBuildId string, assignedBuildId string, ...) *taskqueuespb.TaskVersionDirective
- func MakeUseAssignmentRulesDirective() *taskqueuespb.TaskVersionDirective
- func OverrideIsPinned(override *workflowpb.VersioningOverride) bool
- func PickFinalCurrentAndRamping(current *deploymentspb.DeploymentVersionData, ...) (finalCurrent *deploymentspb.WorkerDeploymentVersion, finalCurrentRev int64, ...)
- func PinnedBuildIdSearchAttribute(version string) string
- func StampFromBuildId(buildId string) *commonpb.WorkerVersionStamp
- func StampFromCapabilities(capabilities *commonpb.WorkerVersionCapabilities, ...) *commonpb.WorkerVersionStamp
- func StampIfUsingVersioning(stamp *commonpb.WorkerVersionStamp) *commonpb.WorkerVersionStamp
- func UnversionedBuildIdSearchAttribute(buildId string) string
- func ValidateDeployment(deployment *deploymentpb.Deployment) error
- func ValidateDeploymentVersion(version *deploymentspb.WorkerDeploymentVersion, maxIDLengthLimit int) error
- func ValidateDeploymentVersionFields(fieldName string, field string, maxIDLengthLimit int) error
- func ValidateDeploymentVersionStringV31(version string) (*deploymentspb.WorkerDeploymentVersion, error)
- func ValidateTaskVersionDirective(directive *taskqueuespb.TaskVersionDirective, ...) error
- func ValidateVersioningOverride(ctx context.Context, override *workflowpb.VersioningOverride, ...) error
- func VersionStampToBuildIdSearchAttribute(stamp *commonpb.WorkerVersionStamp) string
- func VersionedBuildIdSearchAttribute(buildId string) string
- func WorkerDeploymentVersionFromStringV31(s string) (*deploymentspb.WorkerDeploymentVersion, error)
- func WorkerDeploymentVersionFromStringV32(s string) (*deploymentspb.WorkerDeploymentVersion, error)
- func WorkerDeploymentVersionToStringV31(v *deploymentspb.WorkerDeploymentVersion) string
- func WorkerDeploymentVersionToStringV32(v *deploymentspb.WorkerDeploymentVersion) string
- func WorkflowsExistForBuildId(ctx context.Context, visibilityManager manager.VisibilityManager, ...) (bool, error)
- type IsWFTaskQueueInVersionDetector
- type ReactivationSignalCache
- type ReactivationSignalCacheImpl
- type RoutingInfo
- type RoutingInfoCache
- type RoutingInfoCacheImpl
- type VersionMembershipCache
- type VersionMembershipCacheImpl
Constants ¶
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
func AddV31VersioningInfoToV32(info *workflowpb.WorkflowExecutionVersioningInfo) *workflowpb.WorkflowExecutionVersioningInfo
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
AssignedBuildIdSearchAttribute returns the search attribute value for the currently assigned build ID
func BuildIDToStringV32 ¶ added in v1.30.0
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
func CalculateTaskQueueVersioningInfo(deployments *persistencespb.DeploymentData) ( *deploymentspb.WorkerDeploymentVersion, int64, time.Time, *deploymentspb.WorkerDeploymentVersion, bool, float32, int64, time.Time, )
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
func DeploymentOrVersion(d *deploymentpb.Deployment, v *deploymentspb.WorkerDeploymentVersion) *deploymentpb.Deployment
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
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
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 StampFromCapabilities(capabilities *commonpb.WorkerVersionCapabilities, options *deploymentpb.WorkerDeploymentOptions) *commonpb.WorkerVersionStamp
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 ¶
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
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 ¶
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
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
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
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
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, )