278
func NewConfig(
279
dc *dynamicconfig.Collection,
281
>
return &Config{
282
>
PersistenceMaxQPS: dynamicconfig.MatchingPersistenceMaxQPS.Get(dc),
283
>
PersistenceGlobalMaxQPS: dynamicconfig.MatchingPersistenceGlobalMaxQPS.Get(dc),
284
>
PersistenceNamespaceMaxQPS: dynamicconfig.MatchingPersistenceNamespaceMaxQPS.Get(dc),
285
>
PersistenceGlobalNamespaceMaxQPS: dynamicconfig.MatchingPersistenceGlobalNamespaceMaxQPS.Get(dc),
286
>
PersistencePerShardNamespaceMaxQPS: dynamicconfig.DefaultPerShardNamespaceRPSMax,
287
>
PersistenceDynamicRateLimitingParams: dynamicconfig.MatchingPersistenceDynamicRateLimitingParams.Get(dc),
288
>
PersistenceQPSBurstRatio: dynamicconfig.PersistenceQPSBurstRatio.Get(dc),
289
>
SyncMatchWaitDuration: dynamicconfig.MatchingSyncMatchWaitDuration.Get(dc),
290
>
HistoryMaxPageSize: dynamicconfig.MatchingHistoryMaxPageSize.Get(dc),
291
>
EnableDeployments: dynamicconfig.EnableDeployments.Get(dc), // [cleanup-wv-pre-release]
292
>
EnableDeploymentVersions: dynamicconfig.EnableDeploymentVersions.Get(dc),
293
>
UseRevisionNumberForWorkerVersioning: dynamicconfig.UseRevisionNumberForWorkerVersioning.Get(dc),
294
>
MaxTaskQueuesInDeployment: dynamicconfig.MatchingMaxTaskQueuesInDeployment.Get(dc),
295
>
MaxVersionsInTaskQueue: dynamicconfig.MatchingMaxVersionsInTaskQueue.Get(dc),
296
>
RPS: dynamicconfig.MatchingRPS.Get(dc),
297
>
NamespaceRPS: dynamicconfig.MatchingNamespaceRPS.Get(dc),
298
>
OperatorRPSRatio: dynamicconfig.OperatorRPSRatio.Get(dc),
299
>
PollWaitForNamespaceRateLimitToken: dynamicconfig.PollWaitForNamespaceRateLimitToken.Get(dc),
300
>
RangeSize: 100000,
301
>
NewMatcherSub: dynamicconfig.MatchingUseNewMatcher.Subscribe(dc),
302
>
EnableFairnessSub: dynamicconfig.MatchingEnableFairness.Subscribe(dc),
303
>
EnableMigration: dynamicconfig.MatchingEnableMigration.Get(dc),
304
>
AutoEnableV2Sub: dynamicconfig.MatchingAutoEnableV2.Subscribe(dc),
305
>
GetTasksBatchSize: dynamicconfig.MatchingGetTasksBatchSize.Get(dc),
306
>
GetTasksReloadAt: dynamicconfig.MatchingGetTasksReloadAt.Get(dc),
307
>
ForceReadTasksOnWrite: dynamicconfig.MatchingForceReadTasksOnWrite.Get(dc),
308
>
UpdateAckInterval: dynamicconfig.MatchingUpdateAckInterval.Get(dc),
309
>
MetadataUpdateOnAppendInterval: dynamicconfig.MatchingMetadataUpdateOnAppendInterval.Get(dc),
310
>
MaxTaskQueueIdleTime: dynamicconfig.MatchingMaxTaskQueueIdleTime.Get(dc),
311
>
LongPollExpirationInterval: dynamicconfig.MatchingLongPollExpirationInterval.Get(dc),
312
>
BacklogTaskForwardTimeout: dynamicconfig.MatchingBacklogTaskForwardTimeout.Get(dc),
313
>
ForwardPollRetryMaxInterval: dynamicconfig.MatchingForwardPollRetryMaxInterval.Get(dc),
314
>
MinTaskThrottlingBurstSize: dynamicconfig.MatchingMinTaskThrottlingBurstSize.Get(dc),
315
>
MaxTaskDeleteBatchSize: dynamicconfig.MatchingMaxTaskDeleteBatchSize.Get(dc),
316
>
TaskDeleteInterval: dynamicconfig.MatchingTaskDeleteInterval.Get(dc),
317
>
OutstandingTaskAppendsThreshold: dynamicconfig.MatchingOutstandingTaskAppendsThreshold.Get(dc),
318
>
MaxTaskBatchSize: dynamicconfig.MatchingMaxTaskBatchSize.Get(dc),
319
>
ThrottledLogRPS: dynamicconfig.MatchingThrottledLogRPS.Get(dc),
320
>
NumTaskqueueWritePartitions: dynamicconfig.MatchingNumTaskqueueWritePartitions.Get(dc),
321
>
NumTaskqueueReadPartitions: dynamicconfig.MatchingNumTaskqueueReadPartitions.Get(dc),
322
>
NumTaskqueueReadPartitionsSub: dynamicconfig.MatchingNumTaskqueueReadPartitions.Subscribe(dc),
323
>
BreakdownMetricsByTaskQueue: dynamicconfig.MetricsBreakdownByTaskQueue.Get(dc),
324
>
BreakdownMetricsByPartition: dynamicconfig.MetricsBreakdownByPartition.Get(dc),
325
>
BreakdownMetricsByBuildID: dynamicconfig.MetricsBreakdownByBuildID.Get(dc),
326
>
EnableWorkerPluginMetrics: dynamicconfig.MatchingEnableWorkerPluginMetrics.Get(dc),
327
>
EnablePollerAutoscalingMetrics: dynamicconfig.MatchingEnablePollerAutoscalingMetrics.Get(dc),
328
>
ExternalPayloadsEnabled: dynamicconfig.ExternalPayloadsEnabled.Get(dc),
329
>
WorkerRegistryNumBuckets: dynamicconfig.MatchingWorkerRegistryNumBuckets.Get(dc),
330
>
WorkerRegistryEntryTTL: dynamicconfig.MatchingWorkerRegistryEntryTTL.Get(dc),
331
>
WorkerRegistryMinEvictAge: dynamicconfig.MatchingWorkerRegistryMinEvictAge.Get(dc),
332
>
WorkerRegistryMaxEntries: dynamicconfig.MatchingWorkerRegistryMaxEntries.Get(dc),
333
>
WorkerRegistryEvictionInterval: dynamicconfig.MatchingWorkerRegistryEvictionInterval.Get(dc),
334
>
ForwarderMaxOutstandingPolls: dynamicconfig.MatchingForwarderMaxOutstandingPolls.Get(dc),
335
>
ForwarderMaxOutstandingTasks: dynamicconfig.MatchingForwarderMaxOutstandingTasks.Get(dc),
336
>
ForwarderMaxRatePerSecond: dynamicconfig.MatchingForwarderMaxRatePerSecond.Get(dc),
337
>
ForwarderMaxChildrenPerNode: dynamicconfig.MatchingForwarderMaxChildrenPerNode.Get(dc),
338
>
AlignMembershipChange: dynamicconfig.MatchingAlignMembershipChange.Get(dc),
339
>
ShutdownDrainDuration: dynamicconfig.MatchingShutdownDrainDuration.Get(dc),
340
>
VersionCompatibleSetLimitPerQueue: dynamicconfig.VersionCompatibleSetLimitPerQueue.Get(dc),
341
>
VersionBuildIdLimitPerQueue: dynamicconfig.VersionBuildIdLimitPerQueue.Get(dc),
342
>
AssignmentRuleLimitPerQueue: dynamicconfig.AssignmentRuleLimitPerQueue.Get(dc),
343
>
RedirectRuleLimitPerQueue: dynamicconfig.RedirectRuleLimitPerQueue.Get(dc),
344
>
RedirectRuleMaxUpstreamBuildIDsPerQueue: dynamicconfig.RedirectRuleMaxUpstreamBuildIDsPerQueue.Get(dc),
345
>
DeletedRuleRetentionTime: dynamicconfig.MatchingDeletedRuleRetentionTime.Get(dc),
346
>
PollerHistoryTTL: dynamicconfig.PollerHistoryTTL.Get(dc),
347
>
EnableMatchingFanOutForPollCancellation: dynamicconfig.EnableMatchingFanOutForPollCancellation.Get(dc),
348
>
ReachabilityBuildIdVisibilityGracePeriod: dynamicconfig.ReachabilityBuildIdVisibilityGracePeriod.Get(dc),
349
>
ReachabilityCacheOpenWFsTTL: dynamicconfig.ReachabilityCacheOpenWFsTTL.Get(dc),
350
>
ReachabilityCacheClosedWFsTTL: dynamicconfig.ReachabilityCacheClosedWFsTTL.Get(dc),
351
>
TaskQueueLimitPerBuildId: dynamicconfig.TaskQueuesPerBuildIdLimit.Get(dc),
352
>
GetUserDataLongPollTimeout: dynamicconfig.MatchingGetUserDataLongPollTimeout.Get(dc), // Use -10 seconds so that we send back empty response instead of timeout
353
>
GetUserDataRefresh: dynamicconfig.MatchingGetUserDataRefresh.Get(dc),
354
>
EphemeralDataUpdateInterval: dynamicconfig.MatchingEphemeralDataUpdateInterval.Get(dc),
355
>
BacklogMetricsEmitInterval: dynamicconfig.MatchingBacklogMetricsEmitInterval.Get(dc),
356
>
PriorityBacklogForwarding: dynamicconfig.MatchingPriorityBacklogForwarding.Get(dc),
357
>
BacklogNegligibleAge: dynamicconfig.MatchingBacklogNegligibleAge.Get(dc),
358
>
MaxWaitForPollerBeforeFwd: dynamicconfig.MatchingMaxWaitForPollerBeforeFwd.Get(dc),
359
>
QueryPollerUnavailableWindow: dynamicconfig.QueryPollerUnavailableWindow.Get(dc),
360
>
WorkerControllerNoPollerHookWindow: dynamicconfig.WorkerControllerNoPollerHookWindow.Get(dc),
361
>
EmitTaskDispatchLatencyAtPoll: dynamicconfig.MatchingEmitTaskDispatchLatencyAtPoll.Get(dc),
362
>
QueryWorkflowTaskTimeoutLogRate: dynamicconfig.MatchingQueryWorkflowTaskTimeoutLogRate.Get(dc),
363
>
MembershipUnloadDelay: dynamicconfig.MatchingMembershipUnloadDelay.Get(dc),
364
>
TaskQueueInfoByBuildIdTTL: dynamicconfig.TaskQueueInfoByBuildIdTTL.Get(dc),
365
>
PriorityLevels: dynamicconfig.MatchingPriorityLevels.Get(dc),
366
>
RateLimiterRefreshInterval: time.Minute,
367
>
FairnessKeyRateLimitCacheSize: dynamicconfig.MatchingFairnessKeyRateLimitCacheSize.Get(dc),
368
>
MaxFairnessKeyWeightOverrides: dynamicconfig.MatchingMaxFairnessKeyWeightOverrides.Get(dc),
369
>
MaxIDLengthLimit: dynamicconfig.MaxIDLengthLimit.Get(dc),
370
>
371
>
AdminNamespaceToPartitionDispatchRate: dynamicconfig.AdminMatchingNamespaceToPartitionDispatchRate.Get(dc),
372
>
AdminNamespaceToPartitionRateSub: dynamicconfig.AdminMatchingNamespaceToPartitionDispatchRate.Subscribe(dc),
373
>
AdminNamespaceTaskqueueToPartitionDispatchRate: dynamicconfig.AdminMatchingNamespaceTaskqueueToPartitionDispatchRate.Get(dc),
374
>
AdminNamespaceTaskqueueToPartitionRateSub: dynamicconfig.AdminMatchingNamespaceTaskqueueToPartitionDispatchRate.Subscribe(dc),
375
>
376
>
VisibilityPersistenceMaxReadQPS: dynamicconfig.VisibilityPersistenceMaxReadQPS.Get(dc),
377
>
VisibilityPersistenceMaxWriteQPS: dynamicconfig.VisibilityPersistenceMaxWriteQPS.Get(dc),
378
>
VisibilityPersistenceSlowQueryThreshold: dynamicconfig.VisibilityPersistenceSlowQueryThreshold.Get(dc),
379
>
EnableReadFromSecondaryVisibility: dynamicconfig.EnableReadFromSecondaryVisibility.Get(dc),
380
>
VisibilityEnableShadowReadMode: dynamicconfig.VisibilityEnableShadowReadMode.Get(dc),
381
>
VisibilityDisableOrderByClause: dynamicconfig.VisibilityDisableOrderByClause.Get(dc),
382
>
VisibilityEnableManualPagination: dynamicconfig.VisibilityEnableManualPagination.Get(dc),
383
>
VisibilityEnableUnifiedQueryConverter: dynamicconfig.VisibilityEnableUnifiedQueryConverter.Get(dc),
384
>
385
>
ListNexusEndpointsLongPollTimeout: dynamicconfig.MatchingListNexusEndpointsLongPollTimeout.Get(dc),
386
>
NexusEndpointsRefreshInterval: dynamicconfig.MatchingNexusEndpointsRefreshInterval.Get(dc),
387
>
MinDispatchTaskTimeout: nexusoperations.MinDispatchTaskTimeout.Get(dc),
388
>
389
>
PollerScalingBacklogAgeScaleUp: dynamicconfig.MatchingPollerScalingBacklogAgeScaleUp.Get(dc),
390
>
PollerScalingWaitTime: dynamicconfig.MatchingPollerScalingWaitTime.Get(dc),
391
>
PollerScalingDecisionsPerSecond: dynamicconfig.MatchingPollerScalingDecisionsPerSecond.Get(dc),
392
>
PollerScalingTaskAddToDispatchRatio: dynamicconfig.MatchingPollerScalingTaskAddToDispatchRatio.Get(dc),
393
>
EnablePollerScalingDecisionMetrics: dynamicconfig.MatchingEnablePollerScalingDecisionMetrics.Get(dc),
394
>
395
>
FairnessCounter: dynamicconfig.MatchingFairnessCounter.Get(dc),
396
>
FairnessPassDither: dynamicconfig.MatchingFairnessPassDither.Get(dc),
397
>
PartitionScaleAllowedDrift: dynamicconfig.MatchingPartitionScaleAllowedDrift.Get(dc),
398
>
PartitionScaleManagerSettings: dynamicconfig.MatchingPartitionScaleManager.Get(dc),
399
>
400
>
LogAllReqErrors: dynamicconfig.LogAllReqErrors.Get(dc),
401
>
402
>
RateLimitFractionProvider: defaultTaskQueueRateLimitFractionProvider,
403
>
}
404
>
}
405
406
>
func newTaskQueueConfig(tq *tqid.TaskQueue, config *Config, ns namespace.Name) *taskQueueConfig {
config.go
407
>
taskQueueName := tq.Name()
408
>
taskType := tq.TaskType()
409
>
priorityLevels := priorityKey(config.PriorityLevels(ns.String(), taskQueueName, taskType))
410
>
priorityLevels = max(priorityLevels, min(priorityLevels, maxPriorityLevels), 1)
411
>
defaultPriorityKey := (priorityLevels + 1) / 2
412
>
413
>
return &taskQueueConfig{
414
>
RangeSize: config.RangeSize,
415
>
NewMatcherSub: func(cb func(dynamicconfig.GradualChange[bool])) (dynamicconfig.GradualChange[bool], func()) {
416
>
return config.NewMatcherSub(ns.String(), taskQueueName, taskType, cb)
config.go
417
>
},
418
>
EnableFairnessSub: func(cb func(dynamicconfig.GradualChange[bool])) (dynamicconfig.GradualChange[bool], func()) {
419
>
return config.EnableFairnessSub(ns.String(), taskQueueName, taskType, cb)
420
>
},
421
>
EnableMigration: func() bool {
config.go
422
>
return config.EnableMigration(ns.String(), taskQueueName, taskType)
423
>
},
424
AutoEnableV2: func() bool {
425
v, _ := config.AutoEnableV2Sub(ns.String(), taskQueueName, taskType, nil)
426
return v
427
},
428
>
AutoEnableV2Sub: func(cb func(bool)) (bool, func()) {
config.go
429
>
return config.AutoEnableV2Sub(ns.String(), taskQueueName, taskType, cb)
430
>
},
431
>
GetTasksBatchSize: func() int {
config.go
432
>
return config.GetTasksBatchSize(ns.String(), taskQueueName, taskType)
433
>
},
434
>
GetTasksReloadAt: func() int {
config.go
435
>
return config.GetTasksReloadAt(ns.String(), taskQueueName, taskType)
436
>
},
437
ForceReadTasksOnWrite: func() bool {
438
return config.ForceReadTasksOnWrite(ns.String(), taskQueueName, taskType)
439
},
440
>
UpdateAckInterval: func() time.Duration {
config.go
441
>
return config.UpdateAckInterval(ns.String(), taskQueueName, taskType)
442
>
},
443
MetadataUpdateOnAppendInterval: func() time.Duration {
444
return config.MetadataUpdateOnAppendInterval(ns.String(), taskQueueName, taskType)
445
},
446
>
MaxTaskQueueIdleTime: func() time.Duration {
config.go
447
>
return config.MaxTaskQueueIdleTime(ns.String(), taskQueueName, taskType)
448
>
},
449
MinTaskThrottlingBurstSize: func() int {
450
return config.MinTaskThrottlingBurstSize(ns.String(), taskQueueName, taskType)