go.temporal.io/server/service/matching/config.go
608 LOC · 293 covered · 315 uncovered · 54 ranges · 1653 concepts · 47 introducers · 677 tests
File neighbourhood
The centred file is linked to every concept that introduces one of its ranges, every test that runs code from the file, and the gray connector concepts standing between those tests and the file's own introducer concepts. Undirected links join concepts to every file where they introduce source and concepts to the tests they introduce; arrows show specialization between the displayed concepts and bridge only concepts omitted from this view. Concept colors match the source ranges below; connector concepts have no source color and are shown in gray.
Focused file, its introducer and connector concepts, their introduced files, and tests that run code from the file
In the embedded map, ordinary wheel input scrolls the page; use the visible controls to zoom and drag to pan. Open the full-screen map for canvas navigation: wheel pans, Ctrl/Command plus wheel zooms, and arrow keys pan when this region is focused. On touch screens, open the full-screen map to pan or pinch. If JavaScript or WebGL is unavailable, use the related-file, concept, and source links on this page.
Graph controls are ready.
Interactive rendering requires JavaScript and WebGL. Use the related-file, concept, and source links on this page while the interactive map is unavailable.
//go:generate stringer -type loadCause -trimprefix loadCause -output loadcause_string_gen.go
//go:generate stringer -type unloadCause -trimprefix unloadCause -output unloadcause_string_gen.go
package matching
import (
"time"
"go.temporal.io/server/common/backoff"
"go.temporal.io/server/common/dynamicconfig"
"go.temporal.io/server/common/namespace"
"go.temporal.io/server/common/tqid"
"go.temporal.io/server/components/nexusoperations"
"go.temporal.io/server/service/matching/counter"
)
type (
// Config represents configuration for matching service
Config struct {
PersistenceMaxQPS dynamicconfig.IntPropertyFn
PersistenceGlobalMaxQPS dynamicconfig.IntPropertyFn
PersistenceNamespaceMaxQPS dynamicconfig.IntPropertyFnWithNamespaceFilter
PersistenceGlobalNamespaceMaxQPS dynamicconfig.IntPropertyFnWithNamespaceFilter
PersistencePerShardNamespaceMaxQPS dynamicconfig.IntPropertyFnWithNamespaceFilter
PersistenceDynamicRateLimitingParams dynamicconfig.TypedPropertyFn[dynamicconfig.DynamicRateLimitingParams]
PersistenceQPSBurstRatio dynamicconfig.FloatPropertyFn
SyncMatchWaitDuration dynamicconfig.DurationPropertyFnWithTaskQueueFilter
RPS dynamicconfig.IntPropertyFn
NamespaceRPS dynamicconfig.IntPropertyFnWithNamespaceFilter
OperatorRPSRatio dynamicconfig.FloatPropertyFn
PollWaitForNamespaceRateLimitToken dynamicconfig.BoolPropertyFnWithNamespaceFilter
AlignMembershipChange dynamicconfig.DurationPropertyFn
ShutdownDrainDuration dynamicconfig.DurationPropertyFn
HistoryMaxPageSize dynamicconfig.IntPropertyFnWithNamespaceFilter
EnableDeployments dynamicconfig.BoolPropertyFnWithNamespaceFilter // [cleanup-wv-pre-release]
EnableDeploymentVersions dynamicconfig.BoolPropertyFnWithNamespaceFilter
UseRevisionNumberForWorkerVersioning dynamicconfig.BoolPropertyFnWithNamespaceFilter
MaxTaskQueuesInDeployment dynamicconfig.IntPropertyFnWithNamespaceFilter
MaxVersionsInTaskQueue dynamicconfig.IntPropertyFnWithNamespaceFilter
MaxIDLengthLimit dynamicconfig.IntPropertyFn
// task queue configuration
RangeSize int64
NewMatcherSub dynamicconfig.TypedSubscribableWithTaskQueueFilter[dynamicconfig.GradualChange[bool]]
EnableFairnessSub dynamicconfig.TypedSubscribableWithTaskQueueFilter[dynamicconfig.GradualChange[bool]]
EnableMigration dynamicconfig.BoolPropertyFnWithTaskQueueFilter
AutoEnableV2Sub dynamicconfig.TypedSubscribableWithTaskQueueFilter[bool]
GetTasksBatchSize dynamicconfig.IntPropertyFnWithTaskQueueFilter
GetTasksReloadAt dynamicconfig.IntPropertyFnWithTaskQueueFilter
ForceReadTasksOnWrite dynamicconfig.BoolPropertyFnWithTaskQueueFilter
UpdateAckInterval dynamicconfig.DurationPropertyFnWithTaskQueueFilter
MetadataUpdateOnAppendInterval dynamicconfig.DurationPropertyFnWithTaskQueueFilter
MaxTaskQueueIdleTime dynamicconfig.DurationPropertyFnWithTaskQueueFilter
NumTaskqueueWritePartitions dynamicconfig.IntPropertyFnWithTaskQueueFilter
NumTaskqueueReadPartitions dynamicconfig.IntPropertyFnWithTaskQueueFilter
NumTaskqueueReadPartitionsSub dynamicconfig.TypedSubscribableWithTaskQueueFilter[int]
BreakdownMetricsByTaskQueue dynamicconfig.BoolPropertyFnWithTaskQueueFilter
BreakdownMetricsByPartition dynamicconfig.BoolPropertyFnWithTaskQueueFilter
BreakdownMetricsByBuildID dynamicconfig.BoolPropertyFnWithTaskQueueFilter
EnableWorkerPluginMetrics dynamicconfig.BoolPropertyFn
EnablePollerAutoscalingMetrics dynamicconfig.BoolPropertyFn
ExternalPayloadsEnabled dynamicconfig.BoolPropertyFnWithNamespaceFilter
WorkerRegistryNumBuckets dynamicconfig.IntPropertyFn
WorkerRegistryEntryTTL dynamicconfig.DurationPropertyFn
WorkerRegistryMinEvictAge dynamicconfig.DurationPropertyFn
WorkerRegistryMaxEntries dynamicconfig.IntPropertyFn
WorkerRegistryEvictionInterval dynamicconfig.DurationPropertyFn
ForwarderMaxOutstandingPolls dynamicconfig.IntPropertyFnWithTaskQueueFilter
ForwarderMaxOutstandingTasks dynamicconfig.IntPropertyFnWithTaskQueueFilter
ForwarderMaxRatePerSecond dynamicconfig.FloatPropertyFnWithTaskQueueFilter
ForwarderMaxChildrenPerNode dynamicconfig.IntPropertyFnWithTaskQueueFilter
VersionCompatibleSetLimitPerQueue dynamicconfig.IntPropertyFnWithNamespaceFilter
VersionBuildIdLimitPerQueue dynamicconfig.IntPropertyFnWithNamespaceFilter
AssignmentRuleLimitPerQueue dynamicconfig.IntPropertyFnWithNamespaceFilter
RedirectRuleLimitPerQueue dynamicconfig.IntPropertyFnWithNamespaceFilter
RedirectRuleMaxUpstreamBuildIDsPerQueue dynamicconfig.IntPropertyFnWithNamespaceFilter
DeletedRuleRetentionTime dynamicconfig.DurationPropertyFnWithNamespaceFilter
PollerHistoryTTL dynamicconfig.DurationPropertyFnWithNamespaceFilter
EnableMatchingFanOutForPollCancellation dynamicconfig.BoolPropertyFnWithNamespaceFilter
ReachabilityBuildIdVisibilityGracePeriod dynamicconfig.DurationPropertyFnWithNamespaceFilter
ReachabilityCacheOpenWFsTTL dynamicconfig.DurationPropertyFn
ReachabilityCacheClosedWFsTTL dynamicconfig.DurationPropertyFn
TaskQueueLimitPerBuildId dynamicconfig.IntPropertyFnWithNamespaceFilter
GetUserDataLongPollTimeout dynamicconfig.DurationPropertyFn
GetUserDataRefresh dynamicconfig.DurationPropertyFn
EphemeralDataUpdateInterval dynamicconfig.DurationPropertyFnWithTaskQueueFilter
BacklogMetricsEmitInterval dynamicconfig.DurationPropertyFnWithTaskQueueFilter
PriorityBacklogForwarding dynamicconfig.BoolPropertyFnWithTaskQueueFilter
BacklogNegligibleAge dynamicconfig.DurationPropertyFnWithTaskQueueFilter
MaxWaitForPollerBeforeFwd dynamicconfig.DurationPropertyFnWithTaskQueueFilter
QueryPollerUnavailableWindow dynamicconfig.DurationPropertyFn
WorkerControllerNoPollerHookWindow dynamicconfig.DurationPropertyFn
EmitTaskDispatchLatencyAtPoll dynamicconfig.BoolPropertyFnWithTaskQueueFilter
QueryWorkflowTaskTimeoutLogRate dynamicconfig.FloatPropertyFnWithTaskQueueFilter
MembershipUnloadDelay dynamicconfig.DurationPropertyFn
TaskQueueInfoByBuildIdTTL dynamicconfig.DurationPropertyFnWithTaskQueueFilter
PriorityLevels dynamicconfig.IntPropertyFnWithTaskQueueFilter
RateLimitFractionProvider TaskQueueRateLimitFractionProvider
RateLimiterRefreshInterval time.Duration
FairnessKeyRateLimitCacheSize dynamicconfig.IntPropertyFnWithTaskQueueFilter
MaxFairnessKeyWeightOverrides dynamicconfig.IntPropertyFnWithTaskQueueFilter
// Time to hold a poll request before returning an empty response if there are no tasks
LongPollExpirationInterval dynamicconfig.DurationPropertyFnWithTaskQueueFilter
BacklogTaskForwardTimeout dynamicconfig.DurationPropertyFnWithTaskQueueFilter
ForwardPollRetryMaxInterval dynamicconfig.DurationPropertyFnWithTaskQueueFilter
MinTaskThrottlingBurstSize dynamicconfig.IntPropertyFnWithTaskQueueFilter
MaxTaskDeleteBatchSize dynamicconfig.IntPropertyFnWithTaskQueueFilter
TaskDeleteInterval dynamicconfig.DurationPropertyFnWithTaskQueueFilter
// taskWriter configuration
OutstandingTaskAppendsThreshold dynamicconfig.IntPropertyFnWithTaskQueueFilter
MaxTaskBatchSize dynamicconfig.IntPropertyFnWithTaskQueueFilter
ThrottledLogRPS dynamicconfig.IntPropertyFn
AdminNamespaceToPartitionDispatchRate dynamicconfig.FloatPropertyFnWithNamespaceFilter
AdminNamespaceToPartitionRateSub dynamicconfig.TypedSubscribableWithNamespaceFilter[float64]
AdminNamespaceTaskqueueToPartitionDispatchRate dynamicconfig.FloatPropertyFnWithTaskQueueFilter
AdminNamespaceTaskqueueToPartitionRateSub dynamicconfig.TypedSubscribableWithTaskQueueFilter[float64]
VisibilityPersistenceMaxReadQPS dynamicconfig.IntPropertyFn
VisibilityPersistenceMaxWriteQPS dynamicconfig.IntPropertyFn
VisibilityPersistenceSlowQueryThreshold dynamicconfig.DurationPropertyFn
EnableReadFromSecondaryVisibility dynamicconfig.BoolPropertyFnWithNamespaceFilter
VisibilityEnableShadowReadMode dynamicconfig.BoolPropertyFn
VisibilityDisableOrderByClause dynamicconfig.BoolPropertyFnWithNamespaceFilter
VisibilityEnableManualPagination dynamicconfig.BoolPropertyFnWithNamespaceFilter
VisibilityEnableUnifiedQueryConverter dynamicconfig.BoolPropertyFn
ListNexusEndpointsLongPollTimeout dynamicconfig.DurationPropertyFn
NexusEndpointsRefreshInterval dynamicconfig.DurationPropertyFn
MinDispatchTaskTimeout dynamicconfig.DurationPropertyFnWithNamespaceFilter
PollerScalingBacklogAgeScaleUp dynamicconfig.DurationPropertyFnWithTaskQueueFilter
PollerScalingWaitTime dynamicconfig.DurationPropertyFnWithTaskQueueFilter
PollerScalingDecisionsPerSecond dynamicconfig.FloatPropertyFnWithTaskQueueFilter
PollerScalingTaskAddToDispatchRatio dynamicconfig.FloatPropertyFnWithTaskQueueFilter
EnablePollerScalingDecisionMetrics dynamicconfig.BoolPropertyFnWithTaskQueueFilter
FairnessCounter dynamicconfig.TypedPropertyFnWithTaskQueueFilter[counter.CounterParams]
FairnessPassDither dynamicconfig.BoolPropertyFnWithTaskQueueFilter
PartitionScaleAllowedDrift dynamicconfig.TypedPropertyFnWithTaskQueueFilter[dynamicconfig.PartitionScaleAllowedDrift]
PartitionScaleManagerSettings dynamicconfig.TypedPropertyFnWithTaskQueueFilter[dynamicconfig.PartitionScaleManagerSettings]
LogAllReqErrors dynamicconfig.BoolPropertyFnWithNamespaceFilter
}
forwarderConfig struct {
ForwarderMaxOutstandingPolls func() int
ForwarderMaxOutstandingTasks func() int
ForwarderMaxRatePerSecond func() float64
ForwarderMaxChildrenPerNode func() int
}
taskQueueConfig struct {
forwarderConfig
SyncMatchWaitDuration func() time.Duration
EphemeralDataUpdateInterval func() time.Duration
BacklogMetricsEmitInterval func() time.Duration
PriorityBacklogForwarding func() bool
BacklogNegligibleAge func() time.Duration
MaxWaitForPollerBeforeFwd func() time.Duration
QueryPollerUnavailableWindow func() time.Duration
WorkerControllerNoPollerHookWindow func() time.Duration
EmitTaskDispatchLatencyAtPoll func() bool
// Time to hold a poll request before returning an empty response if there are no tasks
LongPollExpirationInterval func() time.Duration
BacklogTaskForwardTimeout func() time.Duration
ForwardPollRetryMaxInterval func() time.Duration
RangeSize int64
NewMatcher bool
NewMatcherSub func(func(dynamicconfig.GradualChange[bool])) (dynamicconfig.GradualChange[bool], func())
EnableFairness bool
EnableFairnessSub func(func(dynamicconfig.GradualChange[bool])) (dynamicconfig.GradualChange[bool], func())
EnableMigration func() bool
AutoEnableV2 func() bool
AutoEnableV2Sub func(func(bool)) (bool, func())
GetTasksBatchSize func() int
GetTasksReloadAt func() int
ForceReadTasksOnWrite func() bool
UpdateAckInterval func() time.Duration
MetadataUpdateOnAppendInterval func() time.Duration
MaxTaskQueueIdleTime func() time.Duration
MinTaskThrottlingBurstSize func() int
MaxTaskDeleteBatchSize func() int
TaskDeleteInterval func() time.Duration
PriorityLevels priorityKey
DefaultPriorityKey priorityKey
GetUserDataLongPollTimeout dynamicconfig.DurationPropertyFn
GetUserDataMinWaitTime time.Duration
GetUserDataReturnBudget time.Duration
GetUserDataInitialRefresh time.Duration
GetUserDataRefresh dynamicconfig.DurationPropertyFn
// taskWriter configuration
OutstandingTaskAppendsThreshold func() int
MaxTaskBatchSize func() int
NumWritePartitions func() int
NumReadPartitions func() int
NumReadPartitionsSub func(func(int)) (int, func())
// partition qps = AdminNamespaceToPartitionDispatchRate(namespace)
AdminNamespaceToPartitionDispatchRate func() float64
AdminNamespaceToPartitionRateSub func(func(float64)) (float64, func())
// partition qps = AdminNamespaceTaskQueueToPartitionDispatchRate(namespace, task_queue)
AdminNamespaceTaskQueueToPartitionDispatchRate func() float64
AdminNamespaceTaskQueueToPartitionRateSub func(func(float64)) (float64, func())
// Retry policy for fetching user data from root partition. Should retry forever.
GetUserDataRetryPolicy backoff.RetryPolicy
// TTL for cache holding TaskQueueInfoByBuildID
TaskQueueInfoByBuildIdTTL func() time.Duration
MaxVersionsInTaskQueue func() int
// Rate limiting
RateLimitFraction func() float64
RateLimiterRefreshInterval time.Duration
FairnessKeyRateLimitCacheSize func() int
MaxFairnessKeyWeightOverrides func() int
BreakdownMetricsByTaskQueue func() bool
BreakdownMetricsByPartition func() bool
BreakdownMetricsByBuildID func() bool
PollerHistoryTTL func() time.Duration
// Poller scaling decisions configuration
PollerScalingBacklogAgeScaleUp func() time.Duration
PollerScalingWaitTime func() time.Duration
PollerScalingDecisionsPerSecond func() float64
PollerScalingTaskAddToDispatchRatio func() float64
EnablePollerScalingDecisionMetrics func() bool
FairnessCounter func() counter.CounterParams
FairnessPassDither func() bool
PartitionScaleAllowedDrift func() dynamicconfig.PartitionScaleAllowedDrift
PartitionScaleManagerSettings func() dynamicconfig.PartitionScaleManagerSettings
loadCause loadCause
}
loadCause int
unloadCause int
)
const (
loadCauseUnspecified loadCause = iota
loadCauseTask
loadCauseQuery
loadCauseDescribe
loadCauseUserData
loadCauseNexusTask
loadCausePoll
loadCauseOtherRead // any other read-only rpc
loadCauseOtherWrite // any other mutating rpc
loadCauseForce // root partition loaded, force load to ensure matching with back logged partitions
)
const (
unloadCauseUnspecified unloadCause = iota
unloadCauseInitError
unloadCauseIdle
unloadCauseMembership // proactive unload due to ownership change
unloadCauseConflict // reactive unload due to other node stealing ownership
unloadCauseShuttingDown
unloadCauseForce
unloadCauseConfigChange
unloadCauseOtherError
)
// NewConfig returns new service config with default values
func NewConfig(
dc *dynamicconfig.Collection,
return &Config{
PersistenceMaxQPS: dynamicconfig.MatchingPersistenceMaxQPS.Get(dc),
PersistenceGlobalMaxQPS: dynamicconfig.MatchingPersistenceGlobalMaxQPS.Get(dc),
PersistenceNamespaceMaxQPS: dynamicconfig.MatchingPersistenceNamespaceMaxQPS.Get(dc),
PersistenceGlobalNamespaceMaxQPS: dynamicconfig.MatchingPersistenceGlobalNamespaceMaxQPS.Get(dc),
PersistencePerShardNamespaceMaxQPS: dynamicconfig.DefaultPerShardNamespaceRPSMax,
PersistenceDynamicRateLimitingParams: dynamicconfig.MatchingPersistenceDynamicRateLimitingParams.Get(dc),
PersistenceQPSBurstRatio: dynamicconfig.PersistenceQPSBurstRatio.Get(dc),
SyncMatchWaitDuration: dynamicconfig.MatchingSyncMatchWaitDuration.Get(dc),
HistoryMaxPageSize: dynamicconfig.MatchingHistoryMaxPageSize.Get(dc),
EnableDeployments: dynamicconfig.EnableDeployments.Get(dc), // [cleanup-wv-pre-release]
EnableDeploymentVersions: dynamicconfig.EnableDeploymentVersions.Get(dc),
UseRevisionNumberForWorkerVersioning: dynamicconfig.UseRevisionNumberForWorkerVersioning.Get(dc),
MaxTaskQueuesInDeployment: dynamicconfig.MatchingMaxTaskQueuesInDeployment.Get(dc),
MaxVersionsInTaskQueue: dynamicconfig.MatchingMaxVersionsInTaskQueue.Get(dc),
RPS: dynamicconfig.MatchingRPS.Get(dc),
NamespaceRPS: dynamicconfig.MatchingNamespaceRPS.Get(dc),
OperatorRPSRatio: dynamicconfig.OperatorRPSRatio.Get(dc),
PollWaitForNamespaceRateLimitToken: dynamicconfig.PollWaitForNamespaceRateLimitToken.Get(dc),
RangeSize: 100000,
NewMatcherSub: dynamicconfig.MatchingUseNewMatcher.Subscribe(dc),
EnableFairnessSub: dynamicconfig.MatchingEnableFairness.Subscribe(dc),
EnableMigration: dynamicconfig.MatchingEnableMigration.Get(dc),
AutoEnableV2Sub: dynamicconfig.MatchingAutoEnableV2.Subscribe(dc),
GetTasksBatchSize: dynamicconfig.MatchingGetTasksBatchSize.Get(dc),
GetTasksReloadAt: dynamicconfig.MatchingGetTasksReloadAt.Get(dc),
ForceReadTasksOnWrite: dynamicconfig.MatchingForceReadTasksOnWrite.Get(dc),
UpdateAckInterval: dynamicconfig.MatchingUpdateAckInterval.Get(dc),
MetadataUpdateOnAppendInterval: dynamicconfig.MatchingMetadataUpdateOnAppendInterval.Get(dc),
MaxTaskQueueIdleTime: dynamicconfig.MatchingMaxTaskQueueIdleTime.Get(dc),
LongPollExpirationInterval: dynamicconfig.MatchingLongPollExpirationInterval.Get(dc),
BacklogTaskForwardTimeout: dynamicconfig.MatchingBacklogTaskForwardTimeout.Get(dc),
ForwardPollRetryMaxInterval: dynamicconfig.MatchingForwardPollRetryMaxInterval.Get(dc),
MinTaskThrottlingBurstSize: dynamicconfig.MatchingMinTaskThrottlingBurstSize.Get(dc),
MaxTaskDeleteBatchSize: dynamicconfig.MatchingMaxTaskDeleteBatchSize.Get(dc),
TaskDeleteInterval: dynamicconfig.MatchingTaskDeleteInterval.Get(dc),
OutstandingTaskAppendsThreshold: dynamicconfig.MatchingOutstandingTaskAppendsThreshold.Get(dc),
MaxTaskBatchSize: dynamicconfig.MatchingMaxTaskBatchSize.Get(dc),
ThrottledLogRPS: dynamicconfig.MatchingThrottledLogRPS.Get(dc),
NumTaskqueueWritePartitions: dynamicconfig.MatchingNumTaskqueueWritePartitions.Get(dc),
NumTaskqueueReadPartitions: dynamicconfig.MatchingNumTaskqueueReadPartitions.Get(dc),
NumTaskqueueReadPartitionsSub: dynamicconfig.MatchingNumTaskqueueReadPartitions.Subscribe(dc),
BreakdownMetricsByTaskQueue: dynamicconfig.MetricsBreakdownByTaskQueue.Get(dc),
BreakdownMetricsByPartition: dynamicconfig.MetricsBreakdownByPartition.Get(dc),
BreakdownMetricsByBuildID: dynamicconfig.MetricsBreakdownByBuildID.Get(dc),
EnableWorkerPluginMetrics: dynamicconfig.MatchingEnableWorkerPluginMetrics.Get(dc),
EnablePollerAutoscalingMetrics: dynamicconfig.MatchingEnablePollerAutoscalingMetrics.Get(dc),
ExternalPayloadsEnabled: dynamicconfig.ExternalPayloadsEnabled.Get(dc),
WorkerRegistryNumBuckets: dynamicconfig.MatchingWorkerRegistryNumBuckets.Get(dc),
WorkerRegistryEntryTTL: dynamicconfig.MatchingWorkerRegistryEntryTTL.Get(dc),
WorkerRegistryMinEvictAge: dynamicconfig.MatchingWorkerRegistryMinEvictAge.Get(dc),
WorkerRegistryMaxEntries: dynamicconfig.MatchingWorkerRegistryMaxEntries.Get(dc),
WorkerRegistryEvictionInterval: dynamicconfig.MatchingWorkerRegistryEvictionInterval.Get(dc),
ForwarderMaxOutstandingPolls: dynamicconfig.MatchingForwarderMaxOutstandingPolls.Get(dc),
ForwarderMaxOutstandingTasks: dynamicconfig.MatchingForwarderMaxOutstandingTasks.Get(dc),
ForwarderMaxRatePerSecond: dynamicconfig.MatchingForwarderMaxRatePerSecond.Get(dc),
ForwarderMaxChildrenPerNode: dynamicconfig.MatchingForwarderMaxChildrenPerNode.Get(dc),
AlignMembershipChange: dynamicconfig.MatchingAlignMembershipChange.Get(dc),
ShutdownDrainDuration: dynamicconfig.MatchingShutdownDrainDuration.Get(dc),
VersionCompatibleSetLimitPerQueue: dynamicconfig.VersionCompatibleSetLimitPerQueue.Get(dc),
VersionBuildIdLimitPerQueue: dynamicconfig.VersionBuildIdLimitPerQueue.Get(dc),
AssignmentRuleLimitPerQueue: dynamicconfig.AssignmentRuleLimitPerQueue.Get(dc),
RedirectRuleLimitPerQueue: dynamicconfig.RedirectRuleLimitPerQueue.Get(dc),
RedirectRuleMaxUpstreamBuildIDsPerQueue: dynamicconfig.RedirectRuleMaxUpstreamBuildIDsPerQueue.Get(dc),
DeletedRuleRetentionTime: dynamicconfig.MatchingDeletedRuleRetentionTime.Get(dc),
PollerHistoryTTL: dynamicconfig.PollerHistoryTTL.Get(dc),
EnableMatchingFanOutForPollCancellation: dynamicconfig.EnableMatchingFanOutForPollCancellation.Get(dc),
ReachabilityBuildIdVisibilityGracePeriod: dynamicconfig.ReachabilityBuildIdVisibilityGracePeriod.Get(dc),
ReachabilityCacheOpenWFsTTL: dynamicconfig.ReachabilityCacheOpenWFsTTL.Get(dc),
ReachabilityCacheClosedWFsTTL: dynamicconfig.ReachabilityCacheClosedWFsTTL.Get(dc),
TaskQueueLimitPerBuildId: dynamicconfig.TaskQueuesPerBuildIdLimit.Get(dc),
GetUserDataLongPollTimeout: dynamicconfig.MatchingGetUserDataLongPollTimeout.Get(dc), // Use -10 seconds so that we send back empty response instead of timeout
GetUserDataRefresh: dynamicconfig.MatchingGetUserDataRefresh.Get(dc),
EphemeralDataUpdateInterval: dynamicconfig.MatchingEphemeralDataUpdateInterval.Get(dc),
BacklogMetricsEmitInterval: dynamicconfig.MatchingBacklogMetricsEmitInterval.Get(dc),
PriorityBacklogForwarding: dynamicconfig.MatchingPriorityBacklogForwarding.Get(dc),
BacklogNegligibleAge: dynamicconfig.MatchingBacklogNegligibleAge.Get(dc),
MaxWaitForPollerBeforeFwd: dynamicconfig.MatchingMaxWaitForPollerBeforeFwd.Get(dc),
QueryPollerUnavailableWindow: dynamicconfig.QueryPollerUnavailableWindow.Get(dc),
WorkerControllerNoPollerHookWindow: dynamicconfig.WorkerControllerNoPollerHookWindow.Get(dc),
EmitTaskDispatchLatencyAtPoll: dynamicconfig.MatchingEmitTaskDispatchLatencyAtPoll.Get(dc),
QueryWorkflowTaskTimeoutLogRate: dynamicconfig.MatchingQueryWorkflowTaskTimeoutLogRate.Get(dc),
MembershipUnloadDelay: dynamicconfig.MatchingMembershipUnloadDelay.Get(dc),
TaskQueueInfoByBuildIdTTL: dynamicconfig.TaskQueueInfoByBuildIdTTL.Get(dc),
PriorityLevels: dynamicconfig.MatchingPriorityLevels.Get(dc),
RateLimiterRefreshInterval: time.Minute,
FairnessKeyRateLimitCacheSize: dynamicconfig.MatchingFairnessKeyRateLimitCacheSize.Get(dc),
MaxFairnessKeyWeightOverrides: dynamicconfig.MatchingMaxFairnessKeyWeightOverrides.Get(dc),
MaxIDLengthLimit: dynamicconfig.MaxIDLengthLimit.Get(dc),
AdminNamespaceToPartitionDispatchRate: dynamicconfig.AdminMatchingNamespaceToPartitionDispatchRate.Get(dc),
AdminNamespaceToPartitionRateSub: dynamicconfig.AdminMatchingNamespaceToPartitionDispatchRate.Subscribe(dc),
AdminNamespaceTaskqueueToPartitionDispatchRate: dynamicconfig.AdminMatchingNamespaceTaskqueueToPartitionDispatchRate.Get(dc),
AdminNamespaceTaskqueueToPartitionRateSub: dynamicconfig.AdminMatchingNamespaceTaskqueueToPartitionDispatchRate.Subscribe(dc),
VisibilityPersistenceMaxReadQPS: dynamicconfig.VisibilityPersistenceMaxReadQPS.Get(dc),
VisibilityPersistenceMaxWriteQPS: dynamicconfig.VisibilityPersistenceMaxWriteQPS.Get(dc),
VisibilityPersistenceSlowQueryThreshold: dynamicconfig.VisibilityPersistenceSlowQueryThreshold.Get(dc),
EnableReadFromSecondaryVisibility: dynamicconfig.EnableReadFromSecondaryVisibility.Get(dc),
VisibilityEnableShadowReadMode: dynamicconfig.VisibilityEnableShadowReadMode.Get(dc),
VisibilityDisableOrderByClause: dynamicconfig.VisibilityDisableOrderByClause.Get(dc),
VisibilityEnableManualPagination: dynamicconfig.VisibilityEnableManualPagination.Get(dc),
VisibilityEnableUnifiedQueryConverter: dynamicconfig.VisibilityEnableUnifiedQueryConverter.Get(dc),
ListNexusEndpointsLongPollTimeout: dynamicconfig.MatchingListNexusEndpointsLongPollTimeout.Get(dc),
NexusEndpointsRefreshInterval: dynamicconfig.MatchingNexusEndpointsRefreshInterval.Get(dc),
MinDispatchTaskTimeout: nexusoperations.MinDispatchTaskTimeout.Get(dc),
PollerScalingBacklogAgeScaleUp: dynamicconfig.MatchingPollerScalingBacklogAgeScaleUp.Get(dc),
PollerScalingWaitTime: dynamicconfig.MatchingPollerScalingWaitTime.Get(dc),
PollerScalingDecisionsPerSecond: dynamicconfig.MatchingPollerScalingDecisionsPerSecond.Get(dc),
PollerScalingTaskAddToDispatchRatio: dynamicconfig.MatchingPollerScalingTaskAddToDispatchRatio.Get(dc),
EnablePollerScalingDecisionMetrics: dynamicconfig.MatchingEnablePollerScalingDecisionMetrics.Get(dc),
FairnessCounter: dynamicconfig.MatchingFairnessCounter.Get(dc),
FairnessPassDither: dynamicconfig.MatchingFairnessPassDither.Get(dc),
PartitionScaleAllowedDrift: dynamicconfig.MatchingPartitionScaleAllowedDrift.Get(dc),
PartitionScaleManagerSettings: dynamicconfig.MatchingPartitionScaleManager.Get(dc),
LogAllReqErrors: dynamicconfig.LogAllReqErrors.Get(dc),
RateLimitFractionProvider: defaultTaskQueueRateLimitFractionProvider,
}
}
func newTaskQueueConfig(tq *tqid.TaskQueue, config *Config, ns namespace.Name) *taskQueueConfig {
config.go ×1
taskQueueName := tq.Name()
taskType := tq.TaskType()
priorityLevels := priorityKey(config.PriorityLevels(ns.String(), taskQueueName, taskType))
priorityLevels = max(priorityLevels, min(priorityLevels, maxPriorityLevels), 1)
defaultPriorityKey := (priorityLevels + 1) / 2
return &taskQueueConfig{
RangeSize: config.RangeSize,
NewMatcherSub: func(cb func(dynamicconfig.GradualChange[bool])) (dynamicconfig.GradualChange[bool], func()) {
return config.NewMatcherSub(ns.String(), taskQueueName, taskType, cb)
task_queue_partition_manager.go ×14
},
EnableFairnessSub: func(cb func(dynamicconfig.GradualChange[bool])) (dynamicconfig.GradualChange[bool], func()) {
return config.EnableFairnessSub(ns.String(), taskQueueName, taskType, cb)
},
return config.EnableMigration(ns.String(), taskQueueName, taskType)
},
v, _ := config.AutoEnableV2Sub(ns.String(), taskQueueName, taskType, nil)
return v
},
return config.AutoEnableV2Sub(ns.String(), taskQueueName, taskType, cb)
},
return config.GetTasksBatchSize(ns.String(), taskQueueName, taskType)
},
return config.GetTasksReloadAt(ns.String(), taskQueueName, taskType)
},
return config.ForceReadTasksOnWrite(ns.String(), taskQueueName, taskType)
},
return config.UpdateAckInterval(ns.String(), taskQueueName, taskType)
},
return config.MetadataUpdateOnAppendInterval(ns.String(), taskQueueName, taskType)
},
return config.MaxTaskQueueIdleTime(ns.String(), taskQueueName, taskType)
},
return config.MinTaskThrottlingBurstSize(ns.String(), taskQueueName, taskType)
},
return config.SyncMatchWaitDuration(ns.String(), taskQueueName, taskType)
},
return config.EphemeralDataUpdateInterval(ns.String(), taskQueueName, taskType)
},
return config.BacklogMetricsEmitInterval(ns.String(), taskQueueName, taskType)
},
PriorityBacklogForwarding: func() bool {
return config.PriorityBacklogForwarding(ns.String(), taskQueueName, taskType)
},
return config.BacklogNegligibleAge(ns.String(), taskQueueName, taskType)
},
return config.MaxWaitForPollerBeforeFwd(ns.String(), taskQueueName, taskType)
},
QueryPollerUnavailableWindow: config.QueryPollerUnavailableWindow,
WorkerControllerNoPollerHookWindow: config.WorkerControllerNoPollerHookWindow,
return config.EmitTaskDispatchLatencyAtPoll(ns.String(), taskQueueName, taskType)
},
return config.LongPollExpirationInterval(ns.String(), taskQueueName, taskType)
},
return config.BacklogTaskForwardTimeout(ns.String(), taskQueueName, taskType)
},
return config.ForwardPollRetryMaxInterval(ns.String(), taskQueueName, taskType)
},
return config.MaxTaskDeleteBatchSize(ns.String(), taskQueueName, taskType)
},
return config.TaskDeleteInterval(ns.String(), taskQueueName, taskType)
},
PriorityLevels: priorityLevels,
DefaultPriorityKey: defaultPriorityKey,
GetUserDataLongPollTimeout: config.GetUserDataLongPollTimeout,
GetUserDataMinWaitTime: 1 * time.Second,
GetUserDataReturnBudget: returnEmptyTaskTimeBudget,
GetUserDataInitialRefresh: ioTimeout,
GetUserDataRefresh: config.GetUserDataRefresh,
return config.OutstandingTaskAppendsThreshold(ns.String(), taskQueueName, taskType)
},
return config.MaxTaskBatchSize(ns.String(), taskQueueName, taskType)
},
return max(1, config.NumTaskqueueWritePartitions(ns.String(), taskQueueName, taskType))
},
return max(1, config.NumTaskqueueReadPartitions(ns.String(), taskQueueName, taskType))
},
return config.NumTaskqueueReadPartitionsSub(ns.String(), taskQueueName, taskType, cb)
},
return config.BreakdownMetricsByTaskQueue(ns.String(), taskQueueName, taskType)
},
return config.BreakdownMetricsByPartition(ns.String(), taskQueueName, taskType)
},
return config.BreakdownMetricsByBuildID(ns.String(), taskQueueName, taskType)
},
AdminNamespaceToPartitionDispatchRate: func() float64 {
return config.AdminNamespaceToPartitionDispatchRate(ns.String())
},
return config.AdminNamespaceToPartitionRateSub(ns.String(), cb)
},
AdminNamespaceTaskQueueToPartitionDispatchRate: func() float64 {
return config.AdminNamespaceTaskqueueToPartitionDispatchRate(ns.String(), taskQueueName, taskType)
},
AdminNamespaceTaskQueueToPartitionRateSub: func(cb func(float64)) (float64, func()) {
config.go ×4
return config.AdminNamespaceTaskqueueToPartitionRateSub(ns.String(), taskQueueName, taskType, cb)
},
forwarderConfig: forwarderConfig{
return config.ForwarderMaxOutstandingPolls(ns.String(), taskQueueName, taskType)
},
ForwarderMaxOutstandingTasks: func() int {
return config.ForwarderMaxOutstandingTasks(ns.String(), taskQueueName, taskType)
},
return config.ForwarderMaxRatePerSecond(ns.String(), taskQueueName, taskType)
},
return max(1, config.ForwarderMaxChildrenPerNode(ns.String(), taskQueueName, taskType))
},
},
GetUserDataRetryPolicy: backoff.NewExponentialRetryPolicy(1 * time.Second).WithMaximumInterval(5 * time.Minute).WithExpirationInterval(backoff.NoInterval),
return config.TaskQueueInfoByBuildIdTTL(ns.String(), taskQueueName, taskType)
},
return config.RateLimitFractionProvider.GetRateLimitFraction(ns, taskQueueName, taskType)
},
RateLimiterRefreshInterval: config.RateLimiterRefreshInterval,
return config.FairnessKeyRateLimitCacheSize(ns.String(), taskQueueName, taskType)
},
MaxFairnessKeyWeightOverrides: func() int {
return config.MaxFairnessKeyWeightOverrides(ns.String(), taskQueueName, taskType)
},
return config.PollerHistoryTTL(ns.String())
},
return config.PollerScalingBacklogAgeScaleUp(ns.String(), taskQueueName, taskType)
},
return config.PollerScalingWaitTime(ns.String(), taskQueueName, taskType)
},
return config.PollerScalingDecisionsPerSecond(ns.String(), taskQueueName, taskType)
},
return config.PollerScalingTaskAddToDispatchRatio(ns.String(), taskQueueName, taskType)
},
return config.EnablePollerScalingDecisionMetrics(ns.String(), taskQueueName, taskType)
},
return config.FairnessCounter(ns.String(), taskQueueName, taskType)
},
return config.FairnessPassDither(ns.String(), taskQueueName, taskType)
},
PartitionScaleAllowedDrift: func() dynamicconfig.PartitionScaleAllowedDrift {
return config.PartitionScaleAllowedDrift(ns.String(), taskQueueName, taskType)
},
PartitionScaleManagerSettings: func() dynamicconfig.PartitionScaleManagerSettings {
service_grpc.pb.go ×20
return config.PartitionScaleManagerSettings(ns.String(), taskQueueName, taskType)
},
MaxVersionsInTaskQueue: func() int { return config.MaxVersionsInTaskQueue(ns.String()) },
}
}
if priority == 0 {
}
priority = min(priority, c.PriorityLevels)
return priority
}
if task.effectivePriority == 0 {
}
}