fx.go ×44

Frontier kind: Code frontier

unlabeled · c_b6dceff7e970

11 tests · 20600 LOC · 595 files · introduces 0 tests · 705 LOC · 31 files

Introduces — evidence that enters the hierarchy at this concept

Code
168 ranges705 lines · 31 files
Tests
0 tests

Contains — complete concept membership

All code (extent)
3969 ranges20600 lines · 595 files · Browse complete extent
All tests (intent)
11 testsBrowse complete intent

Neighbourhood graph

The orange circle is the focus. Violet and green circles are every ancestor and descendant, broader and narrower, at any distance; blue squares and pink diamonds are the introduced files and exact introduced tests of every visible concept, not only the focus's. Arrows point from broader to narrower concepts and bridge only concepts omitted from this view. Undirected links show source or test introduction. Concept and file size follows LOC; exact test nodes use test-count units.

Introduced files, introduced tests, and structurally relevant concept specialization

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 native relationship evidence on this page.

Graph controls are ready.

Interactive rendering requires JavaScript and WebGL. Use the native relationship evidence on this page while the interactive map is unavailable.

Native relationship evidence

Every exact file and test below is linked only from the concept that introduces it.

Introduced tests

Every collected test enters the hierarchy at exactly one concept.

No tests are introduced at this concept. Its intent tests are introduced by other concepts.

Introduced code

Every collected source range enters the hierarchy at exactly one concept.

Showing the top 20 of 31 files by introduced lines: 683 of 705 introduced LOC and 157 of 168 ranges. Expand a file to inspect source; the > gutter marks introduced lines.

go.temporal.io/server/temporal/fx.go 243 introduced LOC · 44 ranges

Open complete file

158 )
159
160 > func NewServerFx(topLevelModule fx.Option, opts ...ServerOption) (*ServerFx, error) { fx.go
161 > var s ServerFx
162 > s.app = fx.New(
163 > topLevelModule,
164 > fx.Supply(opts),
165 > fx.Populate(&s.startupSynchronizationMode),
166 > fx.Populate(&s.logger),
167 > )
168 > if err := s.app.Err(); err != nil {
169 return nil, err
170 }
171 > return &s, nil fx.go
172 }
173
174 > func ServerOptionsProvider(opts []ServerOption) (serverOptionsProvider, error) { fx.go
175 > so := newServerOptions(opts)
176 >
177 > err := so.loadAndValidate()
178 > if err != nil {
179 return serverOptionsProvider{}, err
180 }
181
182 // Logger
183 > logger := so.logger fx.go
184 > if logger == nil {
185 logger = log.NewZapLogger(log.BuildZapLogger(so.config.Log))
186 }
187
188 > persistenceConfig := so.config.Persistence fx.go
189 > err = verifyPersistenceCompatibleVersion(persistenceConfig, so.persistenceServiceResolver, logger)
190 > if err != nil {
191 return serverOptionsProvider{}, err
192 }
193
194 > stopChan := make(chan any) fx.go
195 >
196 > // ClientFactoryProvider
197 > clientFactoryProvider := so.clientFactoryProvider
198 > if clientFactoryProvider == nil {
199 > clientFactoryProvider = client.NewFactoryProvider()
200 > }
201
202 // MetricsHandler
203 > metricHandler := so.metricHandler fx.go
204 > if metricHandler == nil {
205 > metricHandler, err = metrics.MetricsHandlerFromConfig(logger, so.config.Global.Metrics)
206 > if err != nil {
207 return serverOptionsProvider{}, fmt.Errorf("unable to create metrics handler: %w", err)
208 }
212 // if injected, else a no-op provider that discards events. A deployment opts in by injecting a
213 // provider via WithCustomEventLoggerProvider.
214 > eventLoggerProvider := so.eventLoggerProvider fx.go
215 > if eventLoggerProvider == nil {
216 > eventLoggerProvider = lognoop.NewLoggerProvider()
217 > }
218
219 // DynamicConfigClient
220 > dcClient := so.dynamicConfigClient fx.go
221 > if dcClient == nil {
222 dcConfig := so.config.DynamicConfigClient
223 if dcConfig != nil {
234
235 // TLSConfigProvider
236 > tlsConfigProvider := so.tlsConfigProvider fx.go
237 > if tlsConfigProvider == nil {
238 > tlsConfigProvider, err = encryption.NewTLSConfigProviderFromConfig(so.config.Global.TLS, metricHandler, logger, nil)
239 > if err != nil {
240 return serverOptionsProvider{}, err
241 }
243
244 // EsConfig / EsClient
245 > var esConfig *esclient.Config fx.go
246 > var esClient esclient.Client
247 >
248 > if persistenceConfig.SecondaryVisibilityConfigExist() &&
249 > persistenceConfig.DataStores[persistenceConfig.SecondaryVisibilityStore].Elasticsearch != nil {
250 esConfig = persistenceConfig.DataStores[persistenceConfig.SecondaryVisibilityStore].Elasticsearch
251 esConfig.SetHttpClient(so.elasticsearchHttpClient)
252 }
253 > if persistenceConfig.VisibilityConfigExist() && fx.go
254 > persistenceConfig.DataStores[persistenceConfig.VisibilityStore].Elasticsearch != nil {
255 esConfig = persistenceConfig.DataStores[persistenceConfig.VisibilityStore].Elasticsearch
256 esConfig.SetHttpClient(so.elasticsearchHttpClient)
257 }
258
259 > if esConfig != nil { fx.go
260 esHttpClient := so.elasticsearchHttpClient
261 if esHttpClient == nil {
275
276 // check that when static hosts are defined, they are defined for all required hosts
277 > if len(so.hostsByService) > 0 { fx.go
278 for _, service := range DefaultServices {
279 hosts := so.hostsByService[primitives.ServiceName(service)]
284 }
285
286 > if so.config.Global.Authorization.RemoteClusterAuth.Require && so.tokenProvider == nil { fx.go
287 return serverOptionsProvider{}, errors.New("global.authorization.remoteClusterAuth.require is true but no TokenProvider is configured: use WithTokenProvider")
288 }
291 // Coarse check: any remote-cluster TLS entry passes; per-hostname config is still validated
292 // lazily on first dial.
293 > if so.tokenProvider != nil && so.tlsConfigProvider == nil && len(so.config.Global.TLS.RemoteClusters) == 0 { fx.go
294 return serverOptionsProvider{}, errors.New("WithTokenProvider is set but no remote-cluster TLS is configured: supply global.tls.remoteClusters in config, or pass a provider via WithTLSConfigProvider")
295 }
296
297 > return serverOptionsProvider{ fx.go
298 > ServerOptions: so,
299 > StopChan: stopChan,
300 > StartupSynchronizationMode: so.startupSynchronizationMode,
301 >
302 > Config: so.config,
303 > PProfConfig: &so.config.Global.PProf,
304 > LogConfig: so.config.Log,
305 >
306 > ServiceNames: so.serviceNames,
307 > ServiceHosts: so.hostsByService,
308 > NamespaceLogger: so.namespaceLogger,
309 >
310 > ServiceResolver: so.persistenceServiceResolver,
311 > CustomDataStoreFactory: so.customDataStoreFactory,
312 > CustomVisibilityStore: so.customVisibilityStoreFactory,
313 > CustomHistoryArchiverFactory: so.customHistoryArchiverFactory,
314 > CustomVisibilityArchiverFactory: so.customVisibilityArchiverFactory,
315 >
316 > SearchAttributesMapper: so.searchAttributesMapper,
317 > CustomFrontendInterceptors: so.customFrontendInterceptors,
318 > Authorizer: so.authorizer,
319 > ClaimMapper: so.claimMapper,
320 > AudienceGetter: so.audienceGetter,
321 > TokenProvider: so.tokenProvider,
322 >
323 > Logger: logger,
324 > ClientFactoryProvider: clientFactoryProvider,
325 > DynamicConfigClient: dcClient,
326 > TLSConfigProvider: tlsConfigProvider,
327 > EsClient: esClient,
328 > MetricsHandler: metricHandler,
329 > EventLoggerProvider: eventLoggerProvider,
330 > }, nil
331 }
332
333 // Start temporal server.
334 // This function should be called only once, Server doesn't support multiple restarts.
335 > func (s *ServerFx) Start() error { fx.go
336 > err := s.app.Start(context.Background())
337 > if err != nil {
338 return err
339 }
340
341 > if s.startupSynchronizationMode.blockingStart { fx.go
342 // If s.so.interruptCh is nil this will wait forever.
343 interruptSignal := <-s.startupSynchronizationMode.interruptCh
346 }
347
348 > return nil fx.go
349 }
350
405 // into fx providers here. Essentially, we want an `fx.In` object in the server graph, and an `fx.Out` object in the
406 // service graphs. This is a workaround to achieve something similar.
407 > func (params ServiceProviderParamsCommon) GetCommonServiceOptions(serviceName primitives.ServiceName) fx.Option { fx.go
408 > membershipModule := ringpop.MembershipModule
409 > if len(params.StaticServiceHosts) > 0 {
410 membershipModule = static.MembershipModule(params.StaticServiceHosts)
411 }
412
413 > return fx.Options( fx.go
414 > fx.Supply(
415 > serviceName,
416 > params.PersistenceConfig,
417 > params.ClusterMetadata,
418 > params.Cfg,
419 > params.SpanExporters,
420 > ),
421 > fx.Provide(
422 > resource.DefaultSnTaggedLoggerProvider,
423 > params.PersistenceFactoryProvider,
424 > func() persistenceClient.AbstractDataStoreFactory {
425 > return params.DataStoreFactory
426 > },
427 > func() visibility.VisibilityStoreFactory {
428 > return params.VisibilityStoreFactory
429 > },
430 > func() provider.CustomHistoryArchiverFactory {
431 > return params.CustomHistoryArchiverFactory
432 > },
433 > func() provider.CustomVisibilityArchiverFactory {
434 > return params.CustomVisibilityArchiverFactory
435 > },
436 > func() client.FactoryProvider {
437 > return params.ClientFactoryProvider
438 > },
439 > func() authorization.JWTAudienceMapper {
440 > return params.AudienceGetter
441 > },
442 > func() resolver.ServiceResolver {
443 > return params.PersistenceServiceResolver
444 > },
445 > func() searchattribute.Mapper {
446 > return params.SearchAttributesMapper
447 > },
448 > func() authorization.Authorizer {
449 > return params.Authorizer
450 > },
451 func() authorization.ClaimMapper {
452 return params.ClaimMapper
453 },
454 > func() auth.TokenProvider { fx.go
455 > return params.TokenProvider
456 > },
457 > func() encryption.TLSConfigProvider {
458 > return params.TlsConfigProvider
459 > },
460 > func() dynamicconfig.Client {
461 > return params.DynamicConfigClient
462 > },
463 > func() log.Logger {
464 > return params.Logger
465 > },
466 > func() metrics.Handler {
467 > return params.MetricsHandler.WithTags(metrics.ServiceNameTag(serviceName))
468 > },
469 > func() otellog.Logger {
470 > return wideevents.NewLogger(params.EventLoggerProvider, string(serviceName))
471 > },
472 func() esclient.Client {
473 return params.EsClient
474 },
475 > func() resource.NamespaceLogger { fx.go
476 > return params.NamespaceLogger
477 > },
478 > func() tasks.TaskCategoryRegistry {
479 > return params.TaskCategoryRegistry
480 > },
481 ),
482 ServiceTracingModule,
503 }
504
505 > func NewService(app *fx.App, serviceName primitives.ServiceName, logger log.Logger) ServicesGroupOut { fx.go
506 > return ServicesGroupOut{
507 > Services: &ServicesMetadata{
508 > app: app,
509 > serviceName: serviceName,
510 > logger: logger,
511 > },
512 > }
513 > }
514
515 func HistoryServiceProvider(
516 params ServiceProviderParamsCommon,
517 > ) (ServicesGroupOut, error) { fx.go
518 > serviceName := primitives.HistoryService
519 >
520 > if _, ok := params.ServiceNames[serviceName]; !ok {
521 params.Logger.Info("Service is not requested, skipping initialization.", tag.Service(serviceName))
522 return ServicesGroupOut{}, nil
523 }
524
525 > app := fx.New( fx.go
526 > params.GetCommonServiceOptions(serviceName),
527 > history.QueueModule,
528 > history.Module,
529 > replication.Module,
530 > )
531 >
532 > return NewService(app, serviceName, params.Logger), app.Err()
533 }
534
535 func MatchingServiceProvider(
536 params ServiceProviderParamsCommon,
537 > ) (ServicesGroupOut, error) { fx.go
538 > serviceName := primitives.MatchingService
539 >
540 > if _, ok := params.ServiceNames[serviceName]; !ok {
541 params.Logger.Info("Service is not requested, skipping initialization.", tag.Service(serviceName))
542 return ServicesGroupOut{}, nil
543 }
544
545 > app := fx.New( fx.go
546 > params.GetCommonServiceOptions(serviceName),
547 > matching.Module,
548 > )
549 >
550 > return NewService(app, serviceName, params.Logger), app.Err()
551 }
552
553 func FrontendServiceProvider(
554 params ServiceProviderParamsCommon,
555 > ) (ServicesGroupOut, error) { fx.go
556 > return genericFrontendServiceProvider(params, primitives.FrontendService)
557 > }
558
559 func InternalFrontendServiceProvider(
560 params ServiceProviderParamsCommon,
561 > ) (ServicesGroupOut, error) { fx.go
562 > return genericFrontendServiceProvider(params, primitives.InternalFrontendService)
563 > }
564
565 func genericFrontendServiceProvider(
566 params ServiceProviderParamsCommon,
567 serviceName primitives.ServiceName,
568 > ) (ServicesGroupOut, error) { fx.go
569 > if _, ok := params.ServiceNames[serviceName]; !ok {
570 > params.Logger.Info("Service is not requested, skipping initialization.", tag.Service(serviceName))
571 > return ServicesGroupOut{}, nil
572 > }
573
574 > app := fx.New( fx.go
575 > params.GetCommonServiceOptions(serviceName),
576 > fx.Supply(params.CustomFrontendInterceptors),
577 > fx.Supply([]grpc.StreamServerInterceptor{}),
578 > fx.Decorate(func() authorization.ClaimMapper {
579 > switch serviceName {
580 > case primitives.FrontendService:
581 > return params.ClaimMapper
582 case primitives.InternalFrontendService:
583 return authorization.NewInternalClaimMapper()
586 }
587 }),
588 > fx.Decorate(func() log.SnTaggedLogger { fx.go
589 > // Use "frontend" for logs even if serviceName is "internal-frontend", but add an
590 > // extra tag to differentiate.
591 > tags := []tag.Tag{tag.Service(primitives.FrontendService)}
592 > if serviceName == primitives.InternalFrontendService {
593 tags = append(tags, tag.Bool("internal-frontend", true))
594 }
595 > return log.With(params.Logger, tags...) fx.go
596 }),
597 frontend.Module,
598 )
599
600 > return NewService(app, serviceName, params.Logger), app.Err() fx.go
601 }
602
603 func WorkerServiceProvider(
604 params ServiceProviderParamsCommon,
605 > ) (ServicesGroupOut, error) { fx.go
606 > serviceName := primitives.WorkerService
607 >
608 > if _, ok := params.ServiceNames[serviceName]; !ok {
609 params.Logger.Info("Service is not requested, skipping initialization.", tag.Service(serviceName))
610 return ServicesGroupOut{}, nil
611 }
612
613 > app := fx.New( fx.go
614 > params.GetCommonServiceOptions(serviceName),
615 > worker.Module,
616 > )
617 >
618 > return NewService(app, serviceName, params.Logger), app.Err()
619 }
620
728 logger,
729 )
730 > case *serviceerror.NotFound: fx.go
731 > // Initialize current cluster record
732 > if initErr := initCurrentClusterMetadataRecord(
733 > ctx,
734 > clusterMetadataManager,
735 > svc,
736 > indexSearchAttributes,
737 > logger,
738 > ); initErr != nil {
739 return svc.ClusterMetadata, svc.Persistence, initErr
740 }
916 }
917
918 > func PersistenceFactoryProvider() persistenceClient.FactoryProviderFn { fx.go
919 > return persistenceClient.FactoryProvider
920 > }
921
922 func ServerLifetimeHooks(
923 lc fx.Lifecycle,
924 svr *ServerImpl,
925 > ) { fx.go
926 > lc.Append(fx.StartStopHook(svr.Start, svr.Stop))
927 > }
928
929 func verifyPersistenceCompatibleVersion(
931 persistenceServiceResolver resolver.ServiceResolver,
932 logger log.Logger,
933 > ) error { fx.go
934 > // cassandra schema version validation
935 > if err := cassandra.VerifyCompatibleVersion(cfg, persistenceServiceResolver, logger); err != nil {
936 return fmt.Errorf("cassandra schema version compatibility check failed: %w", err)
937 }
938 // sql schema version validation
939 > if err := sql.VerifyCompatibleVersion(cfg, persistenceServiceResolver, logger); err != nil { fx.go
940 return fmt.Errorf("sql schema version compatibility check failed: %w", err)
941 }
942 > return nil fx.go
943 }
944
go.temporal.io/server/common/resource/fx.go 98 introduced LOC · 20 ranges

Open complete file

134 logger log.SnTaggedLogger,
135 fn ThrottledLoggerRpsFn,
136 > ) log.ThrottledLogger { fx.go
137 > return log.NewThrottledLogger(
138 > logger,
139 > quotas.RateFn(fn),
140 > )
141 > }
142
143 func GrpcListenerProvider(factory common.RPCFactory) net.Listener {
355 }
356
357 > func ArchivalMetadataProvider(dc *dynamicconfig.Collection, cfg *config.Config) archiver.ArchivalMetadata { fx.go
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(
373 logger log.SnTaggedLogger,
374 metricsHandler metrics.Handler,
375 > ) provider.ArchiverProvider { fx.go
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(
392 resolver *membership.GRPCResolver,
393 dc *dynamicconfig.Collection,
394 > ) (sdk.ClientFactory, error) { fx.go
395 > frontendURL, _, _, frontendTLSConfig, err := getFrontendConnectionDetails(cfg, tlsConfigProvider, resolver)
396 > if err != nil {
397 return nil, err
398 }
399 > return sdk.NewClientFactory( fx.go
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
409 > return cfg.DCRedirectionPolicy
410 > }
411
412 func PerServiceDialOptionsProvider(
433 dc *dynamicconfig.Collection,
434 tokenProvider auth.TokenProvider,
435 > ) (common.RPCFactory, error) { fx.go
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
442 > if tracingStatsHandler != nil {
443 options = append(options, grpc.WithStatsHandler(tracingStatsHandler))
444 }
445 > enableServerKeepalive := dynamicconfig.EnableInternodeServerKeepAlive.Get(dc)() fx.go
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
477 tlsConfigProvider encryption.TLSConfigProvider,
478 resolver *membership.GRPCResolver,
479 > ) (string, string, int, *tls.Config, error) { fx.go
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
490 > forceTLS = config.ForceTLSConfigFrontend
491 > }
492 }
493
494 > var frontendTLSConfig *tls.Config fx.go
495 > var err error
496 > switch forceTLS {
497 case config.ForceTLSConfigInternode:
498 frontendTLSConfig, err = tlsConfigProvider.GetInternodeClientConfig()
499 > case config.ForceTLSConfigFrontend: fx.go
500 > frontendTLSConfig, err = tlsConfigProvider.GetFrontendClientConfig()
501 default:
502 err = fmt.Errorf("invalid forceTLSConfig")
503 }
504 > if err != nil { fx.go
505 return "", "", 0, nil, fmt.Errorf("unable to load TLS configuration: %w", err)
506 }
507
508 > frontendURL := cfg.PublicClient.HostPort fx.go
509 > if frontendURL == "" {
510 if hasIFE {
511 frontendURL = resolver.MakeURL(primitives.InternalFrontendService)
514 }
515 }
516 > frontendHTTPURL := cfg.PublicClient.HTTPHostPort fx.go
517 > if frontendHTTPURL == "" {
518 > if hasIFE {
519 frontendHTTPURL = resolver.MakeURL(primitives.InternalFrontendService)
520 > } else { fx.go
521 > frontendHTTPURL = resolver.MakeURL(primitives.FrontendService)
522 > }
523 }
524
525 > var frontendHTTPPort int fx.go
526 > if hasIFE {
527 frontendHTTPPort = cfg.Services[string(primitives.InternalFrontendService)].RPC.HTTPPort
528 > } else { fx.go
529 > frontendHTTPPort = cfg.Services[string(primitives.FrontendService)].RPC.HTTPPort
530 > }
531
532 > return frontendURL, frontendHTTPURL, frontendHTTPPort, frontendTLSConfig, nil fx.go
533 }
go.temporal.io/server/temporal/server_impl.go 85 introduced LOC · 10 ranges

Open complete file

65 metricsHandler metrics.Handler,
66 serializer serialization.Serializer,
67 > ) *ServerImpl { server_impl.go
68 > s := &ServerImpl{
69 > so: opts,
70 > stoppedCh: stoppedCh,
71 > logger: logger,
72 > namespaceLogger: namespaceLogger,
73 > persistenceConfig: persistenceConfig,
74 > clusterMetadata: clusterMetadata,
75 > persistenceFactoryProvider: persistenceFactoryProvider,
76 > metricsHandler: metricsHandler,
77 > }
78 > for _, svcMeta := range servicesGroup.Services {
79 > if svcMeta != nil {
80 > s.servicesMetadata = append(s.servicesMetadata, svcMeta)
81 > }
82 }
83 // Store serializer for use in Start()
84 > s.serializer = serializer server_impl.go
85 > return s
86 }
87
88 > func (s *ServerImpl) Start(ctx context.Context) error { server_impl.go
89 > s.logger.Info("Starting server for services", tag.Value(s.so.serviceNames))
90 > s.logger.Debug(s.so.config.String())
91 >
92 > if err := initSystemNamespaces(
93 > ctx,
94 > &s.persistenceConfig,
95 > s.clusterMetadata.CurrentClusterName,
96 > s.so.persistenceServiceResolver,
97 > s.persistenceFactoryProvider,
98 > s.logger,
99 > s.so.customDataStoreFactory,
100 > s.metricsHandler,
101 > s.serializer,
102 > ); err != nil {
103 return fmt.Errorf("unable to initialize system namespace: %w", err)
104 }
105
106 > return s.startServices() server_impl.go
107 }
108
124 }
125
126 > func (s *ServerImpl) startServices() error { server_impl.go
127 > // The membership join time may exceed the configured max join duration.
128 > // Double the service start timeout to make sure there is enough time for start logic.
129 > timeout := max(serviceStartTimeout, 2*s.so.config.Global.Membership.MaxJoinDuration)
130 > ctx, cancel := context.WithTimeout(context.Background(), timeout)
131 > defer cancel()
132 >
133 > svcs := slices.Clone(s.servicesMetadata)
134 > slices.SortFunc(svcs, func(a, b *ServicesMetadata) int {
135 > return cmp.Compare(initOrder[a.serviceName], initOrder[b.serviceName])
136 > })
137
138 > var allErrs error server_impl.go
139 > for _, svc := range svcs {
140 > err := svc.app.Start(ctx)
141 > if err != nil {
142 allErrs = multierr.Append(allErrs, fmt.Errorf("failed to start service %v: %w", svc.serviceName, err))
143 }
144 }
145 > return allErrs server_impl.go
146 }
147
156 metricsHandler metrics.Handler,
157 serializer serialization.Serializer,
158 > ) error { server_impl.go
159 > clusterName := persistenceClient.ClusterName(currentClusterName)
160 > metricsHandler = metricsHandler.WithTags(metrics.ServiceNameTag(primitives.ServerService))
161 > dataStoreFactory := persistenceClient.DataStoreFactoryProvider(
162 > clusterName,
163 > persistenceServiceResolver,
164 > cfg,
165 > customDataStoreFactory,
166 > logger,
167 > metricsHandler,
168 > telemetry.NoopTracerProvider,
169 > serializer,
170 > )
171 > factory := persistenceFactoryProvider(persistenceClient.NewFactoryParams{
172 > DataStoreFactory: dataStoreFactory,
173 > Cfg: cfg,
174 > PersistenceMaxQPS: nil,
175 > PersistenceNamespaceMaxQPS: nil,
176 > ClusterName: persistenceClient.ClusterName(currentClusterName),
177 > MetricsHandler: metricsHandler,
178 > Logger: logger,
179 > Serializer: serializer,
180 > })
181 > defer factory.Close()
182 >
183 > metadataManager, err := factory.NewMetadataManager()
184 > if err != nil {
185 return fmt.Errorf("unable to initialize metadata manager: %w", err)
186 }
187 > defer metadataManager.Close() server_impl.go
188 > ctx, cancel := context.WithTimeout(
189 > headers.SetCallerInfo(ctx, headers.SystemBackgroundHighCallerInfo),
190 > 30*time.Second,
191 > )
192 > defer cancel()
193 >
194 > if err = metadataManager.InitializeSystemNamespaces(ctx, currentClusterName); err != nil {
195 return fmt.Errorf("unable to register system namespace: %w", err)
196 }
197 > return nil server_impl.go
198 }
go.temporal.io/server/common/membership/ringpop/factory.go 54 introduced LOC · 12 ranges

Open complete file

88
89 // getMonitor returns a membership monitor
90 > func (factory *factory) getMonitor() *monitor { factory.go
91 > factory.monOnce.Do(func() {
92 > ctx, cancel := context.WithTimeout(context.Background(), persistenceOperationTimeout)
93 > defer cancel()
94 >
95 > ctx = headers.SetCallerInfo(ctx, headers.SystemBackgroundHighCallerInfo)
96 > currentClusterMetadata, err := factory.MetadataManager.GetCurrentClusterMetadata(ctx)
97 > if err != nil {
98 factory.Logger.Fatal("Failed to get current cluster ID", tag.Error(err))
99 }
100
101 > appName := "temporal" factory.go
102 > if currentClusterMetadata.UseClusterIdMembership {
103 > appName = fmt.Sprintf("temporal-%s", currentClusterMetadata.GetClusterId())
104 > }
105 > rp, err := ringpop.New(appName, ringpop.Channel(factory.getTChannel()), ringpop.AddressResolverFunc(factory.broadcastAddressResolver))
106 > if err != nil {
107 factory.Logger.Fatal("Failed to get new ringpop", tag.Error(err))
108 }
110 // Empirically, ringpop updates usually propagate in under a second even in relatively large clusters.
111 // 3 seconds is an over-estimate to be safer.
112 > maxPropagationTime := dynamicconfig.RingpopApproximateMaxPropagationTime.Get(factory.DC)() factory.go
113 > replicaPoints := dynamicconfig.RingpopReplicaPoints.Get(factory.DC)()
114 >
115 > factory.monitor = newMonitor(
116 > factory.ServiceName,
117 > factory.ServicePortMap,
118 > rp,
119 > factory.Logger,
120 > factory.MetadataManager,
121 > factory.broadcastAddressResolver,
122 > factory.Config.MaxJoinDuration,
123 > maxPropagationTime,
124 > factory.getJoinTime(maxPropagationTime),
125 > replicaPoints,
126 > )
127 })
128
129 > return factory.monitor factory.go
130 }
131
132 > func (factory *factory) getJoinTime(maxPropagationTime time.Duration) time.Time { factory.go
133 > var alignTime time.Duration
134 > switch factory.ServiceName {
135 > case primitives.MatchingService:
136 > alignTime = dynamicconfig.MatchingAlignMembershipChange.Get(factory.DC)()
137 > case primitives.HistoryService:
138 > alignTime = dynamicconfig.HistoryAlignMembershipChange.Get(factory.DC)()
139 }
140 > if alignTime == 0 { factory.go
141 > return time.Time{}
142 > }
143 return util.NextAlignedTime(time.Now().Add(maxPropagationTime), alignTime)
144 }
145
146 > func (factory *factory) broadcastAddressResolver() (string, error) { factory.go
147 > return buildBroadcastHostPort(factory.getTChannel().PeerInfo(), factory.Config.BroadcastAddress)
148 > }
149
150 func (factory *factory) getTChannel() *tchannel.Channel {
218
219 if factory.RPCConfig.BindOnLocalHost {
220 > return net.ParseIP(environment.GetLocalhostIP()) factory.go
221 > }
222
223 if len(factory.RPCConfig.BindOnIP) > 0 {
247 }
248
249 > func (factory *factory) getHostInfoProvider() (membership.HostInfoProvider, error) { factory.go
250 > address, err := factory.broadcastAddressResolver()
251 > if err != nil {
252 return nil, err
253 }
254
255 > servicePort, ok := factory.ServicePortMap[factory.ServiceName] factory.go
256 > if !ok {
257 return nil, membership.ErrUnknownService
258 }
261 // ringpop messages. We use a different port for the service, so we
262 // replace that portion.
263 > serviceAddress, err := replaceServicePort(address, servicePort) factory.go
264 > if err != nil {
265 return nil, err
266 }
267
268 > hostInfo := membership.NewHostInfoFromAddress(serviceAddress) factory.go
269 > return membership.NewHostInfoProvider(hostInfo), nil
270 }
go.temporal.io/server/common/config/persistence.go 33 introduced LOC · 15 ranges

Open complete file

33
34 // Validate validates the persistence config
35 > func (c *Persistence) Validate() error { persistence.go
36 > stores := []string{c.DefaultStore}
37 > if c.VisibilityStore != "" {
38 > stores = append(stores, c.VisibilityStore)
39 > }
40 > if c.SecondaryVisibilityStore != "" {
41 stores = append(stores, c.SecondaryVisibilityStore)
42 }
57 // - visibilityStore (es), secondaryVisibilityStore (advanced sql)
58
59 > if c.VisibilityStore == "" { persistence.go
60 return fmt.Errorf("%w: visibilityStore must be specified", ErrPersistenceConfig)
61 }
62 > if c.SecondaryVisibilityStore != "" { persistence.go
63 isAnyCustom := c.DataStores[c.VisibilityStore].CustomDataStoreConfig != nil ||
64 c.DataStores[c.SecondaryVisibilityStore].CustomDataStoreConfig != nil
84 }
85
86 > for _, st := range stores { persistence.go
87 > ds, ok := c.DataStores[st]
88 > if !ok {
89 return fmt.Errorf("%w: missing config for datastore %q", ErrPersistenceConfig, st)
90 }
91 > if err := ds.Validate(); err != nil { persistence.go
92 return fmt.Errorf("%w: datastore %q: %s", ErrPersistenceConfig, st, err.Error())
93 }
94 }
95 > return nil persistence.go
96 }
97
102
103 // SecondaryVisibilityConfigExist returns whether user specified secondaryVisibilityStore in config
104 > func (c *Persistence) SecondaryVisibilityConfigExist() bool { persistence.go
105 > return c.SecondaryVisibilityStore != ""
106 > }
107
108 func (c *Persistence) IsSQLVisibilityStore() bool {
154
155 // Validate validates the data store config
156 > func (ds *DataStore) Validate() error { persistence.go
157 > storeConfigCount := 0
158 > if ds.SQL != nil {
159 > storeConfigCount++
160 > }
161 > if ds.Cassandra != nil {
162 storeConfigCount++
163 }
164 > if ds.CustomDataStoreConfig != nil { persistence.go
165 storeConfigCount++
166 }
167 > if ds.Elasticsearch != nil { persistence.go
168 storeConfigCount++
169 }
170 > if storeConfigCount != 1 { persistence.go
171 return errors.New(
172 "must provide config for one and only one datastore: " +
175 }
176
177 > if ds.SQL != nil { persistence.go
178 > if ds.SQL.TaskScanPartitions == 0 {
179 > ds.SQL.TaskScanPartitions = 1
180 > }
181 > if err := ds.SQL.validate(); err != nil {
182 return err
183 }
184 }
185 > if ds.Cassandra != nil { persistence.go
186 if err := ds.Cassandra.validate(); err != nil {
187 return err
188 }
189 }
190 > if ds.Elasticsearch != nil { persistence.go
191 if err := ds.Elasticsearch.Validate(); err != nil {
192 return err
193 }
194 }
195 > return nil persistence.go
196 }
197
go.temporal.io/server/temporal/server_options.go 21 introduced LOC · 9 ranges

Open complete file

66 )
67
68 > func newServerOptions(opts []ServerOption) *serverOptions { server_options.go
69 > so := &serverOptions{
70 > // Set defaults here.
71 > persistenceServiceResolver: resolver.NewNoopResolver(),
72 > }
73 > for _, opt := range opts {
74 > opt.apply(so)
75 > }
76
77 > return so server_options.go
78 }
79
80 > func (so *serverOptions) loadAndValidate() error { server_options.go
81 > for serviceName := range so.serviceNames {
82 > if !slices.Contains(Services, string(serviceName)) {
83 return fmt.Errorf("invalid service %q in service list %v", serviceName, so.serviceNames)
84 }
85 }
86
87 > if so.config == nil { server_options.go
88 err := so.loadConfig()
89 if err != nil {
92 }
93
94 > err := so.validateConfig() server_options.go
95 > if err != nil {
96 return fmt.Errorf("config validation error: %w", err)
97 }
98
99 > return nil server_options.go
100 }
101
126 }
127
128 > func (so *serverOptions) validateConfig() error { server_options.go
129 > if err := so.config.Validate(); err != nil {
130 return err
131 }
132
133 > for name := range so.serviceNames { server_options.go
134 > if _, ok := so.config.Services[string(name)]; !ok {
135 return fmt.Errorf("%q service is missing in config", name)
136 }
137 }
138 > return nil server_options.go
139 }
go.temporal.io/server/temporal/server_option.go 19 introduced LOC · 5 ranges

Open complete file

31 )
32
33 > func (f applyFunc) apply(s *serverOptions) { f(s) } server_option.go
34
35 // WithConfig sets a custom configuration
36 > func WithConfig(cfg *config.Config) ServerOption { server_option.go
37 > return applyFunc(func(s *serverOptions) {
38 > s.config = cfg
39 > })
40 }
41
55
56 // ForServices indicates which supplied services (e.g. frontend, history, matching, worker) within the server to start
57 > func ForServices(names []string) ServerOption { server_option.go
58 > return applyFunc(func(s *serverOptions) {
59 > s.serviceNames = make(map[primitives.ServiceName]struct{})
60 > for _, name := range names {
61 > s.serviceNames[primitives.ServiceName(name)] = struct{}{}
62 > }
63 })
64 }
82
83 // WithLogger sets a custom logger
84 > func WithLogger(logger log.Logger) ServerOption { server_option.go
85 > return applyFunc(func(s *serverOptions) {
86 > s.logger = logger
87 > })
88 }
89
138
139 // WithDynamicConfigClient sets custom client for reading dynamic configuration.
140 > func WithDynamicConfigClient(c dynamicconfig.Client) ServerOption { server_option.go
141 > return applyFunc(func(s *serverOptions) {
142 > s.dynamicConfigClient = c
143 > })
144 }
145
go.temporal.io/server/common/rpc/encryption/local_store_tls_provider.go 18 introduced LOC · 4 ranges

Open complete file

122 }
123
124 > func (s *localStoreTlsProvider) GetFrontendClientConfig() (*tls.Config, error) { local_store_tls_provider.go
125 >
126 > var client *config.ClientTLS
127 > var useTLS bool
128 > if isSystemWorker(s.settings) {
129 client = &s.settings.SystemWorker.Client
130 useTLS = true
132 > client = &s.settings.Frontend.Client
133 > useTLS = s.settings.Frontend.IsClientEnabled()
134 > }
135 > return s.getOrCreateConfig(
136 > &s.cachedFrontendClientConfig,
137 > func() (*tls.Config, error) {
138 return newClientTLSConfig(s.workerCertProvider, client.ServerName,
139 useTLS, true, !client.DisableHostVerification)
163 }
164
165 > func (s *localStoreTlsProvider) GetFrontendServerConfig() (*tls.Config, error) { local_store_tls_provider.go
166 > return s.getOrCreateConfig(
167 > &s.cachedFrontendServerConfig,
168 > func() (*tls.Config, error) {
169 return newServerTLSConfig(s.frontendCertProvider, s.frontendPerHostCertProviderMap, &s.settings.Frontend, s.logger)
170 },
218 ) (*tls.Config, error) {
219 if !isEnabled {
220 > return nil, nil local_store_tls_provider.go
221 > }
222
223 // Check if exists under a read lock first
go.temporal.io/server/common/config/config.go 16 introduced LOC · 6 ranges

Open complete file

693
694 // Validate validates this config
695 > func (c *Config) Validate() error { config.go
696 > if err := c.Persistence.Validate(); err != nil {
697 return err
698 }
699
700 > if err := c.Archival.Validate(&c.NamespaceDefaults.Archival); err != nil { config.go
701 return err
702 }
703
704 > _, hasIFE := c.Services[string(primitives.InternalFrontendService)] config.go
705 > if hasIFE && (c.PublicClient.HostPort != "" || c.PublicClient.ForceTLSConfig != "" || c.PublicClient.HTTPHostPort != "") {
706 return fmt.Errorf("when using internal-frontend, publicClient must be empty")
707 }
708
709 > switch c.PublicClient.ForceTLSConfig { config.go
710 > case ForceTLSConfigAuto, ForceTLSConfigInternode, ForceTLSConfigFrontend:
711 default:
712 return fmt.Errorf("invalid value for publicClient.forceTLSConfig: %q", c.PublicClient.ForceTLSConfig)
713 }
714
715 > return nil config.go
716 }
717
718 // String converts the config object into a string
719 > func (c *Config) String() string { config.go
720 > var buf bytes.Buffer
721 > encoder := yaml.NewEncoder(&buf)
722 > encoder.SetIndent(2)
723 > _ = encoder.Encode(c)
724 > maskedYaml, _ := masker.MaskYaml(buf.String(), masker.DefaultYAMLFieldNames)
725 > return maskedYaml
726 > }
727
728 func (r *GroupTLS) IsServerEnabled() bool {
go.temporal.io/server/common/config/fx.go 14 introduced LOC · 4 ranges

Open complete file

15 )
16
17 > func provideRPCConfig(cfg *Config, svcName primitives.ServiceName) *RPC { fx.go
18 > c := cfg.Services[string(svcName)].RPC
19 >
20 > return &c
21 > }
22
23 > func provideMembershipConfig(cfg *Config) *Membership { fx.go
24 > return &cfg.Global.Membership
25 > }
26
27 > func provideServicePortMap(cfg *Config) ServicePortMap { fx.go
28 > servicePortMap := make(ServicePortMap)
29 > for sn, sc := range cfg.Services {
30 > servicePortMap[primitives.ServiceName(sn)] = sc.RPC.GRPCPort
31 > }
32
33 > return servicePortMap fx.go
34 }
go.temporal.io/server/common/membership/ringpop/fx.go 13 introduced LOC · 4 ranges

Open complete file

13 )
14
15 > func provideFactory(lc fx.Lifecycle, params factoryParams) (*factory, error) { fx.go
16 > f, err := newFactory(params)
17 > if err != nil {
18 return nil, err
19 }
20 > lc.Append(fx.StopHook(f.closeTChannel)) fx.go
21 > return f, nil
22 }
23
24 > func provideMembership(lc fx.Lifecycle, f *factory) membership.Monitor { fx.go
25 > m := f.getMonitor()
26 > lc.Append(fx.StopHook(m.Stop))
27 > return m
28 > }
29
30 > func provideHostInfoProvider(lc fx.Lifecycle, f *factory) (membership.HostInfoProvider, error) { fx.go
31 > return f.getHostInfoProvider()
32 > }
go.temporal.io/server/common/config/archival.go 12 introduced LOC · 4 ranges

Open complete file

15
16 // Validate validates the archival config
17 > func (a *Archival) Validate(namespaceDefaults *ArchivalNamespaceDefaults) error { archival.go
18 > if !isArchivalConfigValid(a.History.State, a.History.EnableRead, namespaceDefaults.History.State, namespaceDefaults.History.URI, a.History.Provider != nil) {
19 return errors.New("invalid history archival config")
20 }
21
22 > if !isArchivalConfigValid(a.Visibility.State, a.Visibility.EnableRead, namespaceDefaults.Visibility.State, namespaceDefaults.Visibility.URI, a.Visibility.Provider != nil) { archival.go
23 return errors.New("invalid visibility archival config")
24 }
25
26 > return nil archival.go
27 }
28
33 domianDefaultURI string,
34 specifiedProvider bool,
35 > ) bool { archival.go
36 > archivalEnabled := clusterStatus == ArchivalEnabled
37 > URISet := len(domianDefaultURI) != 0
38 >
39 > validEnable := archivalEnabled && URISet && specifiedProvider
40 > validDisabled := !archivalEnabled && !enableRead && namespaceDefaultStatus != ArchivalEnabled && !URISet && !specifiedProvider
41 > return validEnable || validDisabled
42 > }
go.temporal.io/server/common/rpc/rpc.go 12 introduced LOC · 5 ranges

Open complete file

112
113 if d.tlsFactory != nil {
114 > serverConfig, err := d.tlsFactory.GetFrontendServerConfig() rpc.go
115 > if err != nil {
116 return nil, err
117 }
118 > if serverConfig == nil { rpc.go
119 > return opts, nil
120 > }
121 opts = append(opts, grpc.Creds(credentials.NewTLS(serverConfig)))
122 }
151 }
152 if d.tlsFactory != nil {
153 > serverConfig, err := d.tlsFactory.GetInternodeServerConfig() rpc.go
154 > if err != nil {
155 return nil, err
156 }
157 > if serverConfig == nil { rpc.go
158 > return opts, nil
159 > }
160 opts = append(opts, grpc.Creds(credentials.NewTLS(serverConfig)))
161 }
197
198 if cfg.BindOnLocalHost {
199 > return net.ParseIP(environment.GetLocalhostIP()) rpc.go
200 > }
201
202 if len(cfg.BindOnIP) > 0 {
go.temporal.io/server/common/persistence/persistence_rate_limited_clients.go 9 introduced LOC · 6 ranges

Open complete file

1020 ctx context.Context,
1021 request *GetClusterMembersRequest,
1022 > ) (*GetClusterMembersResponse, error) { persistence_rate_limited_clients.go
1023 > if err := allow(ctx, "GetClusterMembers", CallerSegmentMissing, c.systemRateLimiter, c.namespaceRateLimiter, c.shardRateLimiter); err != nil {
1024 return nil, err
1025 }
1026 > return c.persistence.GetClusterMembers(ctx, request) persistence_rate_limited_clients.go
1027 }
1028
1030 ctx context.Context,
1031 request *UpsertClusterMembershipRequest,
1033 > if err := allow(ctx, "UpsertClusterMembership", CallerSegmentMissing, c.systemRateLimiter, c.namespaceRateLimiter, c.shardRateLimiter); err != nil {
1034 return err
1035 }
1036 > return c.persistence.UpsertClusterMembership(ctx, request) persistence_rate_limited_clients.go
1037 }
1038
1040 ctx context.Context,
1041 request *PruneClusterMembershipRequest,
1043 > if err := allow(ctx, "PruneClusterMembership", CallerSegmentMissing, c.systemRateLimiter, c.namespaceRateLimiter, c.shardRateLimiter); err != nil {
1044 return err
1045 }
1046 > return c.persistence.PruneClusterMembership(ctx, request) persistence_rate_limited_clients.go
1047 }
1048
go.temporal.io/server/common/dynamicconfig/collection.go 8 introduced LOC · 2 ranges

Open complete file

131 c.cancelClientSubscription = notifyingClient.Subscribe(c.keysChanged)
132 } else {
133 > c.poller.Go(c.pollForChanges) collection.go
134 > }
135 }
136
159 }
160
161 > func (c *Collection) pollForChanges(ctx context.Context) error { collection.go
162 > interval := DynamicConfigSubscriptionPollInterval.Get(c)
163 > for ctx.Err() == nil {
164 > util.InterruptibleSleep(ctx, interval())
165 > c.pollOnce()
166 > }
167 return ctx.Err()
168 }
go.temporal.io/server/common/membership/ringpop/service_resolver.go 6 introduced LOC · 3 ranges

Open complete file

175 addrs := ring.LookupN(key, n)
176 if len(addrs) == 0 {
177 > r.RequestRefresh() service_resolver.go
178 > return nil
179 > }
180 return util.MapSlice(addrs, func(addr string) membership.HostInfo { return hosts[addr] })
181 }
go.temporal.io/server/common/persistence/sql/sqlplugin/sqlite/db.go 6 introduced LOC · 1 range

Open complete file

118
119 // VerifyVersion verify schema version is up to date
120 > func (mdb *db) VerifyVersion() error { db.go
121 > return nil
122 > // TODO(jlegrone): implement this
123 > // expectedVersion := mdb.ExpectedVersion()
124 > // return schema.VerifyCompatibleVersion(mdb, mdb.dbName, expectedVersion)
125 > }
go.temporal.io/server/common/pprof/fx.go 6 introduced LOC · 1 range

Open complete file

16 lc fx.Lifecycle,
17 pprof *PProfInitializerImpl,
18 > ) { fx.go
19 > lc.Append(
20 > fx.Hook{
21 > OnStart: func(context.Context) error {
22 > return pprof.Start()
23 > },
24 // todo: refactor pprof to gracefully shutdown http server
25 // OnStop: func(ctx context.Context) error {
go.temporal.io/server/common/pprof/pprof.go 6 introduced LOC · 1 range

Open complete file

31
32 // NewInitializer create a new instance of PProf Initializer
33 > func NewInitializer(cfg *config.PProf, logger log.Logger) *PProfInitializerImpl { pprof.go
34 > return &PProfInitializerImpl{
35 > PProf: cfg,
36 > Logger: logger,
37 > }
38 > }
39
40 // Start the pprof based on config
go.temporal.io/server/common/membership/ringpop/monitor.go 4 introduced LOC · 1 range

Open complete file

226 }
227
228 > func (rpo *monitor) WaitUntilInitialized(ctx context.Context) error { monitor.go
229 > _, err := rpo.initialized.Get(ctx)
230 > return err
231 > }
232
233 func (rpo *monitor) upsertMyMembership(