go.temporal.io/server/common/resource/fx.go

533 LOC · 252 covered · 281 uncovered · 48 ranges · 64 concepts · 4 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 resource
2
3 import (
4 "crypto/tls"
5 "fmt"
6 "net"
7 "os"
8 "time"
9
10 "go.temporal.io/api/workflowservice/v1"
11 "go.temporal.io/server/api/adminservice/v1"
12 "go.temporal.io/server/api/historyservice/v1"
13 "go.temporal.io/server/api/matchingservice/v1"
14 "go.temporal.io/server/client"
15 "go.temporal.io/server/client/admin"
16 "go.temporal.io/server/client/frontend"
17 "go.temporal.io/server/client/history"
18 "go.temporal.io/server/client/matching"
19 "go.temporal.io/server/common"
20 "go.temporal.io/server/common/archiver"
21 "go.temporal.io/server/common/archiver/provider"
22 "go.temporal.io/server/common/clock"
23 "go.temporal.io/server/common/cluster"
24 "go.temporal.io/server/common/config"
25 "go.temporal.io/server/common/deadlock"
26 "go.temporal.io/server/common/dynamicconfig"
27 "go.temporal.io/server/common/log"
28 "go.temporal.io/server/common/log/tag"
29 "go.temporal.io/server/common/membership"
30 "go.temporal.io/server/common/metrics"
31 "go.temporal.io/server/common/namespace"
32 "go.temporal.io/server/common/namespace/nsregistry"
33 commonnexus "go.temporal.io/server/common/nexus"
34 "go.temporal.io/server/common/persistence"
35 persistenceClient "go.temporal.io/server/common/persistence/client"
36 "go.temporal.io/server/common/persistence/serialization"
37 "go.temporal.io/server/common/persistence/visibility"
38 "go.temporal.io/server/common/persistence/visibility/manager"
39 "go.temporal.io/server/common/pingable"
40 "go.temporal.io/server/common/primitives"
41 "go.temporal.io/server/common/quotas"
42 "go.temporal.io/server/common/rpc"
43 "go.temporal.io/server/common/rpc/auth"
44 "go.temporal.io/server/common/rpc/encryption"
45 "go.temporal.io/server/common/rpc/interceptor"
46 "go.temporal.io/server/common/sdk"
47 "go.temporal.io/server/common/searchattribute"
48 "go.temporal.io/server/common/telemetry"
49 "go.temporal.io/server/common/testing/testhooks"
50 "go.uber.org/fx"
51 "google.golang.org/grpc"
52 "google.golang.org/grpc/health"
53 )
54
55 type (
56 ThrottledLoggerRpsFn quotas.RateFn
57 NamespaceLogger log.Logger
58 HostName string
59 InstanceID string
60 ServiceNames map[primitives.ServiceName]struct{}
61
62 HistoryRawClient historyservice.HistoryServiceClient
63 HistoryClient historyservice.HistoryServiceClient
64
65 MatchingRawClient matchingservice.MatchingServiceClient
66 MatchingClient matchingservice.MatchingServiceClient
67
68 RuntimeMetricsReporterParams struct {
69 fx.In
70
71 MetricHandler metrics.Handler
72 Logger log.SnTaggedLogger
73 InstanceID InstanceID `optional:"true"`
74 }
75 )
76
77 // Module
78 // Use fx.Hook and OnStart/OnStop to manage Daemon resource lifecycle
79 // See LifetimeHooksModule for detail
80 var Module = fx.Options(
81 persistenceClient.Module,
82 dynamicconfig.Module,
83 serialization.Module,
84 fx.Provide(HostNameProvider),
85 fx.Provide(TimeSourceProvider),
86 cluster.MetadataLifetimeHooksModule,
87 fx.Provide(SearchAttributeMapperProviderProvider),
88 fx.Provide(SearchAttributeProviderProvider),
89 fx.Provide(SearchAttributeManagerProvider),
90 fx.Provide(NamespaceRegistryProvider),
91 > fx.Provide(func() namespace.NamespaceStateChangedFn { return nsregistry.DefaultNamespaceStateChanged }), fx.go ×44
92 nsregistry.RegistryLifetimeHooksModule,
93 fx.Provide(fx.Annotate(
94 > func(p namespace.Registry) pingable.Pingable { return p }, fx.go ×44
95 fx.ResultTags(`group:"deadlockDetectorRoots"`),
96 )),
97 fx.Provide(ClientFactoryProvider),
98 fx.Provide(ClientBeanProvider),
99 fx.Provide(FrontendClientProvider),
100 fx.Provide(AdminClientProvider),
101 fx.Provide(GrpcListenerProvider),
102 fx.Provide(RuntimeMetricsReporterProvider),
103 metrics.RuntimeMetricsReporterLifetimeHooksModule,
104 fx.Provide(HistoryRawClientProvider),
105 fx.Provide(HistoryClientProvider),
106 fx.Provide(MatchingRawClientProvider),
107 fx.Provide(MatchingClientProvider),
108 membership.GRPCResolverModule,
109 fx.Provide(FrontendHTTPClientCacheProvider),
110 fx.Provide(PersistenceConfigProvider),
111 fx.Provide(health.NewServer),
112 fx.Provide(namespace.NewDefaultReplicationResolverFactory),
113 deadlock.Module,
114 config.Module,
115 testhooks.Module,
116 fx.Provide(commonnexus.NewLoggedHTTPClientTraceProvider),
117 )
118
119 var DefaultOptions = fx.Options(
120 fx.Provide(RPCFactoryProvider),
121 fx.Provide(PerServiceDialOptionsProvider),
122 fx.Provide(ArchivalMetadataProvider),
123 fx.Provide(ArchiverProviderProvider),
124 fx.Provide(ThrottledLoggerProvider),
125 fx.Provide(SdkClientFactoryProvider),
126 fx.Provide(DCRedirectionPolicyProvider),
127 )
128
129 > func DefaultSnTaggedLoggerProvider(logger log.Logger, sn primitives.ServiceName) log.SnTaggedLogger { fx.go ×44
130 > return log.With(logger, tag.Service(sn))
131 > }
132
133 func ThrottledLoggerProvider(
134 logger log.SnTaggedLogger,
135 fn ThrottledLoggerRpsFn,
136 > ) log.ThrottledLogger { fx.go ×44
137 > return log.NewThrottledLogger(
138 > logger,
139 > quotas.RateFn(fn),
140 > )
141 > }
142
143 > func GrpcListenerProvider(factory common.RPCFactory) net.Listener { fx.go ×44
144 > return factory.GetGRPCListener()
145 > }
146
147 > func HostNameProvider() (HostName, error) { fx.go ×44
148 > hn, err := os.Hostname()
149 > return HostName(hn), err
150 > }
151
152 > func TimeSourceProvider() clock.TimeSource { fx.go ×44
153 > return clock.NewRealTimeSource()
154 > }
155
156 func SearchAttributeMapperProviderProvider(
157 saMapper searchattribute.Mapper,
158 namespaceRegistry namespace.Registry,
159 searchAttributeProvider searchattribute.Provider,
160 persistenceConfig *config.Persistence,
161 > ) searchattribute.MapperProvider { fx.go ×44
162 > primaryVisibilityStoreConfig := persistenceConfig.GetVisibilityStoreConfig()
163 > return searchattribute.NewMapperProvider(
164 > saMapper,
165 > namespaceRegistry,
166 > searchAttributeProvider,
167 > primaryVisibilityStoreConfig.GetIndexName(),
168 > )
169 > }
170
171 func SearchAttributeProviderProvider(
172 logger log.SnTaggedLogger,
173 timeSource clock.TimeSource,
174 cmMgr persistence.ClusterMetadataManager,
175 dynamicCollection *dynamicconfig.Collection,
176 > ) searchattribute.Provider { fx.go ×44
177 > return searchattribute.NewManager(
178 > timeSource,
179 > cmMgr,
180 > logger,
181 > dynamicconfig.ForceSearchAttributesCacheRefreshOnRead.Get(dynamicCollection))
182 > }
183
184 func SearchAttributeManagerProvider(
185 logger log.SnTaggedLogger,
186 timeSource clock.TimeSource,
187 cmMgr persistence.ClusterMetadataManager,
188 dynamicCollection *dynamicconfig.Collection,
189 > ) searchattribute.Manager { fx.go ×44
190 > return searchattribute.NewManager(
191 > timeSource,
192 > cmMgr,
193 > logger,
194 > dynamicconfig.ForceSearchAttributesCacheRefreshOnRead.Get(dynamicCollection))
195 > }
196
197 // SearchAttributeValidatorProvider creates a new search attribute validator with the given dependencies. It configures
198 // the validator with dynamic config values for key limits, value size limits, total size limits, visibility allowlist,
199 // and system search attribute error suppression.
200 func SearchAttributeValidatorProvider(
201 saProvider searchattribute.Provider,
202 saMapperProvider searchattribute.MapperProvider,
203 visibilityMgr manager.VisibilityManager,
204 dynamicCollection *dynamicconfig.Collection,
205 metricsHandler metrics.Handler,
206 logger log.Logger,
207 > ) *searchattribute.Validator { fx.go ×44
208 > return searchattribute.NewValidator(
209 > saProvider,
210 > saMapperProvider,
211 > dynamicconfig.SearchAttributesNumberOfKeysLimit.Get(dynamicCollection),
212 > dynamicconfig.SearchAttributesSizeOfValueLimit.Get(dynamicCollection),
213 > dynamicconfig.SearchAttributesTotalSizeLimit.Get(dynamicCollection),
214 > visibilityMgr,
215 > visibility.AllowListForValidation(
216 > visibilityMgr.GetStoreNames(),
217 > dynamicconfig.VisibilityAllowList.Get(dynamicCollection),
218 > ),
219 > dynamicconfig.SuppressErrorSetSystemSearchAttribute.Get(dynamicCollection),
220 > metricsHandler,
221 > logger,
222 > )
223 > }
224
225 type NamespaceRegistryParams struct {
226 fx.In
227
228 Logger log.SnTaggedLogger
229 MetricsHandler metrics.Handler
230 ClusterMetadata cluster.Metadata
231 MetadataManager persistence.MetadataManager
232 DynamicCollection *dynamicconfig.Collection
233 ReplicationResolverFactory namespace.ReplicationResolverFactory
234 NamespaceStateChangedFn namespace.NamespaceStateChangedFn
235 }
236
237 > func NamespaceRegistryProvider(params NamespaceRegistryParams) namespace.Registry { fx.go ×44
238 > return nsregistry.NewRegistry(
239 > params.MetadataManager,
240 > params.ClusterMetadata.IsGlobalNamespaceEnabled(),
241 > params.ClusterMetadata.GetCurrentClusterName(),
242 > dynamicconfig.NamespaceCacheRefreshInterval.Get(params.DynamicCollection),
243 > dynamicconfig.ForceSearchAttributesCacheRefreshOnRead.Get(params.DynamicCollection),
244 > params.MetricsHandler,
245 > params.Logger,
246 > params.ReplicationResolverFactory,
247 > params.NamespaceStateChangedFn,
248 > )
249 > }
250
251 func ClientFactoryProvider(
252 factoryProvider client.FactoryProvider,
253 rpcFactory common.RPCFactory,
254 membershipMonitor membership.Monitor,
255 metricsHandler metrics.Handler,
256 dynamicCollection *dynamicconfig.Collection,
257 testHooks testhooks.TestHooks,
258 persistenceConfig *config.Persistence,
259 logger log.SnTaggedLogger,
260 throttledLogger log.ThrottledLogger,
261 > ) client.Factory { fx.go ×44
262 > return factoryProvider.NewFactory(
263 > rpcFactory,
264 > membershipMonitor,
265 > metricsHandler,
266 > dynamicCollection,
267 > testHooks,
268 > persistenceConfig.NumHistoryShards,
269 > logger,
270 > throttledLogger,
271 > )
272 > }
273
274 func ClientBeanProvider(
275 lc fx.Lifecycle,
276 clientFactory client.Factory,
277 clusterMetadata cluster.Metadata,
278 > ) (client.Bean, error) { fx.go ×44
279 > bean, err := client.NewClientBean(
280 > clientFactory,
281 > clusterMetadata,
282 > )
283 > if err != nil {
284 return nil, err
285 }
286 // Deterministically release the bean's clients (daemon goroutines and
287 // cached gRPC connections) on shutdown.
288 > lc.Append(fx.StopHook(bean.Close)) fx.go ×44
289 > return bean, nil
290 }
291
292 > func FrontendClientProvider(clientBean client.Bean) workflowservice.WorkflowServiceClient { fx.go ×44
293 > frontendRawClient := clientBean.GetFrontendClient()
294 > return frontend.NewRetryableClient(
295 > frontendRawClient,
296 > common.CreateFrontendClientRetryPolicy(),
297 > common.IsServiceClientTransientError,
298 > )
299 > }
300
301 > func AdminClientProvider(clientBean client.Bean, clusterMetadata cluster.Metadata) (adminservice.AdminServiceClient, error) { fx.go ×44
302 > adminRawClient, err := clientBean.GetRemoteAdminClient(clusterMetadata.GetCurrentClusterName())
303 > if err != nil {
304 return nil, err
305 }
306 > return admin.NewRetryableClient( fx.go ×44
307 > adminRawClient,
308 > common.CreateFrontendClientRetryPolicy(),
309 > common.IsServiceClientTransientError,
310 > ), nil
311 }
312
313 func RuntimeMetricsReporterProvider(
314 params RuntimeMetricsReporterParams,
315 > ) *metrics.RuntimeMetricsReporter { fx.go ×44
316 > return metrics.NewRuntimeMetricsReporter(
317 > params.MetricHandler,
318 > time.Minute,
319 > params.Logger,
320 > string(params.InstanceID),
321 > )
322 > }
323
324 > func HistoryRawClientProvider(clientBean client.Bean) HistoryRawClient { fx.go ×44
325 > return clientBean.GetHistoryClient()
326 > }
327
328 > func HistoryClientProvider(historyRawClient HistoryRawClient, dc *dynamicconfig.Collection) HistoryClient { fx.go ×44
329 > return history.NewRetryableClient(
330 > historyRawClient,
331 > common.CreateHistoryClientRetryPolicy(dynamicconfig.RetryUnboundedOnSystemResourceExhausted.Get(dc)),
332 > common.IsServiceClientTransientError,
333 > )
334 > }
335
336 func MatchingRawClientProvider(
337 clientBean client.Bean,
338 namespaceRegistry namespace.Registry,
339 > ) (MatchingRawClient, error) { fx.go ×44
340 > return clientBean.GetMatchingClient(namespaceRegistry.GetNamespaceName)
341 > }
342
343 > func MatchingClientProvider(matchingRawClient MatchingRawClient, dc *dynamicconfig.Collection) MatchingClient { fx.go ×44
344 > return matching.NewRetryableClient(
345 > matchingRawClient,
346 > common.CreateMatchingClientRetryPolicy(dynamicconfig.RetryUnboundedOnSystemResourceExhausted.Get(dc)),
347 > common.CreateMatchingClientLongPollRetryPolicy(),
348 > common.IsServiceClientTransientError,
349 > )
350 > }
351
352 > func PersistenceConfigProvider(persistenceConfig config.Persistence, dc *dynamicconfig.Collection) *config.Persistence { fx.go ×44
353 > persistenceConfig.TransactionSizeLimit = dynamicconfig.TransactionSizeLimit.Get(dc)
354 > return &persistenceConfig
355 > }
356
357 > func ArchivalMetadataProvider(dc *dynamicconfig.Collection, cfg *config.Config) archiver.ArchivalMetadata { fx.go ×44
358 > return archiver.NewArchivalMetadata(
359 > dc,
360 > cfg.Archival.History.State,
361 > cfg.Archival.History.EnableRead,
362 > cfg.Archival.Visibility.State,
363 > cfg.Archival.Visibility.EnableRead,
364 > &cfg.NamespaceDefaults.Archival,
365 > )
366 > }
367
368 func ArchiverProviderProvider(
369 cfg *config.Config,
370 customHistoryArchiverFactory provider.CustomHistoryArchiverFactory,
371 customVisibilityArchiverFactory provider.CustomVisibilityArchiverFactory,
372 persistenceExecutionManager persistence.ExecutionManager,
373 logger log.SnTaggedLogger,
374 metricsHandler metrics.Handler,
375 > ) provider.ArchiverProvider { fx.go ×44
376 > return provider.NewArchiverProvider(
377 > cfg.Archival.History.Provider,
378 > cfg.Archival.Visibility.Provider,
379 > customHistoryArchiverFactory,
380 > customVisibilityArchiverFactory,
381 > persistenceExecutionManager,
382 > logger,
383 > metricsHandler,
384 > )
385 > }
386
387 func SdkClientFactoryProvider(
388 cfg *config.Config,
389 tlsConfigProvider encryption.TLSConfigProvider,
390 metricsHandler metrics.Handler,
391 logger log.SnTaggedLogger,
392 resolver *membership.GRPCResolver,
393 dc *dynamicconfig.Collection,
394 > ) (sdk.ClientFactory, error) { fx.go ×44
395 > frontendURL, _, _, frontendTLSConfig, err := getFrontendConnectionDetails(cfg, tlsConfigProvider, resolver)
396 > if err != nil {
397 return nil, err
398 }
399 > return sdk.NewClientFactory( fx.go ×44
400 > frontendURL,
401 > frontendTLSConfig,
402 > metricsHandler,
403 > logger,
404 > dynamicconfig.WorkerStickyCacheSize.Get(dc),
405 > ), nil
406 }
407
408 > func DCRedirectionPolicyProvider(cfg *config.Config) config.DCRedirectionPolicy { fx.go ×44
409 > return cfg.DCRedirectionPolicy
410 > }
411
412 func PerServiceDialOptionsProvider(
413 logger log.SnTaggedLogger,
414 > ) map[primitives.ServiceName][]grpc.DialOption { fx.go ×44
415 > trailerInterceptor := interceptor.TrailerToContextMetadataInterceptor(logger)
416 > dialOpt := grpc.WithChainUnaryInterceptor(trailerInterceptor)
417 > return map[primitives.ServiceName][]grpc.DialOption{
418 > primitives.HistoryService: {dialOpt},
419 > primitives.MatchingService: {dialOpt},
420 > }
421 > }
422
423 func RPCFactoryProvider(
424 cfg *config.Config,
425 svcName primitives.ServiceName,
426 logger log.Logger,
427 metricsHandler metrics.Handler,
428 tlsConfigProvider encryption.TLSConfigProvider,
429 resolver *membership.GRPCResolver,
430 tracingStatsHandler telemetry.ClientStatsHandler,
431 perServiceDialOptions map[primitives.ServiceName][]grpc.DialOption,
432 monitor membership.Monitor,
433 dc *dynamicconfig.Collection,
434 tokenProvider auth.TokenProvider,
435 > ) (common.RPCFactory, error) { fx.go ×44
436 > frontendURL, frontendHTTPURL, frontendHTTPPort, frontendTLSConfig, err := getFrontendConnectionDetails(cfg, tlsConfigProvider, resolver)
437 > if err != nil {
438 return nil, err
439 }
440
441 > var options []grpc.DialOption fx.go ×44
442 > if tracingStatsHandler != nil {
443 > options = append(options, grpc.WithStatsHandler(tracingStatsHandler)) data_store_factory.go ×29
444 > }
445 > enableServerKeepalive := dynamicconfig.EnableInternodeServerKeepAlive.Get(dc)() fx.go ×44
446 > enableClientKeepalive := dynamicconfig.EnableInternodeClientKeepAlive.Get(dc)()
447 > factory := rpc.NewFactory(
448 > cfg,
449 > svcName,
450 > logger,
451 > metricsHandler,
452 > tlsConfigProvider,
453 > frontendURL,
454 > frontendHTTPURL,
455 > frontendHTTPPort,
456 > frontendTLSConfig,
457 > options,
458 > perServiceDialOptions,
459 > monitor,
460 > tokenProvider,
461 > )
462 > factory.EnableInternodeServerKeepalive = enableServerKeepalive
463 > factory.EnableInternodeClientKeepalive = enableClientKeepalive
464 > logger.Debug(fmt.Sprintf("RPC factory created. enableServerKeepalive: %v, enableClientKeepalive: %v", enableServerKeepalive, enableClientKeepalive))
465 > return factory, nil
466 }
467
468 func FrontendHTTPClientCacheProvider(
469 metadata cluster.Metadata,
470 tlsConfigProvider encryption.TLSConfigProvider,
471 > ) *cluster.FrontendHTTPClientCache { fx.go ×44
472 > return cluster.NewFrontendHTTPClientCache(metadata, tlsConfigProvider)
473 > }
474
475 func getFrontendConnectionDetails(
476 cfg *config.Config,
477 tlsConfigProvider encryption.TLSConfigProvider,
478 resolver *membership.GRPCResolver,
479 > ) (string, string, int, *tls.Config, error) { fx.go ×44
480 > // To simplify the static config, we switch default values based on whether the config
481 > // defines an "internal-frontend" service. The default for TLS config can be overridden
482 > // with publicClient.forceTLSConfig.
483 > _, hasIFE := cfg.Services[string(primitives.InternalFrontendService)]
484 >
485 > forceTLS := cfg.PublicClient.ForceTLSConfig
486 > if forceTLS == config.ForceTLSConfigAuto {
487 > if hasIFE {
488 forceTLS = config.ForceTLSConfigInternode
489 > } else { fx.go ×44
490 > forceTLS = config.ForceTLSConfigFrontend
491 > }
492 }
493
494 > var frontendTLSConfig *tls.Config fx.go ×44
495 > var err error
496 > switch forceTLS {
497 case config.ForceTLSConfigInternode:
498 frontendTLSConfig, err = tlsConfigProvider.GetInternodeClientConfig()
499 > case config.ForceTLSConfigFrontend: fx.go ×44
500 > frontendTLSConfig, err = tlsConfigProvider.GetFrontendClientConfig()
501 default:
502 err = fmt.Errorf("invalid forceTLSConfig")
503 }
504 > if err != nil { fx.go ×44
505 return "", "", 0, nil, fmt.Errorf("unable to load TLS configuration: %w", err)
506 }
507
508 > frontendURL := cfg.PublicClient.HostPort fx.go ×44
509 > if frontendURL == "" {
510 > if hasIFE { request_response.pb.go ×12
511 frontendURL = resolver.MakeURL(primitives.InternalFrontendService)
513 > frontendURL = resolver.MakeURL(primitives.FrontendService)
514 > }
515 }
516 > frontendHTTPURL := cfg.PublicClient.HTTPHostPort fx.go ×44
517 > if frontendHTTPURL == "" {
518 > if hasIFE {
519 frontendHTTPURL = resolver.MakeURL(primitives.InternalFrontendService)
520 > } else { fx.go ×44
521 > frontendHTTPURL = resolver.MakeURL(primitives.FrontendService)
522 > }
523 }
524
525 > var frontendHTTPPort int fx.go ×44
526 > if hasIFE {
527 frontendHTTPPort = cfg.Services[string(primitives.InternalFrontendService)].RPC.HTTPPort
528 > } else { fx.go ×44
529 > frontendHTTPPort = cfg.Services[string(primitives.FrontendService)].RPC.HTTPPort
530 > }
531
532 > return frontendURL, frontendHTTPURL, frontendHTTPPort, frontendTLSConfig, nil fx.go ×44
533 }