go.temporal.io/server/service/matching/fx.go

268 LOC · 124 covered · 144 uncovered · 19 ranges · 64 concepts · 3 introducers · 12 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.

1 package matching
2
3 import (
4 "go.temporal.io/server/chasm"
5 "go.temporal.io/server/common"
6 "go.temporal.io/server/common/cluster"
7 "go.temporal.io/server/common/config"
8 "go.temporal.io/server/common/dynamicconfig"
9 "go.temporal.io/server/common/log"
10 "go.temporal.io/server/common/membership"
11 "go.temporal.io/server/common/metrics"
12 "go.temporal.io/server/common/namespace"
13 "go.temporal.io/server/common/persistence"
14 "go.temporal.io/server/common/persistence/serialization"
15 "go.temporal.io/server/common/persistence/visibility"
16 "go.temporal.io/server/common/persistence/visibility/manager"
17 "go.temporal.io/server/common/primitives"
18 "go.temporal.io/server/common/resolver"
19 "go.temporal.io/server/common/resource"
20 "go.temporal.io/server/common/rpc/interceptor"
21 "go.temporal.io/server/common/searchattribute"
22 "go.temporal.io/server/service"
23 "go.temporal.io/server/service/matching/configs"
24 "go.temporal.io/server/service/matching/workers"
25 "go.temporal.io/server/service/worker/workerdeployment"
26 "go.uber.org/fx"
27 "google.golang.org/grpc"
28 healthpb "google.golang.org/grpc/health/grpc_health_v1"
29 )
30
31 var Module = fx.Options(
32 resource.Module,
33 workerdeployment.Module,
34 fx.Provide(ConfigProvider),
35 fx.Provide(PersistenceRateLimitingParamsProvider),
36 service.PersistenceLazyLoadedServiceResolverModule,
37 fx.Provide(ThrottledLoggerRpsFnProvider),
38 fx.Provide(ServiceErrorInterceptorProvider),
39 fx.Provide(ContextMetadataInterceptorProvider),
40 fx.Provide(RetryableInterceptorProvider),
41 fx.Provide(ErrorHandlerProvider),
42 fx.Provide(TelemetryInterceptorProvider),
43 fx.Provide(NamespaceRateLimitInterceptorProvider),
44 fx.Provide(RateLimitInterceptorProvider),
45 fx.Provide(VisibilityManagerProvider),
46 fx.Provide(WorkersRegistryProvider),
47 fx.Provide(NewHandler),
48 fx.Provide(service.GrpcServerOptionsProvider),
49 fx.Provide(NamespaceReplicationQueueProvider),
50 fx.Provide(ServiceResolverProvider),
51 fx.Provide(ServerProvider),
52 fx.Provide(NewService),
53 fx.Provide(simplePartitionScalerFactoryProvider),
54 fx.Provide(taskQueueRateLimitFractionProviderProvider),
55 fx.Invoke(ServiceLifetimeHooks),
56 )
57
58 > func ServerProvider(grpcServerOptions []grpc.ServerOption) *grpc.Server { fx.go ×44
59 > return grpc.NewServer(grpcServerOptions...)
60 > }
61
62 func ConfigProvider(
63 dc *dynamicconfig.Collection,
64 persistenceConfig config.Persistence,
65 rateLimitFractionProvider TaskQueueRateLimitFractionProvider,
66 > ) *Config { fx.go ×44
67 > cfg := NewConfig(dc)
68 > cfg.RateLimitFractionProvider = rateLimitFractionProvider
69 > return cfg
70 > }
71
72 func ServiceErrorInterceptorProvider(
73 dc *dynamicconfig.Collection,
74 > ) *interceptor.ServiceErrorInterceptor { fx.go ×44
75 > return interceptor.NewServiceErrorInterceptor(
76 > dynamicconfig.MaxServiceErrorMessageLength.Get(dc),
77 > )
78 > }
79
80 > func RetryableInterceptorProvider() *interceptor.RetryableInterceptor { fx.go ×44
81 > return interceptor.NewRetryableInterceptor(
82 > common.CreateMatchingHandlerRetryPolicy(),
83 > common.IsServiceHandlerRetryableError,
84 > )
85 > }
86
87 func ErrorHandlerProvider(
88 logger log.Logger,
89 serviceConfig *Config,
90 > ) *interceptor.RequestErrorHandler { fx.go ×44
91 > return interceptor.NewRequestErrorHandler(
92 > logger,
93 > serviceConfig.LogAllReqErrors,
94 > )
95 > }
96
97 func TelemetryInterceptorProvider(
98 logger log.Logger,
99 namespaceRegistry namespace.Registry,
100 metricsHandler metrics.Handler,
101 serviceConfig *Config,
102 requestErrorHandler *interceptor.RequestErrorHandler,
103 > ) *interceptor.TelemetryInterceptor { fx.go ×44
104 > return interceptor.NewTelemetryInterceptor(
105 > namespaceRegistry,
106 > metricsHandler,
107 > logger,
108 > serviceConfig.LogAllReqErrors,
109 > requestErrorHandler,
110 > )
111 > }
112
113 > func ThrottledLoggerRpsFnProvider(serviceConfig *Config) resource.ThrottledLoggerRpsFn { fx.go ×44
114 > return func() float64 { return float64(serviceConfig.ThrottledLogRPS()) }
115 }
116
117 func NamespaceRateLimitInterceptorProvider(
118 serviceConfig *Config,
119 namespaceRegistry namespace.Registry,
120 metricsHandler metrics.Handler,
121 > ) interceptor.NamespaceRateLimitInterceptor { fx.go ×44
122 >
123 > namespaceRateFn := func(namespaceName string) float64 {
124 > if namespaceRPS := serviceConfig.NamespaceRPS(namespaceName); namespaceRPS > 0 { service_grpc.pb.go ×20
125 return float64(namespaceRPS)
126 }
127 // This fallback to host level rps limit when NamespaceRPS is not configured (i.e. 0)
128 > return float64(serviceConfig.RPS()) service_grpc.pb.go ×20
129 }
130
131 > return interceptor.NewNamespaceRateLimitInterceptor( fx.go ×44
132 > namespaceRegistry,
133 > configs.NewNamespaceRateLimiter(
134 > namespaceRateFn,
135 > serviceConfig.OperatorRPSRatio,
136 > ),
137 > map[string]int{}, // no token overrides
138 > configs.PollTaskAPISet, // set of APIs that will wait for token instead of immediate rejection
139 > serviceConfig.PollWaitForNamespaceRateLimitToken,
140 > metricsHandler,
141 > )
142 }
143
144 func RateLimitInterceptorProvider(
145 serviceConfig *Config,
146 > ) *interceptor.RateLimitInterceptor { fx.go ×44
147 > return interceptor.NewRateLimitInterceptor(
148 > configs.NewPriorityRateLimiter(func() float64 { return float64(serviceConfig.RPS()) }, serviceConfig.OperatorRPSRatio),
149 map[string]int{
150 healthpb.Health_Check_FullMethodName: 0, // exclude health check requests from rate limiting.
151 },
152 )
153 }
154
155 // PersistenceRateLimitingParamsProvider is the same between services but uses different config sources.
156 // if-case comes from resourceImpl.New.
157 func PersistenceRateLimitingParamsProvider(
158 serviceConfig *Config,
159 persistenceLazyLoadedServiceResolver service.PersistenceLazyLoadedServiceResolver,
160 logger log.SnTaggedLogger,
161 > ) service.PersistenceRateLimitingParams { fx.go ×44
162 > return service.NewPersistenceRateLimitingParams(
163 > serviceConfig.PersistenceMaxQPS,
164 > serviceConfig.PersistenceGlobalMaxQPS,
165 > serviceConfig.PersistenceNamespaceMaxQPS,
166 > serviceConfig.PersistenceGlobalNamespaceMaxQPS,
167 > serviceConfig.PersistencePerShardNamespaceMaxQPS,
168 > serviceConfig.OperatorRPSRatio,
169 > serviceConfig.PersistenceQPSBurstRatio,
170 > serviceConfig.PersistenceDynamicRateLimitingParams,
171 > persistenceLazyLoadedServiceResolver,
172 > logger,
173 > )
174 > }
175
176 func ServiceResolverProvider(
177 membershipMonitor membership.Monitor,
178 > ) (membership.ServiceResolver, error) { fx.go ×44
179 > return membershipMonitor.GetResolver(primitives.MatchingService)
180 > }
181
182 // TaskQueueReplicatorNamespaceReplicationQueue is used to ensure the replicator only gets set if global namespaces are
183 // enabled on this cluster. See NamespaceReplicationQueueProvider below.
184 type TaskQueueReplicatorNamespaceReplicationQueue persistence.NamespaceReplicationQueue
185
186 func NamespaceReplicationQueueProvider(
187 namespaceReplicationQueue persistence.NamespaceReplicationQueue,
188 clusterMetadata cluster.Metadata,
189 ) TaskQueueReplicatorNamespaceReplicationQueue {
190 var replicatorNamespaceReplicationQueue persistence.NamespaceReplicationQueue
191 if clusterMetadata.IsGlobalNamespaceEnabled() {
192 replicatorNamespaceReplicationQueue = namespaceReplicationQueue
193 }
194 return replicatorNamespaceReplicationQueue
195 }
196
197 func VisibilityManagerProvider(
198 logger log.Logger,
199 persistenceConfig *config.Persistence,
200 customVisibilityStoreFactory visibility.VisibilityStoreFactory,
201 metricsHandler metrics.Handler,
202 serviceConfig *Config,
203 persistenceServiceResolver resolver.ServiceResolver,
204 searchAttributesMapperProvider searchattribute.MapperProvider,
205 saProvider searchattribute.Provider,
206 namespaceRegistry namespace.Registry,
207 chasmRegistry *chasm.Registry,
208 serializer serialization.Serializer,
209 > ) (manager.VisibilityManager, error) { fx.go ×44
210 > return visibility.NewManager(
211 > *persistenceConfig,
212 > persistenceServiceResolver,
213 > customVisibilityStoreFactory,
214 > nil, // matching visibility never writes
215 > saProvider,
216 > searchAttributesMapperProvider,
217 > namespaceRegistry,
218 > chasmRegistry,
219 > serviceConfig.VisibilityPersistenceMaxReadQPS,
220 > serviceConfig.VisibilityPersistenceMaxWriteQPS,
221 > serviceConfig.OperatorRPSRatio,
222 > serviceConfig.VisibilityPersistenceSlowQueryThreshold,
223 > serviceConfig.EnableReadFromSecondaryVisibility,
224 > serviceConfig.VisibilityEnableShadowReadMode,
225 > dynamicconfig.GetStringPropertyFn(visibility.SecondaryVisibilityWritingModeOff), // matching visibility never writes
226 > serviceConfig.VisibilityDisableOrderByClause,
227 > serviceConfig.VisibilityEnableManualPagination,
228 > serviceConfig.VisibilityEnableUnifiedQueryConverter,
229 > metricsHandler,
230 > logger,
231 > serializer,
232 > )
233 > }
234
235 > func ContextMetadataInterceptorProvider(logger log.Logger) *interceptor.ContextMetadataInterceptor { fx.go ×44
236 > return interceptor.NewContextMetadataInterceptor(true, logger)
237 > }
238
239 > func ServiceLifetimeHooks(lc fx.Lifecycle, svc *Service) { fx.go ×44
240 > lc.Append(fx.StartStopHook(svc.Start, svc.Stop))
241 > }
242
243 func WorkersRegistryProvider(
244 lc fx.Lifecycle,
245 metricsHandler metrics.Handler,
246 serviceConfig *Config,
247 > ) workers.Registry { fx.go ×44
248 > return workers.NewRegistry(lc, workers.RegistryParams{
249 > NumBuckets: serviceConfig.WorkerRegistryNumBuckets,
250 > TTL: serviceConfig.WorkerRegistryEntryTTL,
251 > MinEvictAge: serviceConfig.WorkerRegistryMinEvictAge,
252 > MaxItems: serviceConfig.WorkerRegistryMaxEntries,
253 > EvictionInterval: serviceConfig.WorkerRegistryEvictionInterval,
254 > MetricsHandler: metricsHandler,
255 > MetricsConfig: workers.WorkerMetricsConfig{
256 > EnablePluginMetrics: serviceConfig.EnableWorkerPluginMetrics,
257 > EnablePollerAutoscalingMetrics: serviceConfig.EnablePollerAutoscalingMetrics,
258 > BreakdownMetricsByTaskQueue: serviceConfig.BreakdownMetricsByTaskQueue,
259 > ExternalPayloadsEnabled: serviceConfig.ExternalPayloadsEnabled,
260 > },
261 > })
262 > }
263
264 > func simplePartitionScalerFactoryProvider(dc *dynamicconfig.Collection) PartitionScalerFactory { fx.go ×44
265 > return newSimplePartitionScalerFactory(
266 > dynamicconfig.MatchingPartitionScaler.Get(dc),
267 > )
268 > }