fx.go ×44

Frontier kind: Code frontier

unlabeled · c_e9d5304e7761

12 tests · 18736 LOC · 572 files · introduces 0 tests · 4693 LOC · 176 files

Introduces — evidence that enters the hierarchy at this concept

Code
705 ranges4693 lines · 176 files
Tests
0 tests

Contains — complete concept membership

All code (extent)
3483 ranges18736 lines · 572 files · Browse complete extent
All tests (intent)
12 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 176 files by introduced lines: 2410 of 4693 introduced LOC and 292 of 705 ranges. Expand a file to inspect source; the > gutter marks introduced lines.

go.temporal.io/server/service/frontend/fx.go 411 introduced LOC · 44 ranges

Open complete file

110 fx.Provide(AuthorizationInterceptorProvider),
111 fx.Provide(NamespaceCheckerProvider),
112 > fx.Provide(func(so GrpcServerOptions) *grpc.Server { return grpc.NewServer(so.Options...) }), fx.go
113 fx.Provide(callbackValidatorProvider),
114 fx.Provide(HandlerProvider),
154 metricsHandler metrics.Handler,
155 membershipMonitor membership.Monitor,
156 > ) *Service { fx.go
157 > return NewService(
158 > serviceConfig,
159 > server,
160 > healthServer,
161 > httpAPIServer,
162 > handler,
163 > adminHandler,
164 > operatorHandler,
165 > versionChecker,
166 > visibilityMgr,
167 > logger,
168 > grpcListener,
169 > metricsHandler,
170 > membershipMonitor,
171 > )
172 > }
173
174 // GrpcServerOptions are the options to build the frontend gRPC server along
189 audienceGetter authorization.JWTAudienceMapper,
190 dc *dynamicconfig.Collection,
191 > ) *authorization.Interceptor { fx.go
192 > return authorization.NewInterceptor(
193 > claimMapper,
194 > authorizer,
195 > metricsHandler,
196 > logger,
197 > namespaceChecker,
198 > audienceGetter,
199 > cfg.Global.Authorization.AuthHeaderName,
200 > cfg.Global.Authorization.AuthExtraHeaderName,
201 > serviceConfig.ExposeAuthorizerErrors,
202 > dynamicconfig.EnableCrossNamespaceCommands.Get(dc),
203 > dynamicconfig.EnablePrincipalPropagation.Get(dc),
204 > dynamicconfig.DisableStreamingAuthorizer.Get(dc),
205 > )
206 > }
207
208 > func NamespaceCheckerProvider(registry namespace.Registry) authorization.NamespaceChecker { fx.go
209 > return &namespaceChecker{r: registry}
210 > }
211
212 func (n *namespaceChecker) Exists(name namespace.Name) error {
249 customStreamInterceptors []grpc.StreamServerInterceptor,
250 metricsHandler metrics.Handler,
251 > ) GrpcServerOptions { fx.go
252 > kep := keepalive.EnforcementPolicy{
253 > MinTime: serviceConfig.KeepAliveMinTime(),
254 > PermitWithoutStream: serviceConfig.KeepAlivePermitWithoutStream(),
255 > }
256 > kp := keepalive.ServerParameters{
257 > MaxConnectionIdle: serviceConfig.KeepAliveMaxConnectionIdle(),
258 > MaxConnectionAge: serviceConfig.KeepAliveMaxConnectionAge(),
259 > MaxConnectionAgeGrace: serviceConfig.KeepAliveMaxConnectionAgeGrace(),
260 > Time: serviceConfig.KeepAliveTime(),
261 > Timeout: serviceConfig.KeepAliveTimeout(),
262 > }
263 > var grpcServerOptions []grpc.ServerOption
264 > var err error
265 > switch serviceName {
266 > case primitives.FrontendService:
267 > grpcServerOptions, err = rpcFactory.GetFrontendGRPCServerOptions()
268 case primitives.InternalFrontendService:
269 grpcServerOptions, err = rpcFactory.GetInternodeGRPCServerOptions()
271 err = fmt.Errorf("unexpected frontend service name %q", serviceName)
272 }
273 > if err != nil { fx.go
274 logger.Fatal("creating gRPC server options failed", tag.Error(err))
275 }
276 > unaryInterceptors := []grpc.UnaryServerInterceptor{ fx.go
277 > // Order of interceptors is important
278 > // Mask error interceptor should be the most outer interceptor since it handle the errors format
279 > // Service Error Interceptor should be the next most outer interceptor on error handling
280 > maskInternalErrorDetailsInterceptor.Intercept,
281 > serviceErrorInterceptor.Intercept,
282 > interceptor.NewFrontendServiceErrorInterceptor(logger),
283 > // BusinessID interceptor extracts business ID and adds it to context for use, must be before any interceptor that touches namespaces (namespaceValidator, handoverInterceptor)
284 > businessIDInterceptor.Intercept,
285 > namespaceValidatorInterceptor.NamespaceValidateIntercept,
286 > namespaceLogInterceptor.Intercept, // TODO: Deprecate this with a outer custom interceptor
287 > metrics.NewServerMetricsContextInjectorInterceptor(),
288 > authInterceptor.Intercept,
289 > // Handover interceptor has to above redirection because the request will route to the correct cluster after handover completed.
290 > // And retry cannot be performed before customInterceptors.
291 > namespaceHandoverInterceptor.Intercept,
292 > redirectionInterceptor.Intercept,
293 > // Telemetry interceptor must be after redirection to ensure metrics are recorded in the correct cluster
294 > telemetryInterceptor.UnaryIntercept,
295 > healthInterceptor.Intercept,
296 > namespaceValidatorInterceptor.StateValidationIntercept,
297 > namespaceCountLimiterInterceptor.Intercept,
298 > namespaceRateLimiterInterceptor.Intercept,
299 > rateLimitInterceptor.Intercept,
300 > sdkVersionInterceptor.Intercept,
301 > callerInfoInterceptor.Intercept,
302 > slowRequestLoggerInterceptor.Intercept,
303 > chasmRequestVisibilityInterceptor.Intercept,
304 > contextMetadataInterceptor.Intercept,
305 > }
306 > if len(customInterceptors) > 0 {
307 // TODO: Deprecate WithChainedFrontendGrpcInterceptors and provide a inner custom interceptor
308 unaryInterceptors = append(unaryInterceptors, customInterceptors...)
309 }
310 // retry interceptor should be the most inner interceptor
311 > unaryInterceptors = append(unaryInterceptors, retryableInterceptor.Intercept) fx.go
312 >
313 > streamInterceptor := []grpc.StreamServerInterceptor{
314 > authInterceptor.InterceptStream,
315 > telemetryInterceptor.StreamIntercept,
316 > }
317 > if len(customStreamInterceptors) > 0 {
318 streamInterceptor = append(streamInterceptor, customStreamInterceptors...)
319 }
320
321 > grpcServerOptions = append( fx.go
322 > grpcServerOptions,
323 > grpc.KeepaliveParams(kp),
324 > grpc.KeepaliveEnforcementPolicy(kep),
325 > grpc.ChainUnaryInterceptor(unaryInterceptors...),
326 > grpc.ChainStreamInterceptor(streamInterceptor...),
327 > )
328 >
329 > multiStats := rpc.MultiStatsHandler{}
330 > if traceStatsHandler != nil {
331 multiStats = append(multiStats, traceStatsHandler)
332 }
333 > if metricsStatsHandler != nil { fx.go
334 > multiStats = append(multiStats, metricsStatsHandler)
335 > }
336 > if len(multiStats) > 0 {
337 > grpcServerOptions = append(grpcServerOptions, grpc.StatsHandler(multiStats))
338 > }
339 > return GrpcServerOptions{Options: grpcServerOptions, UnaryInterceptors: unaryInterceptors}
340 }
341
343 dc *dynamicconfig.Collection,
344 persistenceConfig config.Persistence,
345 > ) *Config { fx.go
346 > return NewConfig(
347 > dc,
348 > persistenceConfig.NumHistoryShards,
349 > )
350 > }
351
352 func ServiceErrorInterceptorProvider(
353 dc *dynamicconfig.Collection,
354 > ) *interceptor.ServiceErrorInterceptor { fx.go
355 > return interceptor.NewServiceErrorInterceptor(
356 > dynamicconfig.MaxServiceErrorMessageLength.Get(dc),
357 > )
358 > }
359
360 func ThrottledLoggerRpsFnProvider(serviceConfig *Config) resource.ThrottledLoggerRpsFn {
365 namespaceLogger resource.NamespaceLogger,
366 namespaceRegistry namespace.Registry,
367 > ) *interceptor.NamespaceLogInterceptor { fx.go
368 > return interceptor.NewNamespaceLogInterceptor(
369 > namespaceRegistry,
370 > namespaceLogger)
371 > }
372
373 > func RetryableInterceptorProvider() *interceptor.RetryableInterceptor { fx.go
374 > return interceptor.NewRetryableInterceptor(
375 > common.CreateFrontendHandlerRetryPolicy(),
376 > common.IsServiceHandlerRetryableError,
377 > )
378 > }
379
380 func RedirectionInterceptorProvider(
387 timeSource clock.TimeSource,
388 clusterMetadata cluster.Metadata,
389 > ) *interceptor.Redirection { fx.go
390 > return interceptor.NewRedirection(
391 > configuration.EnableNamespaceNotActiveAutoForwarding,
392 > configuration.ForceNamespaceSelectedAPIAutoForwarding,
393 > namespaceCache,
394 > policy,
395 > logger,
396 > clientBean,
397 > metricsHandler,
398 > timeSource,
399 > clusterMetadata,
400 > )
401 > }
402
403 func BusinessIDInterceptorProvider(
404 extractor interceptor.RoutingKeyExtractor,
405 logger log.Logger,
406 > ) *interceptor.RoutingKeyInterceptor { fx.go
407 > return interceptor.NewRoutingKeyInterceptor(
408 > []interceptor.RoutingKeyExtractorFunc{
409 > interceptor.WorkflowServiceExtractor(extractor),
410 > },
411 > logger,
412 > )
413 > }
414
415 type NamespaceHandoverInterceptorParams struct {
426 func NamespaceHandoverInterceptorProvider(
427 params NamespaceHandoverInterceptorParams,
428 > ) *interceptor.NamespaceHandoverInterceptor { fx.go
429 > return interceptor.NewNamespaceHandoverInterceptor(
430 > params.DynamicConfig,
431 > params.NamespaceRegistry,
432 > params.MetricsHandler,
433 > params.Logger,
434 > params.TimeSource,
435 > params.RequestErrorHandler,
436 > params.AdditionalAllowedMethodsDuringHandover,
437 > )
438 > }
439
440 func ErrorHandlerProvider(
441 logger log.Logger,
442 serviceConfig *Config,
443 > ) *interceptor.RequestErrorHandler { fx.go
444 > return interceptor.NewRequestErrorHandler(
445 > logger,
446 > serviceConfig.LogAllReqErrors,
447 > )
448 > }
449
450 func TelemetryInterceptorProvider(
454 serviceConfig *Config,
455 requestErrorHandler *interceptor.RequestErrorHandler,
456 > ) *interceptor.TelemetryInterceptor { fx.go
457 > return interceptor.NewTelemetryInterceptor(
458 > namespaceRegistry,
459 > metricsHandler,
460 > logger,
461 > serviceConfig.LogAllReqErrors,
462 > requestErrorHandler,
463 > )
464 > }
465
466 func getRateFnWithMetrics(rateFn quotas.RateFn, handler metrics.Handler) quotas.RateFn {
509 logger log.Logger,
510 dc *dynamicconfig.Collection,
511 > ) *interceptor.ContextMetadataInterceptor { fx.go
512 > setTrailer := dynamicconfig.FrontendContextMetadataSetTrailer.Get(dc)()
513 > return interceptor.NewContextMetadataInterceptor(setTrailer, logger)
514 > }
515
516 func MaskInternalErrorDetailsInterceptorProvider(
518 serviceConfig *Config,
519 namespaceRegistry namespace.Registry,
520 > ) *interceptor.MaskInternalErrorDetailsInterceptor { fx.go
521 > return interceptor.NewMaskInternalErrorDetailsInterceptor(
522 > serviceConfig.MaskInternalErrorDetails, namespaceRegistry, logger,
523 > )
524 > }
525
526 func NamespaceRateLimitInterceptorProvider(
609 serviceResolver membership.ServiceResolver,
610 logger log.SnTaggedLogger,
611 > ) *interceptor.ConcurrentRequestLimitInterceptor { fx.go
612 > return interceptor.NewConcurrentRequestLimitInterceptor(
613 > namespaceRegistry,
614 > serviceResolver,
615 > logger,
616 > serviceConfig.MaxConcurrentLongRunningRequestsPerInstance,
617 > serviceConfig.MaxGlobalConcurrentLongRunningRequests,
618 > configs.ExecutionAPICountLimitOverride,
619 > )
620 > }
621
622 type NamespaceValidatorInterceptorParams struct {
629 func NamespaceValidatorInterceptorProvider(
630 params NamespaceValidatorInterceptorParams,
631 > ) *interceptor.NamespaceValidatorInterceptor { fx.go
632 > return interceptor.NewNamespaceValidatorInterceptor(
633 > params.NamespaceRegistry,
634 > params.ServiceConfig.EnableTokenNamespaceEnforcement,
635 > params.ServiceConfig.MaxIDLengthLimit,
636 > params.AdditionalAllowedMethodsDuringHandover,
637 > )
638 > }
639
640 > func SDKVersionInterceptorProvider() *interceptor.SDKVersionInterceptor { fx.go
641 > return interceptor.NewSDKVersionInterceptor()
642 > }
643
644 func CallerInfoInterceptorProvider(
645 namespaceRegistry namespace.Registry,
646 > ) *interceptor.CallerInfoInterceptor { fx.go
647 > return interceptor.NewCallerInfoInterceptor(namespaceRegistry)
648 > }
649
650 func SlowRequestLoggerInterceptorProvider(
651 logger log.Logger,
652 dc *dynamicconfig.Collection,
653 > ) *interceptor.SlowRequestLoggerInterceptor { fx.go
654 > return interceptor.NewSlowRequestLoggerInterceptor(
655 > logger,
656 > dynamicconfig.SlowRequestLoggingThreshold.Get(dc),
657 > )
658 > }
659
660 func PersistenceRateLimitingParamsProvider(
662 persistenceLazyLoadedServiceResolver service.PersistenceLazyLoadedServiceResolver,
663 logger log.SnTaggedLogger,
664 > ) service.PersistenceRateLimitingParams { fx.go
665 > return service.NewPersistenceRateLimitingParams(
666 > serviceConfig.PersistenceMaxQPS,
667 > serviceConfig.PersistenceGlobalMaxQPS,
668 > serviceConfig.PersistenceNamespaceMaxQPS,
669 > serviceConfig.PersistenceGlobalNamespaceMaxQPS,
670 > serviceConfig.PersistencePerShardNamespaceMaxQPS,
671 > serviceConfig.OperatorRPSRatio,
672 > serviceConfig.PersistenceQPSBurstRatio,
673 > serviceConfig.PersistenceDynamicRateLimitingParams,
674 > persistenceLazyLoadedServiceResolver,
675 > logger,
676 > )
677 > }
678
679 func VisibilityManagerProvider(
689 chasmRegistry *chasm.Registry,
690 serializer serialization.Serializer,
691 > ) (manager.VisibilityManager, error) { fx.go
692 > return visibility.NewManager(
693 > *persistenceConfig,
694 > persistenceServiceResolver,
695 > customVisibilityStoreFactory,
696 > nil, // frontend visibility never write
697 > saProvider,
698 > searchAttributesMapperProvider,
699 > namespaceRegistry,
700 > chasmRegistry,
701 > serviceConfig.VisibilityPersistenceMaxReadQPS,
702 > serviceConfig.VisibilityPersistenceMaxWriteQPS,
703 > serviceConfig.OperatorRPSRatio,
704 > serviceConfig.VisibilityPersistenceSlowQueryThreshold,
705 > serviceConfig.EnableReadFromSecondaryVisibility,
706 > serviceConfig.VisibilityEnableShadowReadMode,
707 > dynamicconfig.GetStringPropertyFn(visibility.SecondaryVisibilityWritingModeOff), // frontend visibility never write
708 > serviceConfig.VisibilityDisableOrderByClause,
709 > serviceConfig.VisibilityEnableManualPagination,
710 > serviceConfig.VisibilityEnableUnifiedQueryConverter,
711 > metricsHandler,
712 > logger,
713 > serializer,
714 > )
715 > }
716
717 func FEReplicatorNamespaceReplicationQueueProvider(
718 namespaceReplicationQueue persistence.NamespaceReplicationQueue,
719 clusterMetadata cluster.Metadata,
720 > ) FEReplicatorNamespaceReplicationQueue { fx.go
721 > var replicatorNamespaceReplicationQueue persistence.NamespaceReplicationQueue
722 > if clusterMetadata.IsGlobalNamespaceEnabled() {
723 replicatorNamespaceReplicationQueue = namespaceReplicationQueue
724 }
725 > return replicatorNamespaceReplicationQueue fx.go
726 }
727
729 membershipMonitor membership.Monitor,
730 serviceName primitives.ServiceName,
731 > ) (membership.ServiceResolver, error) { fx.go
732 > return membershipMonitor.GetResolver(serviceName)
733 > }
734
735 func AdminHandlerProvider(
766 schedulerClient schedulerpb.SchedulerServiceClient,
767 namespaceDLQHandler nsreplication.DLQMessageHandler,
768 > ) *AdminHandler { fx.go
769 > args := NewAdminHandlerArgs{
770 > persistenceConfig,
771 > configuration,
772 > namespaceReplicationQueue,
773 > replicatorNamespaceReplicationQueue,
774 > visibilityMgr,
775 > logger,
776 > taskManager,
777 > fairTaskManager,
778 > persistenceExecutionManager,
779 > clusterMetadataManager,
780 > persistenceMetadataManager,
781 > clientFactory,
782 > clientBean,
783 > historyClient,
784 > sdkClientFactory,
785 > membershipMonitor,
786 > hostInfoProvider,
787 > metricsHandler,
788 > namespaceRegistry,
789 > saProvider,
790 > saManager,
791 > saMapperProvider,
792 > clusterMetadata,
793 > healthServer,
794 > eventSerializer,
795 > timeSource,
796 > chasmRegistry,
797 > namespaceDataMerger,
798 > schedulerClient,
799 > taskCategoryRegistry,
800 > matchingClient,
801 > }
802 > return NewAdminHandler(args, namespaceDLQHandler)
803 > }
804
805 // NamespaceDLQHandlerProvider provides the default namespace DLQ message handler.
842 namespaceRegistry namespace.Registry,
843 nexusEndpointClient *NexusEndpointClient,
844 > ) *OperatorHandlerImpl { fx.go
845 > args := NewOperatorHandlerImplArgs{
846 > configuration,
847 > logger,
848 > sdkClientFactory,
849 > metricsHandler,
850 > visibilityMgr,
851 > saManager,
852 > healthServer,
853 > historyClient,
854 > clusterMetadataManager,
855 > clusterMetadata,
856 > clientFactory,
857 > namespaceRegistry,
858 > nexusEndpointClient,
859 > }
860 > return NewOperatorHandlerImpl(args)
861 > }
862
863 // callbackValidatorProvider creates a callback Validator using the production dynamic config keys
864 // so that existing operator configurations (callback.allowedAddresses) are honored.
865 > func callbackValidatorProvider(dc *dynamicconfig.Collection) callback.Validator { fx.go
866 > return callback.NewValidator(
867 > callback.MaxPerExecution.Get(dc),
868 > dynamicconfig.FrontendCallbackURLMaxLength.Get(dc),
869 > dynamicconfig.FrontendCallbackHeaderMaxSize.Get(dc),
870 > callback.AllowedAddresses.Get(dc),
871 > )
872 > }
873
874 func HandlerProvider(
911 registry *chasm.Registry,
912 frontendServiceResolver membership.ServiceResolver,
913 > ) Handler { fx.go
914 > workerDeploymentReadRateLimiter := configs.NewGlobalNamespaceRateLimiter(
915 > frontendServiceResolver,
916 > serviceConfig.GlobalWorkerDeploymentReadRPS,
917 > serviceConfig.GlobalWorkerDeploymentReadBurstRatio,
918 > log.With(logger, tag.ComponentRPCHandler, tag.ScopeNamespace),
919 > )
920 >
921 > wfHandler := NewWorkflowHandler(
922 > callbackValidator,
923 > serviceConfig,
924 > namespaceReplicationQueue,
925 > visibilityMgr,
926 > logger,
927 > throttledLogger,
928 > persistenceExecutionManager.GetName(),
929 > clusterMetadataManager,
930 > persistenceMetadataManager,
931 > historyClient,
932 > matchingClient,
933 > workerDeploymentStoreClient,
934 > schedulerClient,
935 > archiverProvider,
936 > payloadSerializer,
937 > namespaceRegistry,
938 > saMapperProvider,
939 > saProvider,
940 > saValidator,
941 > clusterMetadata,
942 > archivalMetadata,
943 > healthServer,
944 > timeSource,
945 > membershipMonitor,
946 > healthInterceptor,
947 > scheduleSpecBuilder,
948 > httpEnabled(cfg, serviceName),
949 > activityHandler,
950 > nexusOperationHandler,
951 > registry,
952 > workerDeploymentReadRateLimiter,
953 > chasmworkflow.NewValidator(
954 > chasmworkflow.NewConfig(dc),
955 > saMapperProvider,
956 > saValidator,
957 > ),
958 > )
959 > return wfHandler
960 > }
961
962 func RegisterNexusOperationHTTPHandler(
963 h *NexusOperationHTTPHandler,
964 router *mux.Router,
965 > ) { fx.go
966 > h.RegisterRoutes(router)
967 > }
968
969 func RegisterNexusCompletionHTTPHandler(
970 h *nexusCompletionHTTPHandler,
971 router *mux.Router,
972 > ) { fx.go
973 > h.RegisterRoutes(router)
974 > }
975
976 func RegisterOpenAPIHTTPHandler(
978 logger log.Logger,
979 router *mux.Router,
980 > ) *OpenAPIHTTPHandler { fx.go
981 > h := NewOpenAPIHTTPHandler(
982 > rateLimitInterceptor,
983 > logger,
984 > )
985 > h.RegisterRoutes(router)
986 > return h
987 > }
988
989 > func MuxRouterProvider() *mux.Router { fx.go
990 > // Instantiate a router to support additional route prefixes.
991 > return mux.NewRouter().UseEncodedPath()
992 > }
993
994 > func httpEnabled(cfg *config.Config, serviceName primitives.ServiceName) bool { fx.go
995 > // If the service is not the frontend service, HTTP API is disabled
996 > if serviceName != primitives.FrontendService && serviceName != primitives.InternalFrontendService {
997 return false
998 }
999 // If HTTP API port is 0, it is disabled
1000 > return cfg.Services[string(serviceName)].RPC.HTTPPort != 0 fx.go
1001 }
1002
1016 logger log.Logger,
1017 router *mux.Router,
1018 > ) (*HTTPAPIServer, error) { fx.go
1019 > if !httpEnabled(cfg, serviceName) {
1020 return nil, nil
1021 }
1042 nexusEndpointManager persistence.NexusEndpointManager,
1043 logger log.Logger,
1044 > ) *NexusEndpointClient { fx.go
1045 > clientConfig := newNexusEndpointClientConfig(dc)
1046 > return newNexusEndpointClient(
1047 > clientConfig,
1048 > namespaceRegistry,
1049 > matchingClient,
1050 > nexusEndpointManager,
1051 > logger,
1052 > )
1053 > }
1054
1055 > func ServiceLifetimeHooks(lc fx.Lifecycle, svc *Service) { fx.go
1056 > lc.Append(fx.StartStopHook(svc.Start, svc.Stop))
1057 > }
go.temporal.io/server/service/history/fx.go 193 introduced LOC · 32 ranges

Open complete file

108 chasmRegistry *chasm.Registry,
109 testHooks testhooks.TestHooks,
110 > ) { fx.go
111 > if hook, ok := testhooks.Get(
112 > testHooks,
113 > testhooks.HistoryChasmRuntimeProvider,
114 > testhooks.GlobalScope,
115 > ); ok {
116 hook(chasmEngine, chasmVisibilityManager, chasmRegistry)
117 }
129 )
130
131 > func ServerProvider(grpcServerOptions []grpc.ServerOption) *grpc.Server { fx.go
132 > return grpc.NewServer(grpcServerOptions...)
133 > }
134
135 > func HistoryServiceServerProvider(handler *Handler) historyservice.HistoryServiceServer { fx.go
136 > return handler
137 > }
138
139 func ServiceResolverProvider(
140 membershipMonitor membership.Monitor,
141 > ) (membership.ServiceResolver, error) { fx.go
142 > return membershipMonitor.GetResolver(primitives.HistoryService)
143 > }
144
145 func HandlerProvider(args NewHandlerArgs, lc fx.Lifecycle) (*Handler, error) {
202 lc.Append(fx.Hook{
203 OnStart: func(_ context.Context) error {
204 > h, err := buildNexusHandler(args.ChasmRegistry) fx.go
205 > if err != nil {
206 return err
207 }
208 > handler.nexusHandler = h fx.go
209 > return nil
210 },
211 })
214 }
215
216 > func buildNexusHandler(chasmRegistry *chasm.Registry) (nexus.Handler, error) { fx.go
217 > nexusServices := chasmRegistry.NexusServices()
218 > if len(nexusServices) == 0 {
219 return nil, nil
220 }
221 > serviceRegistry := nexus.NewServiceRegistry() fx.go
222 > for _, svc := range nexusServices {
223 > // No chance of collision here since the registry would have errored out earlier.
224 > serviceRegistry.MustRegister(svc)
225 > }
226
227 > return serviceRegistry.NewHandler() fx.go
228 }
229
230 func HistoryEngineFactoryProvider(
231 params HistoryEngineFactoryParams,
232 > ) shard.EngineFactory { fx.go
233 > return &historyEngineFactory{
234 > HistoryEngineFactoryParams: params,
235 > }
236 > }
237
238 func ConfigProvider(
239 dc *dynamicconfig.Collection,
240 persistenceConfig config.Persistence,
241 > ) *configs.Config { fx.go
242 > return configs.NewConfig(
243 > dc,
244 > persistenceConfig.NumHistoryShards,
245 > )
246 > }
247
248 func ServiceErrorInterceptorProvider(
249 dc *dynamicconfig.Collection,
250 > ) *interceptor.ServiceErrorInterceptor { fx.go
251 > return interceptor.NewServiceErrorInterceptor(
252 > dynamicconfig.MaxServiceErrorMessageLength.Get(dc),
253 > )
254 > }
255
256 func ThrottledLoggerRpsFnProvider(serviceConfig *configs.Config) resource.ThrottledLoggerRpsFn {
258 }
259
260 > func RetryableInterceptorProvider() *interceptor.RetryableInterceptor { fx.go
261 > return interceptor.NewRetryableInterceptor(
262 > common.CreateHistoryHandlerRetryPolicy(),
263 > api.IsRetryableError,
264 > )
265 > }
266
267 func ErrorHandlerProvider(
268 logger log.Logger,
269 serviceConfig *configs.Config,
270 > ) *interceptor.RequestErrorHandler { fx.go
271 > return interceptor.NewRequestErrorHandler(
272 > logger,
273 > serviceConfig.LogAllReqErrors,
274 > )
275 > }
276
277 func TelemetryInterceptorProvider(
281 serviceConfig *configs.Config,
282 requestErrorHandler *interceptor.RequestErrorHandler,
283 > ) *interceptor.TelemetryInterceptor { fx.go
284 > return interceptor.NewTelemetryInterceptor(
285 > namespaceRegistry,
286 > metricsHandler,
287 > logger,
288 > serviceConfig.LogAllReqErrors,
289 > requestErrorHandler,
290 > )
291 > }
292
293 func HealthSignalAggregatorProvider(
294 dynamicCollection *dynamicconfig.Collection,
295 logger log.ThrottledLogger,
296 > ) interceptor.HealthSignalAggregator { fx.go
297 > return interceptor.NewHealthSignalAggregator(
298 > logger,
299 > dynamicconfig.HistoryHealthSignalMetricsEnabled.Get(dynamicCollection),
300 > dynamicconfig.HistoryHealthSignalUsePercentiles.Get(dynamicCollection),
301 > dynamicconfig.PersistenceHealthSignalWindowSize.Get(dynamicCollection)(),
302 > dynamicconfig.PersistenceHealthSignalBufferSize.Get(dynamicCollection)(),
303 > dynamicconfig.HistoryHealthSignalLatencyWindowSize.Get(dynamicCollection)(),
304 > dynamicconfig.HistoryHealthSignalLatencyWindowCount.Get(dynamicCollection)(),
305 > )
306 > }
307
308 func HealthCheckInterceptorProvider(
309 healthSignalAggregator interceptor.HealthSignalAggregator,
310 > ) *interceptor.HealthCheckInterceptor { fx.go
311 > return interceptor.NewHealthCheckInterceptor(
312 > healthSignalAggregator,
313 > )
314 > }
315
316 > func ContextMetadataInterceptorProvider(logger log.Logger) *interceptor.ContextMetadataInterceptor { fx.go
317 > return interceptor.NewContextMetadataInterceptor(true, logger)
318 > }
319
320 func HistoryAdditionalInterceptorsProvider(
322 chasmRequestEngineInterceptor *chasm.ChasmEngineInterceptor,
323 chasmRequestVisibilityInterceptor *chasm.ChasmVisibilityInterceptor,
324 > ) []grpc.UnaryServerInterceptor { fx.go
325 > return []grpc.UnaryServerInterceptor{
326 > healthCheckInterceptor.UnaryIntercept,
327 > chasmRequestEngineInterceptor.Intercept,
328 > chasmRequestVisibilityInterceptor.Intercept,
329 > }
330 > }
331
332 func NamespaceRateLimitInterceptorProvider(
334 namespaceRegistry namespace.Registry,
335 metricsHandler metrics.Handler,
336 > ) interceptor.NamespaceRateLimitInterceptor { fx.go
337 >
338 > namespaceRateFn := func(namespaceName string) float64 {
339 if namespaceRPS := serviceConfig.NamespaceRPS(namespaceName); namespaceRPS > 0 {
340 return float64(namespaceRPS)
344 }
345
346 > return interceptor.NewNamespaceRateLimitInterceptor( fx.go
347 > namespaceRegistry,
348 > configs.NewNamespaceRateLimiter(
349 > namespaceRateFn,
350 > serviceConfig.OperatorRPSRatio,
351 > ),
352 > map[string]int{}, // no token overrides
353 > map[string]struct{}{}, // no long polls on history service
354 > dynamicconfig.GetBoolPropertyFnFilteredByNamespace(false), // no long poll methods
355 > metricsHandler,
356 > )
357 }
358
359 func RateLimitInterceptorProvider(
360 serviceConfig *configs.Config,
361 > ) *interceptor.RateLimitInterceptor { fx.go
362 > return interceptor.NewRateLimitInterceptor(
363 > configs.NewPriorityRateLimiter(func() float64 { return float64(serviceConfig.RPS()) }, serviceConfig.OperatorRPSRatio),
364 map[string]int{
365 healthpb.Health_Check_FullMethodName: 0, // exclude health check requests from rate limiting.
371 func ESProcessorConfigProvider(
372 serviceConfig *configs.Config,
373 > ) *elasticsearch.ProcessorConfig { fx.go
374 > return &elasticsearch.ProcessorConfig{
375 > IndexerConcurrency: serviceConfig.IndexerConcurrency,
376 > ESProcessorNumOfWorkers: serviceConfig.ESProcessorNumOfWorkers,
377 > ESProcessorBulkActions: serviceConfig.ESProcessorBulkActions,
378 > ESProcessorBulkSize: serviceConfig.ESProcessorBulkSize,
379 > ESProcessorFlushInterval: serviceConfig.ESProcessorFlushInterval,
380 > ESProcessorAckTimeout: serviceConfig.ESProcessorAckTimeout,
381 > }
382 > }
383
384 func PersistenceRateLimitingParamsProvider(
387 ownershipBasedQuotaScaler shard.LazyLoadedOwnershipBasedQuotaScaler,
388 logger log.SnTaggedLogger,
389 > ) service.PersistenceRateLimitingParams { fx.go
390 > hostCalculator := calculator.NewLoggedCalculator(
391 > shard.NewOwnershipAwareQuotaCalculator(
392 > ownershipBasedQuotaScaler,
393 > persistenceLazyLoadedServiceResolver,
394 > serviceConfig.PersistenceMaxQPS,
395 > serviceConfig.PersistenceGlobalMaxQPS,
396 > ),
397 > log.With(logger, tag.ComponentPersistence, tag.ScopeHost),
398 > )
399 > namespaceCalculator := calculator.NewLoggedNamespaceCalculator(
400 > shard.NewOwnershipAwareNamespaceQuotaCalculator(
401 > ownershipBasedQuotaScaler,
402 > persistenceLazyLoadedServiceResolver,
403 > serviceConfig.PersistenceNamespaceMaxQPS,
404 > serviceConfig.PersistenceGlobalNamespaceMaxQPS,
405 > ),
406 > log.With(logger, tag.ComponentPersistence, tag.ScopeNamespace),
407 > )
408 > return service.PersistenceRateLimitingParams{
409 > PersistenceMaxQps: func() int {
410 > return int(hostCalculator.GetQuota())
411 > },
412 PersistenceNamespaceMaxQps: func(namespace string) int {
413 return int(namespaceCalculator.GetQuota(namespace))
433 chasmRegistry *chasm.Registry,
434 serializer serialization.Serializer,
435 > ) (manager.VisibilityManager, error) { fx.go
436 > return visibility.NewManager(
437 > *persistenceConfig,
438 > persistenceServiceResolver,
439 > customVisibilityStoreFactory,
440 > esProcessorConfig,
441 > saProvider,
442 > searchAttributesMapperProvider,
443 > namespaceRegistry,
444 > chasmRegistry,
445 > serviceConfig.VisibilityPersistenceMaxReadQPS,
446 > serviceConfig.VisibilityPersistenceMaxWriteQPS,
447 > serviceConfig.OperatorRPSRatio,
448 > serviceConfig.VisibilityPersistenceSlowQueryThreshold,
449 > serviceConfig.EnableReadFromSecondaryVisibility,
450 > serviceConfig.VisibilityEnableShadowReadMode,
451 > serviceConfig.SecondaryVisibilityWritingMode,
452 > serviceConfig.VisibilityDisableOrderByClause,
453 > serviceConfig.VisibilityEnableManualPagination,
454 > serviceConfig.VisibilityEnableUnifiedQueryConverter,
455 > metricsHandler,
456 > logger,
457 > serializer,
458 > )
459 > }
460
461 func ChasmVisibilityManagerProvider(
475 metricsHandler metrics.Handler,
476 config *configs.Config,
477 > ) events.Notifier { fx.go
478 > return events.NewNotifier(
479 > timeSource,
480 > metricsHandler,
481 > config.GetShardID,
482 > )
483 > }
484
485 > func ServiceLifetimeHooks(lc fx.Lifecycle, svc *Service) { fx.go
486 > lc.Append(fx.StartStopHook(svc.Start, svc.Stop))
487 > }
488
489 func ReplicationProgressCacheProvider(
491 logger log.Logger,
492 handler metrics.Handler,
493 > ) replication.ProgressCache { fx.go
494 > return replication.NewProgressCache(serviceConfig, logger, handler)
495 > }
496
497 func VersionMembershipCacheProvider(
499 serviceConfig *configs.Config,
500 metricsHandler metrics.Handler,
501 > ) worker_versioning.VersionMembershipAndReactivationStatusCache { fx.go
502 > c := commoncache.New(serviceConfig.VersionMembershipCacheMaxSize(), &commoncache.Options{
503 > TTL: max(1*time.Second, serviceConfig.VersionMembershipCacheTTL()),
504 > })
505 > lc.Append(fx.Hook{
506 > OnStop: func(context.Context) error {
507 c.Stop()
508 return nil
509 },
510 })
511 > return worker_versioning.NewVersionMembershipAndReactivationStatusCache(c, metricsHandler) fx.go
512 }
513
516 serviceConfig *configs.Config,
517 metricsHandler metrics.Handler,
518 > ) worker_versioning.RoutingInfoCache { fx.go
519 > c := commoncache.New(serviceConfig.RoutingInfoCacheMaxSize(), &commoncache.Options{
520 > TTL: max(1*time.Second, serviceConfig.RoutingInfoCacheTTL()),
521 > })
522 > lc.Append(fx.Hook{
523 > OnStop: func(context.Context) error {
524 c.Stop()
525 return nil
526 },
527 })
528 > return worker_versioning.NewRoutingInfoCache(c, metricsHandler) fx.go
529 }
go.temporal.io/server/temporal/fx.go 163 introduced LOC · 33 ranges

Open complete file

632 metricsHandler metrics.Handler,
633 serializer serialization.Serializer,
634 > ) (*cluster.Config, config.Persistence, error) { fx.go
635 > ctx := context.TODO()
636 > logger = log.With(logger, tag.ComponentMetadataInitializer)
637 > metricsHandler = metricsHandler.WithTags(metrics.ServiceNameTag(primitives.ServerService))
638 > clusterName := persistenceClient.ClusterName(svc.ClusterMetadata.CurrentClusterName)
639 > dataStoreFactory := persistenceClient.DataStoreFactoryProvider(
640 > clusterName,
641 > persistenceServiceResolver,
642 > &svc.Persistence,
643 > customDataStoreFactory,
644 > logger,
645 > metricsHandler,
646 > telemetry.NoopTracerProvider,
647 > serializer,
648 > )
649 > factory := persistenceFactoryProvider(persistenceClient.NewFactoryParams{
650 > DataStoreFactory: dataStoreFactory,
651 > Cfg: &svc.Persistence,
652 > PersistenceMaxQPS: nil,
653 > PersistenceNamespaceMaxQPS: nil,
654 > ClusterName: persistenceClient.ClusterName(svc.ClusterMetadata.CurrentClusterName),
655 > MetricsHandler: metricsHandler,
656 > Logger: logger,
657 > Serializer: serializer,
658 > })
659 > defer factory.Close()
660 >
661 > clusterMetadataManager, err := factory.NewClusterMetadataManager()
662 > if err != nil {
663 return svc.ClusterMetadata, svc.Persistence, fmt.Errorf("error initializing cluster metadata manager: %w", err)
664 }
665 > defer clusterMetadataManager.Close() fx.go
666 >
667 > visCSAOverride := map[enumspb.IndexedValueType]int{}
668 > for tpName, value := range svc.Visibility.PersistenceCustomSearchAttributes {
669 saType, ok := enumspb.IndexedValueType_shorthandValue[tpName]
670 if !ok {
684 }
685
686 > visDataStores := []config.DataStore{ fx.go
687 > svc.Persistence.GetVisibilityStoreConfig(),
688 > svc.Persistence.GetSecondaryVisibilityStoreConfig(),
689 > }
690 > indexSearchAttributes := make(map[string]*persistencespb.IndexSearchAttributes)
691 > for _, ds := range visDataStores {
692 > indexSearchAttributes[ds.GetIndexName()] = sadefs.GetDBIndexSearchAttributes(visCSAOverride)
693 > }
694
695 > clusterMetadata := svc.ClusterMetadata fx.go
696 > if len(clusterMetadata.ClusterInformation) > 1 {
697 logger.Warn(
698 "All remote cluster settings under ClusterMetadata.ClusterInformation config will be ignored. "+
700 tag.Key("clusterInformation"))
701 }
702 > if _, ok := clusterMetadata.ClusterInformation[clusterMetadata.CurrentClusterName]; !ok { fx.go
703 logger.Error("Current cluster setting is missing under clusterMetadata.ClusterInformation",
704 tag.ClusterName(clusterMetadata.CurrentClusterName))
705 return svc.ClusterMetadata, svc.Persistence, missingCurrentClusterMetadataErr
706 }
707 > ctx = headers.SetCallerInfo(ctx, headers.SystemOperatorCallerInfo) fx.go
708 > resp, err := clusterMetadataManager.GetClusterMetadata(
709 > ctx,
710 > &persistence.GetClusterMetadataRequest{ClusterName: clusterMetadata.CurrentClusterName},
711 > )
712 > switch err.(type) {
713 case nil:
714 // Update current record
743 }
744
745 > clusterLoader := NewClusterMetadataLoader(clusterMetadataManager, logger) fx.go
746 > err = clusterLoader.LoadAndMergeWithStaticConfig(ctx, svc)
747 > if err != nil {
748 return svc.ClusterMetadata, svc.Persistence, fmt.Errorf("error while loading metadata from cluster: %w", err)
749 }
750 > return svc.ClusterMetadata, svc.Persistence, nil fx.go
751 }
752
956 // - []go.opentelemetry.io/otel/sdk/trace.SpanExporter
957 var TraceExportModule = fx.Options(
958 > fx.Provide(func(inputs SpanExporterInputs) ([]otelsdktrace.SpanExporter, error) { fx.go
959 > var tracingReady atomic.Bool
960 > otel.SetErrorHandler(otel.ErrorHandlerFunc(func(err error) {
961 if tracingReady.Load() { // ignore errors during startup
962 inputs.Logger.Warn("OTEL error", tag.Error(err), tag.ServiceErrorType(err))
965
966 // (1) Exporters from config.
967 > exportersByType := map[telemetry.SpanExporterType]otelsdktrace.SpanExporter{} fx.go
968 > if inputs.Config != nil {
969 > var err error
970 > exportersByType, err = inputs.Config.ExporterConfig.SpanExporters()
971 > if err != nil {
972 return nil, err
973 }
975
976 // (2) Exporters from env variables.
977 > exportersByTypeFromEnv, err := telemetry.SpanExportersFromEnv(os.LookupEnv) fx.go
978 > if err != nil {
979 return nil, err
980 }
981
982 // (3) Exporters from code (ie from testing).
983 > customExportersByType := inputs.Config.ExporterConfig.CustomExporters fx.go
984 >
985 > // Merge exporters.
986 > maps.Copy(exportersByType, exportersByTypeFromEnv) // env overrides config
987 > maps.Copy(exportersByType, customExportersByType) // custom overrides all
988 > exporters := expmaps.Values(exportersByType)
989 >
990 > // Configure exporters' lifecycle hooks.
991 > inputs.Lifecycyle.Append(fx.Hook{
992 > OnStart: func(ctx context.Context) error {
993 > err = startAll(exporters)(ctx)
994 > tracingReady.Store(true)
995 > return err
996 > },
997 OnStop: shutdownAll(exporters),
998 })
999 > return exporters, nil fx.go
1000 }),
1001 )
1020 fx.Provide(
1021 fx.Annotate(
1022 > func(exps []otelsdktrace.SpanExporter, opts []otelsdktrace.BatchSpanProcessorOption) []otelsdktrace.SpanProcessor { fx.go
1023 > sps := make([]otelsdktrace.SpanProcessor, 0, len(exps))
1024 > for _, exp := range exps {
1025 sps = append(sps, otelsdktrace.NewBatchSpanProcessor(exp, opts...))
1026 }
1027 > return sps fx.go
1028 },
1029 fx.ParamTags(`optional:"true"`, ``),
1032 fx.Provide(
1033 fx.Annotate(
1034 > func(rsn primitives.ServiceName, rsi resource.InstanceID) (*otelresource.Resource, error) { fx.go
1035 > attrs := []attribute.KeyValue{
1036 > semconv.ServiceNameKey.String(telemetry.ResourceServiceName(rsn, os.LookupEnv)),
1037 > semconv.ServiceVersionKey.String(headers.ServerVersion),
1038 > }
1039 > if rsi != "" {
1040 attrs = append(attrs, semconv.ServiceInstanceIDKey.String(string(rsi)))
1041 }
1042
1043 > return otelresource.New(context.Background(), fx.go
1044 > otelresource.WithProcess(),
1045 > otelresource.WithOS(),
1046 > otelresource.WithHost(),
1047 > otelresource.WithContainer(),
1048 > otelresource.WithAttributes(attrs...),
1049 > )
1050 },
1051 fx.ParamTags(``, `optional:"true"`),
1052 ),
1053 ),
1054 > fx.Provide(func(lc fx.Lifecycle, r *otelresource.Resource, sps []otelsdktrace.SpanProcessor) trace.TracerProvider { fx.go
1055 > if len(sps) == 0 {
1056 return telemetry.NoopTracerProvider
1057 }
1082 }),
1083 // Haven't had use for baggage propagation yet
1084 > fx.Provide(func() propagation.TextMapPropagator { return propagation.TraceContext{} }), fx.go
1085 fx.Provide(telemetry.NewServerStatsHandler),
1086 fx.Provide(telemetry.NewClientStatsHandler),
1088 )
1089
1090 > func startAll(exporters []otelsdktrace.SpanExporter) func(ctx context.Context) error { fx.go
1091 > type starter interface{ Start(context.Context) error }
1092 > return func(ctx context.Context) error {
1093 > for _, e := range exporters {
1094 if starter, ok := e.(starter); ok {
1095 err := starter.Start(ctx)
1099 }
1100 }
1101 > return nil fx.go
1102 }
1103 }
1104
1105 > func shutdownAll(exporters []otelsdktrace.SpanExporter) func(ctx context.Context) error { fx.go
1106 > return func(ctx context.Context) error {
1107 shutdownCtx, cancel := context.WithTimeout(context.Background(), 1*time.Second)
1108 defer cancel()
1120 }
1121
1122 > var FxLogAdapter = fx.WithLogger(func(logger log.Logger) fxevent.Logger { fx.go
1123 > return &fxLogAdapter{logger: logger}
1124 > })
1125
1126 type fxLogAdapter struct {
1128 }
1129
1130 > func (l *fxLogAdapter) LogEvent(e fxevent.Event) { fx.go
1131 > switch e := e.(type) {
1132 > case *fxevent.OnStartExecuting:
1133 > l.logger.Debug("OnStart hook executing",
1134 > tag.ComponentFX,
1135 > tag.String("callee", e.FunctionName),
1136 > tag.String("caller", e.CallerName),
1137 > )
1138 > case *fxevent.OnStartExecuted:
1139 > if e.Err != nil {
1140 l.logger.Error("OnStart hook failed",
1141 tag.ComponentFX,
1144 tag.Error(e.Err),
1145 )
1146 > } else { fx.go
1147 > l.logger.Debug("OnStart hook executed",
1148 > tag.ComponentFX,
1149 > tag.String("callee", e.FunctionName),
1150 > tag.String("caller", e.CallerName),
1151 > tag.Stringer("runtime", e.Runtime),
1152 > )
1153 > }
1154 case *fxevent.OnStopExecuting:
1155 l.logger.Debug("OnStop hook executing",
1174 )
1175 }
1176 > case *fxevent.Supplied: fx.go
1177 > if e.Err != nil {
1178 l.logger.Error("supplied",
1179 tag.ComponentFX,
1182 tag.Error(e.Err))
1183 }
1184 > case *fxevent.Provided: fx.go
1185 > if e.Err != nil {
1186 l.logger.Error("error encountered while applying options",
1187 tag.ComponentFX,
1196 tag.Error(e.Err))
1197 }
1198 > case *fxevent.Decorated: fx.go
1199 > if e.Err != nil {
1200 l.logger.Error("error encountered while applying options",
1201 tag.ComponentFX,
1203 tag.Error(e.Err))
1204 }
1205 > case *fxevent.Run: fx.go
1206 > if e.Err != nil {
1207 l.logger.Error("error returned",
1208 tag.ComponentFX,
1213 )
1214 }
1215 > case *fxevent.Invoking: fx.go
1216 > // Do not log stack as it will make logs hard to read.
1217 > l.logger.Debug("invoking",
1218 > tag.ComponentFX,
1219 > tag.String("function", e.FunctionName),
1220 > tag.String("module", e.ModuleName),
1221 > )
1222 > case *fxevent.Invoked:
1223 > if e.Err != nil {
1224 l.logger.Error("invoke failed",
1225 tag.ComponentFX,
1244 l.logger.Error("rollback failed", tag.ComponentFX, tag.Error(e.Err))
1245 }
1246 > case *fxevent.Started: fx.go
1247 > if e.Err != nil {
1248 l.logger.Error("start failed", tag.ComponentFX, tag.Error(e.Err))
1249 > } else { fx.go
1250 > l.logger.Debug("started", tag.ComponentFX)
1251 > }
1252 > case *fxevent.LoggerInitialized:
1253 > if e.Err != nil {
1254 l.logger.Error("custom logger initialization failed", tag.ComponentFX, tag.Error(e.Err))
1255 > } else { fx.go
1256 > l.logger.Debug("initialized custom fxevent.Logger",
1257 > tag.ComponentFX,
1258 > tag.String("function", e.ConstructorName))
1259 > }
1260 > case *fxevent.BeforeRun:
1261 > l.logger.Debug("before run",
1262 > tag.ComponentFX,
1263 > tag.String("name", e.Name),
1264 > tag.String("kind", e.Kind),
1265 > tag.String("module", e.ModuleName),
1266 > )
1267 default:
1268 l.logger.Warn("unknown fx log type, update fxLogAdapter",
go.temporal.io/server/common/resource/fx.go 148 introduced LOC · 25 ranges

Open complete file

89 fx.Provide(SearchAttributeManagerProvider),
90 fx.Provide(NamespaceRegistryProvider),
91 > fx.Provide(func() namespace.NamespaceStateChangedFn { return nsregistry.DefaultNamespaceStateChanged }), fx.go
92 nsregistry.RegistryLifetimeHooksModule,
93 fx.Provide(fx.Annotate(
94 > func(p namespace.Registry) pingable.Pingable { return p }, fx.go
95 fx.ResultTags(`group:"deadlockDetectorRoots"`),
96 )),
127 )
128
129 > func DefaultSnTaggedLoggerProvider(logger log.Logger, sn primitives.ServiceName) log.SnTaggedLogger { fx.go
130 > return log.With(logger, tag.Service(sn))
131 > }
132
133 func ThrottledLoggerProvider(
141 }
142
143 > func GrpcListenerProvider(factory common.RPCFactory) net.Listener { fx.go
144 > return factory.GetGRPCListener()
145 > }
146
147 > func HostNameProvider() (HostName, error) { fx.go
148 > hn, err := os.Hostname()
149 > return HostName(hn), err
150 > }
151
152 > func TimeSourceProvider() clock.TimeSource { fx.go
153 > return clock.NewRealTimeSource()
154 > }
155
156 func SearchAttributeMapperProviderProvider(
159 searchAttributeProvider searchattribute.Provider,
160 persistenceConfig *config.Persistence,
161 > ) searchattribute.MapperProvider { fx.go
162 > primaryVisibilityStoreConfig := persistenceConfig.GetVisibilityStoreConfig()
163 > return searchattribute.NewMapperProvider(
164 > saMapper,
165 > namespaceRegistry,
166 > searchAttributeProvider,
167 > primaryVisibilityStoreConfig.GetIndexName(),
168 > )
169 > }
170
171 func SearchAttributeProviderProvider(
174 cmMgr persistence.ClusterMetadataManager,
175 dynamicCollection *dynamicconfig.Collection,
176 > ) searchattribute.Provider { fx.go
177 > return searchattribute.NewManager(
178 > timeSource,
179 > cmMgr,
180 > logger,
181 > dynamicconfig.ForceSearchAttributesCacheRefreshOnRead.Get(dynamicCollection))
182 > }
183
184 func SearchAttributeManagerProvider(
187 cmMgr persistence.ClusterMetadataManager,
188 dynamicCollection *dynamicconfig.Collection,
189 > ) searchattribute.Manager { fx.go
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
205 metricsHandler metrics.Handler,
206 logger log.Logger,
207 > ) *searchattribute.Validator { fx.go
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 {
235 }
236
237 > func NamespaceRegistryProvider(params NamespaceRegistryParams) namespace.Registry { fx.go
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(
259 logger log.SnTaggedLogger,
260 throttledLogger log.ThrottledLogger,
261 > ) client.Factory { fx.go
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(
276 clientFactory client.Factory,
277 clusterMetadata cluster.Metadata,
278 > ) (client.Bean, error) { fx.go
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
289 > return bean, nil
290 }
291
292 > func FrontendClientProvider(clientBean client.Bean) workflowservice.WorkflowServiceClient { fx.go
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
302 > adminRawClient, err := clientBean.GetRemoteAdminClient(clusterMetadata.GetCurrentClusterName())
303 > if err != nil {
304 return nil, err
305 }
306 > return admin.NewRetryableClient( fx.go
307 > adminRawClient,
308 > common.CreateFrontendClientRetryPolicy(),
309 > common.IsServiceClientTransientError,
310 > ), nil
311 }
312
313 func RuntimeMetricsReporterProvider(
314 params RuntimeMetricsReporterParams,
315 > ) *metrics.RuntimeMetricsReporter { fx.go
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
325 > return clientBean.GetHistoryClient()
326 > }
327
328 > func HistoryClientProvider(historyRawClient HistoryRawClient, dc *dynamicconfig.Collection) HistoryClient { fx.go
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
340 > return clientBean.GetMatchingClient(namespaceRegistry.GetNamespaceName)
341 > }
342
343 > func MatchingClientProvider(matchingRawClient MatchingRawClient, dc *dynamicconfig.Collection) MatchingClient { fx.go
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
353 > persistenceConfig.TransactionSizeLimit = dynamicconfig.TransactionSizeLimit.Get(dc)
354 > return &persistenceConfig
355 > }
356
357 func ArchivalMetadataProvider(dc *dynamicconfig.Collection, cfg *config.Config) archiver.ArchivalMetadata {
412 func PerServiceDialOptionsProvider(
413 logger log.SnTaggedLogger,
414 > ) map[primitives.ServiceName][]grpc.DialOption { fx.go
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(
469 metadata cluster.Metadata,
470 tlsConfigProvider encryption.TLSConfigProvider,
471 > ) *cluster.FrontendHTTPClientCache { fx.go
472 > return cluster.NewFrontendHTTPClientCache(metadata, tlsConfigProvider)
473 > }
474
475 func getFrontendConnectionDetails(
go.temporal.io/server/service/history/replication/fx.go 141 introduced LOC · 22 ranges

Open complete file

42 var Module = fx.Provide(
43 NewTaskFetcherFactory,
44 > func(m persistence.ExecutionManager) ExecutionManager { fx.go
45 > return m
46 > },
47 nsreplication.NewNoopDataMerger,
48 nsreplication.NewDefaultAdmitter,
51 ServerSchedulerRateLimiterProvider,
52 PersistenceRateLimiterProvider,
53 > func(serializer serialization.Serializer) TaskSerializer { fx.go
54 > return serializer
55 > },
56 replicationTaskConverterFactoryProvider,
57 replicationTaskExecutorProvider,
86 metricsHandler metrics.Handler,
87 testHooks testhooks.TestHooks,
88 > ) EagerNamespaceRefresher { fx.go
89 > return NewEagerNamespaceRefresher(
90 > metadataManager,
91 > namespaceRegistry,
92 > logger,
93 > clientBean,
94 > nsreplication.NewTaskExecutor(
95 > clusterMetadata.GetCurrentClusterName(),
96 > metadataManager,
97 > dataMerger,
98 > admitter,
99 > logger,
100 > testHooks,
101 > ),
102 > clusterMetadata.GetCurrentClusterName(),
103 > metricsHandler,
104 > )
105 > }
106
107 func replicationTaskConverterFactoryProvider(
108 config *configs.Config,
109 replicationTaskSerializer TaskSerializer,
110 > ) SourceTaskConverterProvider { fx.go
111 > return func(
112 > historyEngine historyi.Engine,
113 > shardContext historyi.ShardContext,
114 > clientClusterName string,
115 > serializer serialization.Serializer,
116 > ) SourceTaskConverter {
117 return NewSourceTaskConverter(
118 historyEngine,
124 }
125
126 > func replicationTaskExecutorProvider() TaskExecutorProvider { fx.go
127 > return func(params TaskExecutorParams) TaskExecutor {
128 return NewTaskExecutor(
129 params.RemoteCluster,
141 queueFactory ctasks.SequentialTaskQueueFactory[TrackableExecutableTask],
142 lc fx.Lifecycle,
143 > ) ctasks.Scheduler[TrackableExecutableTask] { fx.go
144 > // SequentialScheduler has panic wrapper when executing task,
145 > // if changing the executor, please make sure other executor has panic wrapper
146 > scheduler := ctasks.NewSequentialScheduler[TrackableExecutableTask](
147 > &ctasks.SequentialSchedulerOptions{
148 > QueueSize: config.ReplicationProcessorSchedulerQueueSize(),
149 > WorkerCount: config.ReplicationProcessorSchedulerWorkerCount,
150 > },
151 > WorkflowKeyHashFn,
152 > queueFactory,
153 > logger,
154 > )
155 > taskChannelKeyFn := func(e TrackableExecutableTask) ClusterChannelKey {
156 return ClusterChannelKey{
157 ClusterName: e.SourceClusterName(),
158 }
159 }
160 > channelWeightFn := func(key ClusterChannelKey) int { fx.go
161 return 1
162 }
163 // This creates a per cluster channel.
164 // They share the same weight so it just does a round-robin on all clusters' tasks.
165 > rrScheduler := ctasks.NewInterleavedWeightedRoundRobinScheduler( fx.go
166 > ctasks.InterleavedWeightedRoundRobinSchedulerOptions[TrackableExecutableTask, ClusterChannelKey]{
167 > TaskChannelKeyFn: taskChannelKeyFn,
168 > ChannelWeightFn: channelWeightFn,
169 > },
170 > scheduler,
171 > logger,
172 > )
173 > lc.Append(fx.StartStopHook(rrScheduler.Start, rrScheduler.Stop))
174 > return rrScheduler
175 }
176
183 metricsHandler metrics.Handler,
184 lc fx.Lifecycle,
185 > ) ctasks.Scheduler[TrackableExecutableTask] { fx.go
186 > // P-way parallelism for executions of the same workflow (per ReplicationLowPriorityTaskParallelism)
187 > // is modeled as P distinct per-namespace-workflow queue IDs. We bucket by execution (RunID) so all
188 > // low-priority tasks for one execution share a queue; the third field stores the slot index, not
189 > // the run UUID.
190 > queueFactory := func(task TrackableExecutableTask) ctasks.SequentialTaskQueue[TrackableExecutableTask] {
191 item := task.QueueID()
192 workflowKey, ok := item.(definition.WorkflowKey)
204 // SequentialScheduler has panic wrapper when executing task,
205 // if changing the executor, please make sure other executor has panic wrapper
206 > scheduler := ctasks.NewSequentialScheduler[TrackableExecutableTask]( fx.go
207 > &ctasks.SequentialSchedulerOptions{
208 > QueueSize: config.ReplicationProcessorSchedulerQueueSize(),
209 > WorkerCount: config.ReplicationLowPriorityProcessorSchedulerWorkerCount,
210 > },
211 > WorkflowKeyHashFn,
212 > queueFactory,
213 > logger,
214 > )
215 > taskChannelKeyFn := func(e TrackableExecutableTask) ClusterChannelKey {
216 return ClusterChannelKey{
217 ClusterName: e.SourceClusterName(),
218 }
219 }
220 > channelWeightFn := func(key ClusterChannelKey) int { fx.go
221 return 1
222 }
223 > taskQuotaRequestFn := func(t TrackableExecutableTask) quotas.Request { fx.go
224 var taskType string
225 var nsName namespace.Name
245 "")
246 }
247 > taskMetricsTagsFn := func(t TrackableExecutableTask) []metrics.Tag { fx.go
248 replicationTask := t.ReplicationTask()
249 var taskType string
268 // This creates a per cluster channel.
269 // They share the same weight so it just does a round-robin on all clusters' tasks.
270 > rrScheduler := ctasks.NewInterleavedWeightedRoundRobinScheduler( fx.go
271 > ctasks.InterleavedWeightedRoundRobinSchedulerOptions[TrackableExecutableTask, ClusterChannelKey]{
272 > TaskChannelKeyFn: taskChannelKeyFn,
273 > ChannelWeightFn: channelWeightFn,
274 > },
275 > scheduler,
276 > logger,
277 > )
278 > ts := ctasks.NewRateLimitedScheduler[TrackableExecutableTask](
279 > rrScheduler,
280 > rateLimiter,
281 > timeSource,
282 > taskQuotaRequestFn,
283 > taskMetricsTagsFn,
284 > ctasks.RateLimitedSchedulerOptions{
285 > Enabled: config.ReplicationEnableRateLimit,
286 > EnableShadowMode: config.ReplicationEnableRateLimitShadowMode,
287 > },
288 > logger,
289 > metricsHandler,
290 > )
291 > lc.Append(fx.StartStopHook(ts.Start, ts.Stop))
292 > return ts
293 }
294
297 metricsHandler metrics.Handler,
298 config *configs.Config,
299 > ) ctasks.SequentialTaskQueueFactory[TrackableExecutableTask] { fx.go
300 > return func(task TrackableExecutableTask) ctasks.SequentialTaskQueue[TrackableExecutableTask] {
301 if config.EnableReplicationTaskBatching() {
302 return NewSequentialBatchableTaskQueue(task, nil, logger, metricsHandler)
315 func executableTaskConverterProvider(
316 processToolBox ProcessToolBox,
317 > ) ExecutableTaskConverter { fx.go
318 > return NewExecutableTaskConverter(processToolBox)
319 > }
320
321 func streamReceiverMonitorProvider(
322 processToolBox ProcessToolBox,
323 taskConverter ExecutableTaskConverter,
324 > ) StreamReceiverMonitor { fx.go
325 > return NewStreamReceiverMonitor(
326 > processToolBox,
327 > taskConverter,
328 > processToolBox.Config.EnableReplicationStream(),
329 > )
330 > }
331
332 func resendHandlerProvider(
340 logger log.Logger,
341 importer eventhandler.EventImporter,
342 > ) eventhandler.ResendHandler { fx.go
343 > return eventhandler.NewResendHandler(
344 > namespaceRegistry,
345 > clientBean,
346 > serializer,
347 > clusterMetadata,
348 > func(ctx context.Context, namespaceId namespace.ID, workflowId string) (historyi.Engine, error) {
349 shardContext, err := shardController.GetShardByNamespaceWorkflow(
350 namespaceId,
368 serializer serialization.Serializer,
369 logger log.Logger,
370 > ) eventhandler.EventImporter { fx.go
371 > return eventhandler.NewEventImporter(
372 > historyFetcher,
373 > func(ctx context.Context, namespaceId namespace.ID, workflowId string) (historyi.Engine, error) {
374 shardContext, err := shardController.GetShardByNamespaceWorkflow(
375 namespaceId,
390 replicationTaskSerializer TaskSerializer,
391 clusterMetadata cluster.Metadata,
392 > ) *DLQWriterAdapter { fx.go
393 > return NewDLQWriterAdapter(dlqWriter, replicationTaskSerializer, clusterMetadata.GetCurrentClusterName())
394 > }
395
396 func historyEventsHandlerProvider(
399 shardController shard.Controller,
400 logger log.Logger,
401 > ) eventhandler.HistoryEventsHandler { fx.go
402 > return eventhandler.NewHistoryEventsHandler(
403 > clusterMetadata,
404 > importer,
405 > shardController,
406 > logger,
407 > )
408 > }
409
410 func historyPaginatedFetcherProvider(
413 serializer serialization.Serializer,
414 logger log.Logger,
415 > ) eventhandler.HistoryPaginatedFetcher { fx.go
416 > return eventhandler.NewHistoryPaginatedFetcher(
417 > namespaceRegistry,
418 > clientBean,
419 > serializer,
420 > logger,
421 > )
422 > }
go.temporal.io/server/service/history/history_engine.go 138 introduced LOC · 8 ranges

Open complete file

182 testHooks testhooks.TestHooks,
183 chasmEngine chasm.Engine,
184 > ) historyi.Engine { history_engine.go
185 > currentClusterName := shard.GetClusterMetadata().GetCurrentClusterName()
186 >
187 > logger := shard.GetLogger()
188 > executionManager := shard.GetExecutionManager()
189 >
190 > workflowDeleteManager := deletemanager.NewDeleteManager(
191 > shard,
192 > workflowCache,
193 > config,
194 > shard.GetTimeSource(),
195 > persistenceVisibilityMgr,
196 > )
197 > syncStateRetriever := replication.NewSyncStateRetriever(
198 > shard,
199 > workflowCache,
200 > workflowConsistencyChecker,
201 > eventBlobCache,
202 > shard.GetLogger(),
203 > )
204 >
205 > historyEngImpl := &historyEngineImpl{
206 > status: common.DaemonStatusInitialized,
207 > currentClusterName: currentClusterName,
208 > shardContext: shard,
209 > clusterMetadata: shard.GetClusterMetadata(),
210 > timeSource: shard.GetTimeSource(),
211 > executionManager: executionManager,
212 > tokenSerializer: tasktoken.NewSerializer(),
213 > logger: log.With(logger, tag.ComponentHistoryEngine),
214 > throttledLogger: log.With(shard.GetThrottledLogger(), tag.ComponentHistoryEngine),
215 > metricsHandler: shard.GetMetricsHandler(),
216 > eventNotifier: eventNotifier,
217 > config: config,
218 > sdkClientFactory: sdkClientFactory,
219 > matchingClient: matchingClient,
220 > rawMatchingClient: rawMatchingClient,
221 > persistenceVisibilityMgr: persistenceVisibilityMgr,
222 > workflowDeleteManager: workflowDeleteManager,
223 > serializer: serializer,
224 > workflowConsistencyChecker: workflowConsistencyChecker,
225 > versionChecker: headers.NewDefaultVersionChecker(),
226 > tracer: tracerProvider.Tracer(consts.LibraryName),
227 > taskCategoryRegistry: taskCategoryRegistry,
228 > commandHandlerRegistry: commandHandlerRegistry,
229 > chasmWorkflowRegistry: chasmWorkflowRegistry,
230 > workflowCache: workflowCache,
231 > replicationProgressCache: replicationProgressCache,
232 > syncStateRetriever: syncStateRetriever,
233 > outboundQueueCBPool: outboundQueueCBPool,
234 > testHooks: testHooks,
235 > chasmEngine: chasmEngine,
236 > versionCache: versionCache,
237 > workerDeploymentClient: workerDeploymentClient,
238 > routingInfoCache: routingInfoCache,
239 > }
240 >
241 > historyEngImpl.queueProcessors = make(map[tasks.Category]queues.Queue)
242 > for _, factory := range queueProcessorFactories {
243 > processor := factory.CreateQueue(shard)
244 > historyEngImpl.queueProcessors[processor.Category()] = processor
245 > }
246
247 > historyEngImpl.eventsReapplier = ndc.NewEventsReapplier(shard.StateMachineRegistry(), shard.ChasmWorkflowRegistry(), shard.GetMetricsHandler(), logger) history_engine.go
248 >
249 > if shard.GetClusterMetadata().IsGlobalNamespaceEnabled() {
250 historyEngImpl.replicationAckMgr = replication.NewAckManager(
251 shard,
289 )
290 }
291 > historyEngImpl.workflowRebuilder = NewWorkflowRebuilder( history_engine.go
292 > shard,
293 > workflowCache,
294 > logger,
295 > )
296 > historyEngImpl.workflowResetter = ndc.NewWorkflowResetter(
297 > shard,
298 > workflowCache,
299 > logger,
300 > )
301 >
302 > historyEngImpl.searchAttributesValidator = searchattribute.NewValidator(
303 > shard.GetSearchAttributesProvider(),
304 > shard.GetSearchAttributesMapperProvider(),
305 > config.SearchAttributesNumberOfKeysLimit,
306 > config.SearchAttributesSizeOfValueLimit,
307 > config.SearchAttributesTotalSizeLimit,
308 > persistenceVisibilityMgr,
309 > visibility.AllowListForValidation(
310 > persistenceVisibilityMgr.GetStoreNames(),
311 > config.VisibilityAllowList,
312 > ),
313 > config.SuppressErrorSetSystemSearchAttribute,
314 > shard.GetMetricsHandler(),
315 > logger,
316 > )
317 >
318 > historyEngImpl.replicationDLQHandler = replication.NewLazyDLQHandler(
319 > shard,
320 > workflowDeleteManager,
321 > workflowCache,
322 > clientBean,
323 > replicationTaskExecutorProvider,
324 > )
325 > historyEngImpl.replicationProcessorMgr = replication.NewTaskProcessorManager(
326 > config,
327 > shard,
328 > historyEngImpl,
329 > workflowCache,
330 > workflowDeleteManager,
331 > clientBean,
332 > serializer,
333 > replicationTaskFetcherFactory,
334 > replicationTaskExecutorProvider,
335 > testHooks,
336 > dlqWriter,
337 > )
338 >
339 > return historyEngImpl
340 }
341
343 // Make sure all the components are loaded lazily so start can return immediately. This is important because
344 // ShardController calls start sequentially for all the shards for a given host during startup.
345 > func (e *historyEngineImpl) Start() { history_engine.go
346 > if !atomic.CompareAndSwapInt32(
347 > &e.status,
348 > common.DaemonStatusInitialized,
349 > common.DaemonStatusStarted,
350 > ) {
351 return
352 }
353
354 > e.logger.Info("", tag.LifeCycleStarting) history_engine.go
355 > defer e.logger.Info("", tag.LifeCycleStarted)
356 >
357 > e.registerNamespaceStateChangeCallback()
358 >
359 > for _, queueProcessor := range e.queueProcessors {
360 > queueProcessor.Start()
361 > }
362 > e.replicationProcessorMgr.Start()
363 }
364
384 }
385
386 > func (e *historyEngineImpl) registerNamespaceStateChangeCallback() { history_engine.go
387 >
388 > e.shardContext.GetNamespaceRegistry().RegisterStateChangeCallback(e, func(ns *namespace.Namespace, deletedFromDb bool) {
389 > if e.shardContext.GetClusterMetadata().IsGlobalNamespaceEnabled() {
390 e.shardContext.UpdateHandoverNamespace(ns, deletedFromDb)
391 }
392
393 > if deletedFromDb { history_engine.go
394 return
395 }
396
397 > if ns.IsGlobalNamespace() && history_engine.go
398 > ns.ReplicationPolicy() == namespace.ReplicationPolicyMultiCluster &&
399 > //nolint:forbidigo // namespace state-change callback; FailoverNamespace operates per-namespace, no workflow context
400 > ns.ActiveInCluster(e.currentClusterName) {
401
402 for _, queueProcessor := range e.queueProcessors {
go.temporal.io/server/service/history/timer_queue_factory.go 123 introduced LOC · 2 ranges

Open complete file

81 func (f *timerQueueFactory) CreateQueue(
82 shardContext historyi.ShardContext,
83 > ) queues.Queue { timer_queue_factory.go
84 > logger := log.With(shardContext.GetLogger(), tag.ComponentTimerQueue)
85 > metricsHandler := f.MetricsHandler.WithTags(metrics.OperationTag(metrics.OperationTimerQueueProcessorScope))
86 >
87 > currentClusterName := f.ClusterMetadata.GetCurrentClusterName()
88 > workflowDeleteManager := deletemanager.NewDeleteManager(
89 > shardContext,
90 > f.WorkflowCache,
91 > f.Config,
92 > shardContext.GetTimeSource(),
93 > f.VisibilityManager,
94 > )
95 >
96 > shardScheduler := queues.NewRateLimitedScheduler(
97 > f.HostScheduler,
98 > queues.RateLimitedSchedulerOptions{
99 > Enabled: f.Config.TaskSchedulerEnableRateLimiter,
100 > EnableShadowMode: f.Config.TaskSchedulerEnableRateLimiterShadowMode,
101 > StartupDelay: f.Config.TaskSchedulerRateLimiterStartupDelay,
102 > },
103 > currentClusterName,
104 > f.NamespaceRegistry,
105 > f.SchedulerRateLimiter,
106 > f.TimeSource,
107 > f.ChasmRegistry,
108 > logger,
109 > metricsHandler,
110 > )
111 >
112 > rescheduler := queues.NewRescheduler(
113 > shardScheduler,
114 > shardContext.GetTimeSource(),
115 > logger,
116 > metricsHandler,
117 > )
118 >
119 > activeExecutor := newTimerQueueActiveTaskExecutor(
120 > shardContext,
121 > f.WorkflowCache,
122 > workflowDeleteManager,
123 > logger,
124 > f.MetricsHandler,
125 > f.Config,
126 > f.MatchingRawClient,
127 > f.ChasmEngine,
128 > )
129 >
130 > standbyExecutor := newTimerQueueStandbyTaskExecutor(
131 > shardContext,
132 > f.WorkflowCache,
133 > workflowDeleteManager,
134 > f.MatchingRawClient,
135 > f.ChasmEngine,
136 > logger,
137 > f.MetricsHandler,
138 > // note: the cluster name is for calculating time for standby tasks,
139 > // here we are basically using current cluster time
140 > // this field will be deprecated soon, currently exists so that
141 > // we have the option of revert to old behavior
142 > currentClusterName,
143 > f.Config,
144 > f.ClientBean,
145 > )
146 >
147 > executor := queues.NewActiveStandbyExecutor(
148 > currentClusterName,
149 > f.NamespaceRegistry,
150 > activeExecutor,
151 > standbyExecutor,
152 > logger,
153 > )
154 > if f.ExecutorWrapper != nil {
155 executor = f.ExecutorWrapper.Wrap(executor)
156 }
157
158 > factory := queues.NewExecutableFactory( timer_queue_factory.go
159 > executor,
160 > shardScheduler,
161 > rescheduler,
162 > f.HostPriorityAssigner,
163 > shardContext.GetTimeSource(),
164 > shardContext.GetNamespaceRegistry(),
165 > shardContext.GetClusterMetadata(),
166 > f.ChasmRegistry,
167 > queues.GetTaskTypeTagValue,
168 > logger,
169 > metricsHandler,
170 > f.Tracer,
171 > f.DLQWriter,
172 > f.Config.TaskDLQEnabled,
173 > f.Config.TaskDLQUnexpectedErrorAttempts,
174 > f.Config.TaskDLQInternalErrors,
175 > f.Config.TaskDLQErrorPattern,
176 > )
177 > return queues.NewScheduledQueue(
178 > shardContext,
179 > tasks.CategoryTimer,
180 > shardScheduler,
181 > rescheduler,
182 > factory,
183 > &queues.Options{
184 > ReaderOptions: queues.ReaderOptions{
185 > BatchSize: f.Config.TimerTaskBatchSize,
186 > MaxPendingTasksCount: f.Config.QueuePendingTaskMaxCount,
187 > PollBackoffInterval: f.Config.TimerProcessorPollBackoffInterval,
188 > MaxPredicateSize: f.Config.QueueMaxPredicateSize,
189 > },
190 > MonitorOptions: queues.MonitorOptions{
191 > PendingTasksCriticalCount: f.Config.QueuePendingTaskCriticalCount,
192 > ReaderStuckCriticalAttempts: f.Config.QueueReaderStuckCriticalAttempts,
193 > SliceCountCriticalThreshold: f.Config.QueueCriticalSlicesCount,
194 > },
195 > MaxPollRPS: f.Config.TimerProcessorMaxPollRPS,
196 > MaxPollInterval: f.Config.TimerProcessorMaxPollInterval,
197 > MaxPollIntervalJitterCoefficient: f.Config.TimerProcessorMaxPollIntervalJitterCoefficient,
198 > CheckpointInterval: f.Config.TimerProcessorUpdateAckInterval,
199 > CheckpointIntervalJitterCoefficient: f.Config.TimerProcessorUpdateAckIntervalJitterCoefficient,
200 > MaxReaderCount: f.Config.TimerQueueMaxReaderCount,
201 > MoveGroupTaskCountBase: f.Config.QueueMoveGroupTaskCountBase,
202 > MoveGroupTaskCountMultiplier: f.Config.QueueMoveGroupTaskCountMultiplier,
203 > ShrinkPredicateMaxPendingKeys: f.Config.QueueShrinkPredicateMaxPendingKeys,
204 > },
205 > f.HostReaderRateLimiter,
206 > logger,
207 > metricsHandler,
208 > )
209 }
go.temporal.io/server/service/matching/fx.go 120 introduced LOC · 16 ranges

Open complete file

56 )
57
58 > func ServerProvider(grpcServerOptions []grpc.ServerOption) *grpc.Server { fx.go
59 > return grpc.NewServer(grpcServerOptions...)
60 > }
61
62 func ConfigProvider(
64 persistenceConfig config.Persistence,
65 rateLimitFractionProvider TaskQueueRateLimitFractionProvider,
66 > ) *Config { fx.go
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
75 > return interceptor.NewServiceErrorInterceptor(
76 > dynamicconfig.MaxServiceErrorMessageLength.Get(dc),
77 > )
78 > }
79
80 > func RetryableInterceptorProvider() *interceptor.RetryableInterceptor { fx.go
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
91 > return interceptor.NewRequestErrorHandler(
92 > logger,
93 > serviceConfig.LogAllReqErrors,
94 > )
95 > }
96
97 func TelemetryInterceptorProvider(
101 serviceConfig *Config,
102 requestErrorHandler *interceptor.RequestErrorHandler,
103 > ) *interceptor.TelemetryInterceptor { fx.go
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 {
119 namespaceRegistry namespace.Registry,
120 metricsHandler metrics.Handler,
121 > ) interceptor.NamespaceRateLimitInterceptor { fx.go
122 >
123 > namespaceRateFn := func(namespaceName string) float64 {
124 if namespaceRPS := serviceConfig.NamespaceRPS(namespaceName); namespaceRPS > 0 {
125 return float64(namespaceRPS)
129 }
130
131 > return interceptor.NewNamespaceRateLimitInterceptor( fx.go
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
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.
159 persistenceLazyLoadedServiceResolver service.PersistenceLazyLoadedServiceResolver,
160 logger log.SnTaggedLogger,
161 > ) service.PersistenceRateLimitingParams { fx.go
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
179 > return membershipMonitor.GetResolver(primitives.MatchingService)
180 > }
181
182 // TaskQueueReplicatorNamespaceReplicationQueue is used to ensure the replicator only gets set if global namespaces are
207 chasmRegistry *chasm.Registry,
208 serializer serialization.Serializer,
209 > ) (manager.VisibilityManager, error) { fx.go
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
236 > return interceptor.NewContextMetadataInterceptor(true, logger)
237 > }
238
239 > func ServiceLifetimeHooks(lc fx.Lifecycle, svc *Service) { fx.go
240 > lc.Append(fx.StartStopHook(svc.Start, svc.Stop))
241 > }
242
243 func WorkersRegistryProvider(
245 metricsHandler metrics.Handler,
246 serviceConfig *Config,
247 > ) workers.Registry { fx.go
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
265 > return newSimplePartitionScalerFactory(
266 > dynamicconfig.MatchingPartitionScaler.Get(dc),
267 > )
268 > }
go.temporal.io/server/service/history/transfer_queue_factory.go 118 introduced LOC · 2 ranges

Open complete file

88 func (f *transferQueueFactory) CreateQueue(
89 shardContext historyi.ShardContext,
90 > ) queues.Queue { transfer_queue_factory.go
91 > logger := log.With(shardContext.GetLogger(), tag.ComponentTransferQueue)
92 > metricsHandler := f.MetricsHandler.WithTags(metrics.OperationTag(metrics.OperationTransferQueueProcessorScope))
93 >
94 > currentClusterName := f.ClusterMetadata.GetCurrentClusterName()
95 >
96 > shardScheduler := queues.NewRateLimitedScheduler(
97 > f.HostScheduler,
98 > queues.RateLimitedSchedulerOptions{
99 > Enabled: f.Config.TaskSchedulerEnableRateLimiter,
100 > EnableShadowMode: f.Config.TaskSchedulerEnableRateLimiterShadowMode,
101 > StartupDelay: f.Config.TaskSchedulerRateLimiterStartupDelay,
102 > },
103 > currentClusterName,
104 > f.NamespaceRegistry,
105 > f.SchedulerRateLimiter,
106 > f.TimeSource,
107 > f.ChasmRegistry,
108 > logger,
109 > metricsHandler,
110 > )
111 >
112 > rescheduler := queues.NewRescheduler(
113 > shardScheduler,
114 > shardContext.GetTimeSource(),
115 > logger,
116 > metricsHandler,
117 > )
118 >
119 > activeExecutor := newTransferQueueActiveTaskExecutor(
120 > shardContext,
121 > f.WorkflowCache,
122 > f.SdkClientFactory,
123 > logger,
124 > f.MetricsHandler,
125 > f.Config,
126 > f.HistoryRawClient,
127 > f.MatchingRawClient,
128 > f.VisibilityManager,
129 > f.ChasmEngine,
130 > f.VersionMembershipCache,
131 > f.TestHooks,
132 > )
133 >
134 > standbyExecutor := newTransferQueueStandbyTaskExecutor(
135 > shardContext,
136 > f.WorkflowCache,
137 > logger,
138 > f.MetricsHandler,
139 > currentClusterName,
140 > f.HistoryRawClient,
141 > f.MatchingRawClient,
142 > f.VisibilityManager,
143 > f.ChasmEngine,
144 > f.ClientBean,
145 > )
146 >
147 > executor := queues.NewActiveStandbyExecutor(
148 > currentClusterName,
149 > f.NamespaceRegistry,
150 > activeExecutor,
151 > standbyExecutor,
152 > logger,
153 > )
154 > if f.ExecutorWrapper != nil {
155 executor = f.ExecutorWrapper.Wrap(executor)
156 }
157
158 > factory := queues.NewExecutableFactory( transfer_queue_factory.go
159 > executor,
160 > shardScheduler,
161 > rescheduler,
162 > f.HostPriorityAssigner,
163 > shardContext.GetTimeSource(),
164 > shardContext.GetNamespaceRegistry(),
165 > shardContext.GetClusterMetadata(),
166 > f.ChasmRegistry,
167 > queues.GetTaskTypeTagValue,
168 > logger,
169 > metricsHandler,
170 > f.Tracer,
171 > f.DLQWriter,
172 > f.Config.TaskDLQEnabled,
173 > f.Config.TaskDLQUnexpectedErrorAttempts,
174 > f.Config.TaskDLQInternalErrors,
175 > f.Config.TaskDLQErrorPattern,
176 > )
177 > return queues.NewImmediateQueue(
178 > shardContext,
179 > tasks.CategoryTransfer,
180 > shardScheduler,
181 > rescheduler,
182 > &queues.Options{
183 > ReaderOptions: queues.ReaderOptions{
184 > BatchSize: f.Config.TransferTaskBatchSize,
185 > MaxPendingTasksCount: f.Config.QueuePendingTaskMaxCount,
186 > PollBackoffInterval: f.Config.TransferProcessorPollBackoffInterval,
187 > MaxPredicateSize: f.Config.QueueMaxPredicateSize,
188 > },
189 > MonitorOptions: queues.MonitorOptions{
190 > PendingTasksCriticalCount: f.Config.QueuePendingTaskCriticalCount,
191 > ReaderStuckCriticalAttempts: f.Config.QueueReaderStuckCriticalAttempts,
192 > SliceCountCriticalThreshold: f.Config.QueueCriticalSlicesCount,
193 > },
194 > MaxPollRPS: f.Config.TransferProcessorMaxPollRPS,
195 > MaxPollInterval: f.Config.TransferProcessorMaxPollInterval,
196 > MaxPollIntervalJitterCoefficient: f.Config.TransferProcessorMaxPollIntervalJitterCoefficient,
197 > CheckpointInterval: f.Config.TransferProcessorUpdateAckInterval,
198 > CheckpointIntervalJitterCoefficient: f.Config.TransferProcessorUpdateAckIntervalJitterCoefficient,
199 > MaxReaderCount: f.Config.TransferQueueMaxReaderCount,
200 > MoveGroupTaskCountBase: f.Config.QueueMoveGroupTaskCountBase,
201 > MoveGroupTaskCountMultiplier: f.Config.QueueMoveGroupTaskCountMultiplier,
202 > ShrinkPredicateMaxPendingKeys: f.Config.QueueShrinkPredicateMaxPendingKeys,
203 > },
204 > f.HostReaderRateLimiter,
205 > queues.GrouperNamespaceID{},
206 > logger,
207 > metricsHandler,
208 > factory,
209 > nil, // taskPostProcessor
210 > )
211 }
go.temporal.io/server/service/worker/service.go 110 introduced LOC · 12 ranges

Open complete file

134 grpcListener net.Listener,
135 healthServer *health.Server,
136 > ) (*Service, error) { service.go
137 > workerServiceResolver, err := membershipMonitor.GetResolver(primitives.WorkerService)
138 > if err != nil {
139 return nil, err
140 }
141
142 > s := &Service{ service.go
143 > config: serviceConfig,
144 > sdkClientFactory: sdkClientFactory,
145 > logger: logger,
146 > clusterMetadata: clusterMetadata,
147 > clientBean: clientBean,
148 > clusterMetadataManager: clusterMetadataManager,
149 > namespaceRegistry: namespaceRegistry,
150 > executionManager: executionManager,
151 > workerServiceResolver: workerServiceResolver,
152 > membershipMonitor: membershipMonitor,
153 > hostInfo: hostInfoProvider.HostInfo(),
154 > namespaceReplicationQueue: namespaceReplicationQueue,
155 > metricsHandler: metricsHandler,
156 > metadataManager: metadataManager,
157 > taskManager: taskManager,
158 > historyClient: historyClient,
159 > visibilityManager: visibilityManager,
160 >
161 > workerManager: workerManager,
162 > perNamespaceWorkerManager: perNamespaceWorkerManager,
163 > matchingClient: matchingClient,
164 > namespaceReplicationTaskExecutor: namespaceReplicationTaskExecutor,
165 >
166 > server: server,
167 > grpcListener: grpcListener,
168 > healthServer: healthServer,
169 > }
170 > if err := s.initScanner(serializer); err != nil {
171 return nil, err
172 }
173 > return s, nil service.go
174 }
175
245
246 // Start is called to start the service
247 > func (s *Service) Start() { service.go
248 > s.logger.Info(
249 > "worker starting",
250 > tag.ComponentWorker,
251 > )
252 >
253 > metrics.RestartCount.With(s.metricsHandler).Record(1)
254 >
255 > s.membershipMonitor.Start()
256 >
257 > s.ensureSystemNamespaceExists(context.TODO())
258 > s.startScanner()
259 >
260 > if s.clusterMetadata.IsGlobalNamespaceEnabled() {
261 s.startReplicator()
262 }
263 > if s.config.EnableParentClosePolicyWorker() { service.go
264 > s.startParentClosePolicyProcessor()
265 > }
266
267 > s.workerManager.Start() service.go
268 > s.perNamespaceWorkerManager.Start(
269 > // TODO: get these from fx instead of passing through Start
270 > s.hostInfo,
271 > s.workerServiceResolver,
272 > )
273 >
274 > healthpb.RegisterHealthServer(s.server, s.healthServer)
275 > s.healthServer.SetServingStatus(ServiceName, healthpb.HealthCheckResponse_SERVING)
276 >
277 > reflection.Register(s.server)
278 >
279 > go func() {
280 > s.logger.Info("Starting to serve on worker listener")
281 > if err := s.server.Serve(s.grpcListener); err != nil {
282 s.logger.Fatal("Failed to serve on worker listener", tag.Error(err))
283 }
284 }()
285
286 > s.logger.Info( service.go
287 > "worker service started",
288 > tag.ComponentWorker,
289 > tag.Address(s.hostInfo.GetAddress()),
290 > )
291 }
292
309 }
310
311 > func (s *Service) startParentClosePolicyProcessor() { service.go
312 > params := &parentclosepolicy.BootstrapParams{
313 > Config: *s.config.ParentCloseCfg,
314 > SdkClientFactory: s.sdkClientFactory,
315 > MetricsHandler: s.metricsHandler,
316 > Logger: s.logger,
317 > ClientBean: s.clientBean,
318 > CurrentCluster: s.clusterMetadata.GetCurrentClusterName(),
319 > HostInfo: s.hostInfo,
320 > }
321 > processor := parentclosepolicy.New(params)
322 > if err := processor.Start(); err != nil {
323 s.logger.Fatal(
324 "error starting parentclosepolicy processor",
328 }
329
330 > func (s *Service) initScanner(serializer serialization.Serializer) error { service.go
331 > currentCluster := s.clusterMetadata.GetCurrentClusterName()
332 > adminClient, err := s.clientBean.GetRemoteAdminClient(currentCluster)
333 > if err != nil {
334 return err
335 }
336 > s.scanner = scanner.New( service.go
337 > s.logger,
338 > s.config.ScannerCfg,
339 > s.sdkClientFactory,
340 > s.metricsHandler,
341 > s.executionManager,
342 > s.metadataManager,
343 > s.visibilityManager,
344 > s.taskManager,
345 > s.historyClient,
346 > adminClient,
347 > s.matchingClient,
348 > s.namespaceRegistry,
349 > currentCluster,
350 > s.hostInfo,
351 > serializer,
352 > )
353 > return nil
354 }
355
356 > func (s *Service) startScanner() { service.go
357 > if err := s.scanner.Start(); err != nil {
358 s.logger.Fatal(
359 "error starting scanner",
384 func (s *Service) ensureSystemNamespaceExists(
385 ctx context.Context,
386 > ) { service.go
387 > _, err := s.metadataManager.GetNamespace(ctx, &persistence.GetNamespaceRequest{Name: primitives.SystemLocalNamespace})
388 > switch err.(type) {
389 > case nil:
390 // noop
391 case *serviceerror.NamespaceNotFound:
go.temporal.io/server/service/history/visibility_queue_factory.go 92 introduced LOC · 2 ranges

Open complete file

77 func (f *visibilityQueueFactory) CreateQueue(
78 shard historyi.ShardContext,
79 > ) queues.Queue { visibility_queue_factory.go
80 > logger := log.With(shard.GetLogger(), tag.ComponentVisibilityQueue)
81 > metricsHandler := f.MetricsHandler.WithTags(metrics.OperationTag(metrics.OperationVisibilityQueueProcessorScope))
82 >
83 > shardScheduler := queues.NewRateLimitedScheduler(
84 > f.HostScheduler,
85 > queues.RateLimitedSchedulerOptions{
86 > Enabled: f.Config.TaskSchedulerEnableRateLimiter,
87 > EnableShadowMode: f.Config.TaskSchedulerEnableRateLimiterShadowMode,
88 > StartupDelay: f.Config.TaskSchedulerRateLimiterStartupDelay,
89 > },
90 > f.ClusterMetadata.GetCurrentClusterName(),
91 > f.NamespaceRegistry,
92 > f.SchedulerRateLimiter,
93 > f.TimeSource,
94 > f.ChasmRegistry,
95 > logger,
96 > metricsHandler,
97 > )
98 >
99 > rescheduler := queues.NewRescheduler(
100 > shardScheduler,
101 > shard.GetTimeSource(),
102 > logger,
103 > metricsHandler,
104 > )
105 >
106 > executor := newVisibilityQueueTaskExecutor(
107 > shard,
108 > f.WorkflowCache,
109 > f.VisibilityMgr,
110 > logger,
111 > f.MetricsHandler,
112 > f.Config.VisibilityProcessorEnsureCloseBeforeDelete,
113 > f.Config.VisibilityProcessorEnableCloseWorkflowCleanup,
114 > f.Config.VisibilityProcessorRelocateAttributesMinBlobSize,
115 > f.Config.ExternalPayloadsEnabled,
116 > )
117 > if f.ExecutorWrapper != nil {
118 executor = f.ExecutorWrapper.Wrap(executor)
119 }
120
121 > factory := queues.NewExecutableFactory( visibility_queue_factory.go
122 > executor,
123 > shardScheduler,
124 > rescheduler,
125 > f.HostPriorityAssigner,
126 > shard.GetTimeSource(),
127 > shard.GetNamespaceRegistry(),
128 > shard.GetClusterMetadata(),
129 > f.ChasmRegistry,
130 > queues.GetTaskTypeTagValue,
131 > logger,
132 > metricsHandler,
133 > f.Tracer,
134 > f.DLQWriter,
135 > f.Config.TaskDLQEnabled,
136 > f.Config.TaskDLQUnexpectedErrorAttempts,
137 > f.Config.TaskDLQInternalErrors,
138 > f.Config.TaskDLQErrorPattern,
139 > )
140 > return queues.NewImmediateQueue(
141 > shard,
142 > tasks.CategoryVisibility,
143 > shardScheduler,
144 > rescheduler,
145 > &queues.Options{
146 > ReaderOptions: queues.ReaderOptions{
147 > BatchSize: f.Config.VisibilityTaskBatchSize,
148 > MaxPendingTasksCount: f.Config.QueuePendingTaskMaxCount,
149 > PollBackoffInterval: f.Config.VisibilityProcessorPollBackoffInterval,
150 > MaxPredicateSize: f.Config.QueueMaxPredicateSize,
151 > },
152 > MonitorOptions: queues.MonitorOptions{
153 > PendingTasksCriticalCount: f.Config.QueuePendingTaskCriticalCount,
154 > ReaderStuckCriticalAttempts: f.Config.QueueReaderStuckCriticalAttempts,
155 > SliceCountCriticalThreshold: f.Config.QueueCriticalSlicesCount,
156 > },
157 > MaxPollRPS: f.Config.VisibilityProcessorMaxPollRPS,
158 > MaxPollInterval: f.Config.VisibilityProcessorMaxPollInterval,
159 > MaxPollIntervalJitterCoefficient: f.Config.VisibilityProcessorMaxPollIntervalJitterCoefficient,
160 > CheckpointInterval: f.Config.VisibilityProcessorUpdateAckInterval,
161 > CheckpointIntervalJitterCoefficient: f.Config.VisibilityProcessorUpdateAckIntervalJitterCoefficient,
162 > MaxReaderCount: f.Config.VisibilityQueueMaxReaderCount,
163 > MoveGroupTaskCountBase: f.Config.QueueMoveGroupTaskCountBase,
164 > MoveGroupTaskCountMultiplier: f.Config.QueueMoveGroupTaskCountMultiplier,
165 > ShrinkPredicateMaxPendingKeys: f.Config.QueueShrinkPredicateMaxPendingKeys,
166 > },
167 > f.HostReaderRateLimiter,
168 > queues.GrouperNamespaceID{},
169 > logger,
170 > metricsHandler,
171 > factory,
172 > nil, // taskPostProcessor
173 > )
174 }
go.temporal.io/server/client/clientfactory.go 81 introduced LOC · 13 ranges

Open complete file

73
74 // NewFactoryProvider creates a default implementation of FactoryProvider.
75 > func NewFactoryProvider() FactoryProvider { clientfactory.go
76 > return &factoryProviderImpl{}
77 > }
78
79 // NewFactory creates an instance of client factory that knows how to dispatch RPC calls.
87 logger log.Logger,
88 throttledLogger log.Logger,
89 > ) Factory { clientfactory.go
90 > return &rpcClientFactory{
91 > rpcFactory: rpcFactory,
92 > monitor: monitor,
93 > metricsHandler: metricsHandler,
94 > dynConfig: dc,
95 > testHooks: testHooks,
96 > numberOfHistoryShards: numberOfHistoryShards,
97 > logger: logger,
98 > throttledLogger: throttledLogger,
99 > }
100 > }
101
102 > func (cf *rpcClientFactory) NewHistoryClientWithTimeout(timeout time.Duration) (historyservice.HistoryServiceClient, error) { clientfactory.go
103 > resolver, err := cf.monitor.GetResolver(primitives.HistoryService)
104 > if err != nil {
105 return nil, err
106 }
107 > client := history.NewClient( clientfactory.go
108 > cf.dynConfig,
109 > resolver,
110 > cf.logger,
111 > cf.numberOfHistoryShards,
112 > cf.rpcFactory,
113 > timeout,
114 > )
115 > if cf.metricsHandler != nil {
116 > client = history.NewMetricClient(client, cf.metricsHandler, cf.logger, cf.throttledLogger)
117 > }
118 > return client, nil
119 }
120
123 timeout time.Duration,
124 longPollTimeout time.Duration,
125 > ) (matchingservice.MatchingServiceClient, error) { clientfactory.go
126 > resolver, err := cf.monitor.GetResolver(primitives.MatchingService)
127 > if err != nil {
128 return nil, err
129 }
130
131 > keyResolver := newServiceKeyResolver(resolver) clientfactory.go
132 > clientProvider := func(clientKey string) (any, func() error, error) {
133 connection := cf.rpcFactory.CreateMatchingGRPCConnection(clientKey)
134 return matchingservice.NewMatchingServiceClient(connection), connection.Close, nil
135 }
136 > client := matching.NewClient( clientfactory.go
137 > timeout,
138 > longPollTimeout,
139 > common.NewClientCache(keyResolver, clientProvider, cf.logger),
140 > cf.metricsHandler,
141 > cf.logger,
142 > matching.NewLoadBalancer(namespaceIDToName, cf.dynConfig, cf.testHooks),
143 > dynamicconfig.MatchingSpreadRoutingBatchSize.Get(cf.dynConfig),
144 > resolver,
145 > dynamicconfig.MatchingConnectionCloseDelay.Get(cf.dynConfig),
146 > )
147 >
148 > if cf.metricsHandler != nil {
149 > client = matching.NewMetricClient(client, cf.metricsHandler, cf.logger, cf.throttledLogger)
150 > }
151 > return client, nil
152
153 }
166 timeout time.Duration,
167 longPollTimeout time.Duration,
168 > ) (grpc.ClientConnInterface, workflowservice.WorkflowServiceClient, error) { clientfactory.go
169 > connection := cf.rpcFactory.CreateLocalFrontendGRPCConnection()
170 > client := workflowservice.NewWorkflowServiceClient(connection)
171 > return connection, cf.newFrontendClient(client, timeout, longPollTimeout), nil
172 > }
173
174 func (cf *rpcClientFactory) NewRemoteAdminClientWithTimeout(
185 timeout time.Duration,
186 longPollTimeout time.Duration,
187 > ) (adminservice.AdminServiceClient, error) { clientfactory.go
188 > connection := cf.rpcFactory.CreateLocalFrontendGRPCConnection()
189 > client := adminservice.NewAdminServiceClient(connection)
190 > return cf.newAdminClient(client, timeout, longPollTimeout), nil
191 > }
192
193 func (cf *rpcClientFactory) newAdminClient(
195 timeout time.Duration,
196 longPollTimeout time.Duration,
197 > ) adminservice.AdminServiceClient { clientfactory.go
198 > client = admin.NewClient(timeout, longPollTimeout, client)
199 > if cf.metricsHandler != nil {
200 > client = admin.NewMetricClient(client, cf.metricsHandler, cf.throttledLogger)
201 > }
202 > return client
203 }
204
207 timeout time.Duration,
208 longPollTimeout time.Duration,
209 > ) workflowservice.WorkflowServiceClient { clientfactory.go
210 > client = frontend.NewClient(timeout, longPollTimeout, client)
211 > if cf.metricsHandler != nil {
212 > client = frontend.NewMetricClient(client, cf.metricsHandler, cf.throttledLogger)
213 > }
214 > return client
215 }
216
217 > func newServiceKeyResolver(resolver membership.ServiceResolver) *serviceKeyResolverImpl { clientfactory.go
218 > return &serviceKeyResolverImpl{
219 > resolver: resolver,
220 > }
221 > }
222
223 // Lookup returns the address for a node within a batch. key contains the key (including batch
224 // number), and index is the index within the batch. If not using batches, index should be 0.
225 // Note that Lookup(key) and LookupN(key, n)[0] are equal.
226 > func (r *serviceKeyResolverImpl) Lookup(key string, index int) (string, error) { clientfactory.go
227 > hosts := r.resolver.LookupN(key, index+1)
228 > if len(hosts) == 0 {
229 return "", membership.ErrInsufficientHosts
230 }
go.temporal.io/server/common/persistence/client/fx.go 81 introduced LOC · 10 ranges

Open complete file

85 )
86
87 > func ClusterNameProvider(config *cluster.Config) ClusterName { fx.go
88 > return ClusterName(config.CurrentClusterName)
89 > }
90
91 func EventBlobCacheProvider(
93 logger log.Logger,
94 serializer serialization.Serializer,
95 > ) persistence.XDCCache { fx.go
96 > return persistence.NewEventsBlobCache(
97 > dynamicconfig.XDCCacheMaxSizeBytes.Get(dc)(),
98 > 20*time.Second,
99 > logger,
100 > )
101 > }
102
103 func EnableDataLossMetricsProvider(
104 dc *dynamicconfig.Collection,
105 > ) EnableDataLossMetrics { fx.go
106 > return EnableDataLossMetrics(dynamicconfig.EnableDataLossMetrics.Get(dc))
107 > }
108
109 func EnableBestEffortDeleteTasksOnWorkflowUpdateProvider(
110 dc *dynamicconfig.Collection,
111 > ) EnableBestEffortDeleteTasksOnWorkflowUpdate { fx.go
112 > return EnableBestEffortDeleteTasksOnWorkflowUpdate(dynamicconfig.EnableBestEffortDeleteTasksOnWorkflowUpdate.Get(dc))
113 > }
114
115 func FactoryProvider(
116 params NewFactoryParams,
117 > ) Factory { fx.go
118 > var systemRequestRateLimiter, namespaceRequestRateLimiter, shardRequestRateLimiter quotas.RequestRateLimiter
119 > if params.PersistenceMaxQPS != nil && params.PersistenceMaxQPS() > 0 {
120 > systemRequestRateLimiter = NewPriorityRateLimiter(
121 > params.PersistenceMaxQPS,
122 > RequestPriorityFn,
123 > params.OperatorRPSRatio,
124 > params.PersistenceBurstRatio,
125 > params.HealthSignals,
126 > params.DynamicRateLimitingParams,
127 > params.MetricsHandler,
128 > params.Logger,
129 > )
130 > namespaceRequestRateLimiter = NewPriorityNamespaceRateLimiter(
131 > params.PersistenceMaxQPS,
132 > params.PersistenceNamespaceMaxQPS,
133 > RequestPriorityFn,
134 > params.OperatorRPSRatio,
135 > params.PersistenceBurstRatio,
136 > )
137 > shardRequestRateLimiter = NewPriorityNamespaceShardRateLimiter(
138 > params.PersistenceMaxQPS,
139 > params.PersistencePerShardNamespaceMaxQPS,
140 > RequestPriorityFn,
141 > params.OperatorRPSRatio,
142 > params.PersistenceBurstRatio,
143 > )
144 > }
145
146 > return NewFactory( fx.go
147 > params.DataStoreFactory,
148 > params.Cfg,
149 > systemRequestRateLimiter,
150 > namespaceRequestRateLimiter,
151 > shardRequestRateLimiter,
152 > params.Serializer,
153 > params.EventBlobCache,
154 > string(params.ClusterName),
155 > params.MetricsHandler,
156 > params.Logger,
157 > params.HealthSignals,
158 > params.EnableDataLossMetrics,
159 > params.EnableBestEffortDeleteTasksOnWorkflowUpdate,
160 > )
161 }
162
166 metricsHandler metrics.Handler,
167 logger log.ThrottledLogger,
168 > ) persistence.HealthSignalAggregator { fx.go
169 > if dynamicconfig.PersistenceHealthSignalMetricsEnabled.Get(dynamicCollection)() {
170 > aggregator := persistence.NewHealthSignalAggregator(
171 > dynamicconfig.PersistenceHealthSignalAggregationEnabled.Get(dynamicCollection)(),
172 > dynamicconfig.PersistenceHealthSignalPercentilesEnabled.Get(dynamicCollection),
173 > dynamicconfig.PersistenceHealthSignalWindowSize.Get(dynamicCollection)(),
174 > dynamicconfig.PersistenceHealthSignalBufferSize.Get(dynamicCollection)(),
175 > metricsHandler,
176 > logger,
177 > dynamicconfig.PersistenceHealthSignalLatencyWindowSize.Get(dynamicCollection)(),
178 > dynamicconfig.PersistenceHealthSignalLatencyWindowCount.Get(dynamicCollection)(),
179 > )
180 > lc.Append(fx.StopHook(aggregator.Stop))
181 > return aggregator
182 > }
183
184 return persistence.NoopHealthSignalAggregator
220 }
221
222 > func DataStoreFactoryLifetimeHooks(lc fx.Lifecycle, f persistence.DataStoreFactory) { fx.go
223 > lc.Append(fx.StopHook(f.Close))
224 > }
225
226 func managerProvider[T persistence.Closeable](newManagerFn func(Factory) (T, error)) func(Factory, fx.Lifecycle) (T, error) {
227 return func(f Factory, lc fx.Lifecycle) (T, error) {
228 > manager, err := newManagerFn(f) // passing receiver (Factory) as first argument. fx.go
229 > if err != nil {
230 var unimpl *serviceerror.Unimplemented
231 if errors.As(err, &unimpl) {
236 return nilT, err
237 }
238 > lc.Append(fx.StopHook(manager.Close)) fx.go
239 > return manager, nil
240 }
241 }
go.temporal.io/server/common/persistence/client/quotas.go 79 introduced LOC · 10 ranges

Open complete file

70 metricsHandler metrics.Handler,
71 logger log.Logger,
72 > ) quotas.RequestRateLimiter { quotas.go
73 > hostRateFn := func() float64 { return float64(hostMaxQPS()) }
74
75 > return quotas.NewMultiRequestRateLimiter( quotas.go
76 > // host-level dynamic rate limiter
77 > newPriorityDynamicRateLimiter(
78 > hostRateFn,
79 > requestPriorityFn,
80 > operatorRPSRatio,
81 > burstRatio,
82 > healthSignals,
83 > dynamicParams,
84 > metricsHandler,
85 > logger,
86 > ),
87 > // basic host-level rate limiter
88 > newPriorityRateLimiter(
89 > hostRateFn,
90 > requestPriorityFn,
91 > operatorRPSRatio,
92 > burstRatio,
93 > ),
94 > )
95 }
96
101 operatorRPSRatio OperatorRPSRatio,
102 burstRatio PersistenceBurstRatio,
103 > ) quotas.RequestRateLimiter { quotas.go
104 >
105 > return newPriorityNamespaceRateLimiter(
106 > namespaceMaxQPS,
107 > hostMaxQPS,
108 > requestPriorityFn,
109 > operatorRPSRatio,
110 > burstRatio,
111 > )
112 > }
113
114 func NewPriorityNamespaceShardRateLimiter(
118 operatorRPSRatio OperatorRPSRatio,
119 burstRatio PersistenceBurstRatio,
120 > ) quotas.RequestRateLimiter { quotas.go
121 >
122 > return newPerShardPerNamespacePriorityRateLimiter(
123 > perShardNamespaceMaxQPS,
124 > hostMaxQPS,
125 > requestPriorityFn,
126 > operatorRPSRatio,
127 > burstRatio,
128 > )
129 > }
130
131 func newPerShardPerNamespacePriorityRateLimiter(
149 )
150 }
151 > return quotas.NoopRequestRateLimiter quotas.go
152 },
153 perShardPerNamespaceKeyFn,
189 )
190 }
191 > return quotas.NoopRequestRateLimiter quotas.go
192 })
193 }
233 metricsHandler metrics.Handler,
234 logger log.Logger,
235 > ) quotas.RequestRateLimiter { quotas.go
236 > rateLimiters := make(map[int]quotas.RequestRateLimiter)
237 > for priority := range RequestPrioritiesOrdered {
238 > // TODO: refactor this so dynamic rate adjustment is global for all priorities
239 > if priority == CallerTypeDefaultPriority[headers.CallerTypeOperator] {
240 > rateLimiters[priority] = NewHealthRequestRateLimiterImpl(
241 > healthSignals,
242 > operatorRateFn(rateFn, operatorRPSRatio),
243 > dynamicParams,
244 > burstRatio,
245 > metricsHandler,
246 > logger,
247 > )
248 > } else {
249 > rateLimiters[priority] = NewHealthRequestRateLimiterImpl(
250 > healthSignals,
251 > rateFn,
252 > dynamicParams,
253 > burstRatio,
254 > metricsHandler,
255 > logger,
256 > )
257 > }
258 }
259
260 > return quotas.NewPriorityRateLimiter( quotas.go
261 > requestPriorityFn,
262 > rateLimiters,
263 > )
264 }
265
273 }
274 return CallerTypeDefaultPriority[req.CallerType]
275 > case headers.CallerTypeBackgroundHigh, headers.CallerTypeBackgroundLow: quotas.go
276 > if priority, ok := BackgroundTypeAPIPriorityOverride[req.API]; ok {
277 > return priority
278 > }
279 > return CallerTypeDefaultPriority[req.CallerType]
280 case headers.CallerTypePreemptable:
281 return CallerTypeDefaultPriority[req.CallerType]
282 > default: quotas.go
283 > // default requests to API priority to be consistent with existing behavior
284 > return CallerTypeDefaultPriority[headers.CallerTypeAPI]
285 }
286 }
go.temporal.io/server/service/worker/deletenamespace/fx.go 77 introduced LOC · 10 ranges

Open complete file

57 func newComponent(
58 params componentParams,
59 > ) workercommon.WorkerComponent { fx.go
60 > return &deleteNamespaceComponent{
61 > atWorkerCfg: dynamicconfig.WorkerDeleteNamespaceActivityLimits.Get(params.DynamicCollection)(),
62 > visibilityManager: params.VisibilityManager,
63 > metadataManager: params.MetadataManager,
64 > clusterMetadata: params.ClusterMetadata,
65 > nexusEndpointManager: params.NexusEndpointManager,
66 > historyClient: params.HistoryClient,
67 > metricsHandler: params.MetricsHandler,
68 > logger: params.Logger,
69 > protectedNamespaces: dynamicconfig.ProtectedNamespaces.Get(params.DynamicCollection),
70 > allowDeleteNamespaceIfNexusEndpointTarget: dynamicconfig.AllowDeleteNamespaceIfNexusEndpointTarget.Get(params.DynamicCollection),
71 > nexusEndpointListDefaultPageSize: dynamicconfig.NexusEndpointListDefaultPageSize.Get(params.DynamicCollection),
72 > deleteActivityRPS: dynamicconfig.DeleteNamespaceDeleteActivityRPS.Subscribe(params.DynamicCollection),
73 > useChasmDeleteExecution: dynamicconfig.DeleteNamespaceUseChasmDeleteExecution.Get(params.DynamicCollection),
74 > namespaceCacheRefreshInterval: dynamicconfig.NamespaceCacheRefreshInterval.Get(params.DynamicCollection),
75 > }
76 > }
77
78 > func (wc *deleteNamespaceComponent) RegisterWorkflow(registry sdkworker.Registry) { fx.go
79 > registry.RegisterWorkflowWithOptions(DeleteNamespaceWorkflow, workflow.RegisterOptions{Name: WorkflowName})
80 > registry.RegisterActivity(wc.deleteNamespaceLocalActivities())
81 >
82 > registry.RegisterWorkflowWithOptions(reclaimresources.ReclaimResourcesWorkflow, workflow.RegisterOptions{Name: reclaimresources.WorkflowName})
83 > registry.RegisterActivity(wc.reclaimResourcesLocalActivities())
84 >
85 > registry.RegisterWorkflowWithOptions(deleteexecutions.DeleteExecutionsWorkflow, workflow.RegisterOptions{Name: deleteexecutions.WorkflowName})
86 > registry.RegisterActivity(wc.deleteExecutionsLocalActivities())
87 > }
88
89 > func (wc *deleteNamespaceComponent) DedicatedWorkflowWorkerOptions() *workercommon.DedicatedWorkerOptions { fx.go
90 > // use default worker
91 > return nil
92 > }
93
94 > func (wc *deleteNamespaceComponent) RegisterActivities(registry sdkworker.Registry) { fx.go
95 > registry.RegisterActivity(wc.reclaimResourcesActivities())
96 > registry.RegisterActivity(wc.deleteExecutionsActivities())
97 > }
98
99 > func (wc *deleteNamespaceComponent) DedicatedActivityWorkerOptions() *workercommon.DedicatedWorkerOptions { fx.go
100 > return &workercommon.DedicatedWorkerOptions{
101 > TaskQueue: primitives.DeleteNamespaceActivityTQ,
102 > Options: sdkworker.Options{
103 > BackgroundActivityContext: headers.SetCallerType(context.Background(), headers.CallerTypePreemptable),
104 > MaxConcurrentActivityExecutionSize: wc.atWorkerCfg.MaxConcurrentActivityExecutionSize,
105 > TaskQueueActivitiesPerSecond: wc.atWorkerCfg.TaskQueueActivitiesPerSecond,
106 > WorkerActivitiesPerSecond: wc.atWorkerCfg.WorkerActivitiesPerSecond,
107 > MaxConcurrentActivityTaskPollers: wc.atWorkerCfg.MaxConcurrentActivityTaskPollers,
108 > },
109 > }
110 > }
111
112 > func (wc *deleteNamespaceComponent) deleteNamespaceLocalActivities() *localActivities { fx.go
113 > return newLocalActivities(
114 > wc.metadataManager,
115 > wc.clusterMetadata,
116 > wc.nexusEndpointManager,
117 > wc.logger,
118 > wc.protectedNamespaces,
119 > wc.allowDeleteNamespaceIfNexusEndpointTarget,
120 > wc.nexusEndpointListDefaultPageSize)
121 > }
122
123 > func (wc *deleteNamespaceComponent) reclaimResourcesActivities() *reclaimresources.Activities { fx.go
124 > return reclaimresources.NewActivities(wc.visibilityManager, wc.logger)
125 > }
126
127 > func (wc *deleteNamespaceComponent) reclaimResourcesLocalActivities() *reclaimresources.LocalActivities { fx.go
128 > return reclaimresources.NewLocalActivities(wc.visibilityManager, wc.metadataManager, wc.namespaceCacheRefreshInterval, wc.logger)
129 > }
130
131 > func (wc *deleteNamespaceComponent) deleteExecutionsActivities() *deleteexecutions.Activities { fx.go
132 > return deleteexecutions.NewActivities(
133 > wc.visibilityManager,
134 > wc.historyClient,
135 > wc.deleteActivityRPS,
136 > wc.useChasmDeleteExecution,
137 > wc.metricsHandler,
138 > wc.logger,
139 > )
140 > }
141
142 > func (wc *deleteNamespaceComponent) deleteExecutionsLocalActivities() *deleteexecutions.LocalActivities { fx.go
143 > return deleteexecutions.NewLocalActivities(wc.visibilityManager, wc.metricsHandler, wc.logger)
144 > }
go.temporal.io/server/service/worker/fx.go 77 introduced LOC · 10 ranges

Open complete file

59 fx.Provide(schedulerpb.NewSchedulerServiceLayeredClient),
60 fx.Provide(
61 > func(c resource.HistoryClient) dlq.HistoryClient { fx.go
62 > return c
63 > },
64 > func(m cluster.Metadata) dlq.CurrentClusterName {
65 > return dlq.CurrentClusterName(m.GetCurrentClusterName())
66 > },
67 > func(b client.Bean) dlq.TaskClientDialer {
68 > return dlq.TaskClientDialerFn(func(_ context.Context, address string) (dlq.TaskClient, error) {
69 c, err := b.GetRemoteAdminClient(address)
70 if err != nil {
94 logger log.Logger,
95 testHooks testhooks.TestHooks,
96 > ) nsreplication.TaskExecutor { fx.go
97 > return nsreplication.NewTaskExecutor(
98 > clusterMetadata.GetCurrentClusterName(),
99 > metadataManager,
100 > dataMerger,
101 > admitter,
102 > logger,
103 > testHooks,
104 > )
105 > }),
106 fx.Provide(nsreplication.NewNoopDataMerger),
107 fx.Provide(nsreplication.NewDefaultAdmitter),
121 persistenceLazyLoadedServiceResolver service.PersistenceLazyLoadedServiceResolver,
122 logger log.SnTaggedLogger,
123 > ) service.PersistenceRateLimitingParams { fx.go
124 > return service.NewPersistenceRateLimitingParams(
125 > serviceConfig.PersistenceMaxQPS,
126 > serviceConfig.PersistenceGlobalMaxQPS,
127 > serviceConfig.PersistenceNamespaceMaxQPS,
128 > serviceConfig.PersistenceGlobalNamespaceMaxQPS,
129 > serviceConfig.PersistencePerShardNamespaceMaxQPS,
130 > serviceConfig.OperatorRPSRatio,
131 > serviceConfig.PersistenceQPSBurstRatio,
132 > serviceConfig.PersistenceDynamicRateLimitingParams,
133 > persistenceLazyLoadedServiceResolver,
134 > logger,
135 > )
136 > }
137
138 > func HostInfoProvider() (membership.HostInfo, error) { fx.go
139 > hn, err := os.Hostname()
140 > return membership.NewHostInfoFromAddress(hn), err
141 > }
142
143 func ServiceResolverProvider(
144 membershipMonitor membership.Monitor,
145 > ) (membership.ServiceResolver, error) { fx.go
146 > return membershipMonitor.GetResolver(primitives.WorkerService)
147 > }
148
149 func ConfigProvider(
150 dc *dynamicconfig.Collection,
151 persistenceConfig *config.Persistence,
152 > ) *Config { fx.go
153 > return NewConfig(
154 > dc,
155 > persistenceConfig,
156 > )
157 > }
158
159 func VisibilityManagerProvider(
169 chasmRegistry *chasm.Registry,
170 serializer serialization.Serializer,
171 > ) (manager.VisibilityManager, error) { fx.go
172 > return visibility.NewManager(
173 > *persistenceConfig,
174 > persistenceServiceResolver,
175 > customVisibilityStoreFactory,
176 > nil, // worker visibility never write
177 > saProvider,
178 > searchAttributesMapperProvider,
179 > namespaceRegistry,
180 > chasmRegistry,
181 > serviceConfig.VisibilityPersistenceMaxReadQPS,
182 > serviceConfig.VisibilityPersistenceMaxWriteQPS,
183 > serviceConfig.OperatorRPSRatio,
184 > serviceConfig.VisibilityPersistenceSlowQueryThreshold,
185 > serviceConfig.EnableReadFromSecondaryVisibility,
186 > serviceConfig.VisibilityEnableShadowReadMode,
187 > dynamicconfig.GetStringPropertyFn(visibility.SecondaryVisibilityWritingModeOff), // worker visibility never write
188 > serviceConfig.VisibilityDisableOrderByClause,
189 > serviceConfig.VisibilityEnableManualPagination,
190 > serviceConfig.VisibilityEnableUnifiedQueryConverter,
191 > metricsHandler,
192 > logger,
193 > serializer,
194 > )
195 > }
196
197 > func ServiceLifetimeHooks(lc fx.Lifecycle, svc *Service) { fx.go
198 > lc.Append(fx.StartStopHook(svc.Start, svc.Stop))
199 > }
200
201 type perNamespaceWorkerManagerInitParams struct {
223 }
224
225 > func ServerProvider(rpcFactory common.RPCFactory, logger log.Logger) *grpc.Server { fx.go
226 > opts, err := rpcFactory.GetInternodeGRPCServerOptions()
227 > if err != nil {
228 logger.Fatal("Failed to get gRPC server options", tag.Error(err))
229 }
230 > return grpc.NewServer(opts...) fx.go
231 }
go.temporal.io/server/service/fx.go 71 introduced LOC · 11 ranges

Open complete file

55
56 var PersistenceLazyLoadedServiceResolverModule = fx.Options(
57 > fx.Provide(func() PersistenceLazyLoadedServiceResolver { fx.go
58 > return PersistenceLazyLoadedServiceResolver{
59 > Value: &atomic.Value{},
60 > }
61 > }),
62 fx.Invoke(initPersistenceLazyLoadedServiceResolver),
63 )
68 serviceResolver membership.ServiceResolver,
69 lazyLoadedServiceResolver PersistenceLazyLoadedServiceResolver,
70 > ) { fx.go
71 > lazyLoadedServiceResolver.Store(serviceResolver)
72 > logger.Info("Initialized service resolver for persistence rate limiting", tag.Service(serviceName))
73 > }
74
75 func (p PersistenceLazyLoadedServiceResolver) AvailableMemberCount() int {
91 lazyLoadedServiceResolver PersistenceLazyLoadedServiceResolver,
92 logger log.Logger,
93 > ) PersistenceRateLimitingParams { fx.go
94 > hostCalculator := calculator.NewLoggedCalculator(
95 > calculator.ClusterAwareQuotaCalculator{
96 > MemberCounter: lazyLoadedServiceResolver,
97 > PerInstanceQuota: maxQps,
98 > GlobalQuota: globalMaxQps,
99 > },
100 > log.With(logger, tag.ComponentPersistence, tag.ScopeHost),
101 > )
102 > namespaceCalculator := calculator.NewLoggedNamespaceCalculator(
103 > calculator.ClusterAwareNamespaceQuotaCalculator{
104 > MemberCounter: lazyLoadedServiceResolver,
105 > PerInstanceQuota: namespaceMaxQps,
106 > GlobalQuota: globalNamespaceMaxQps,
107 > },
108 > log.With(logger, tag.ComponentPersistence, tag.ScopeNamespace),
109 > )
110 > return PersistenceRateLimitingParams{
111 > PersistenceMaxQps: func() int {
112 > return int(hostCalculator.GetQuota())
113 > },
114 PersistenceNamespaceMaxQps: func(namespace string) int {
115 return int(namespaceCalculator.GetQuota(namespace))
124 func GrpcServerOptionsProvider(
125 params GrpcServerOptionsParams,
126 > ) []grpc.ServerOption { fx.go
127 >
128 > grpcServerOptions, err := params.RPCFactory.GetInternodeGRPCServerOptions()
129 > if err != nil {
130 params.Logger.Fatal("creating gRPC server options failed", tag.Error(err))
131 }
132
133 > multiStats := rpc.MultiStatsHandler{} fx.go
134 > if params.TracingStatsHandler != nil {
135 multiStats = append(multiStats, params.TracingStatsHandler)
136 }
137 > if params.MetricsStatsHandler != nil { fx.go
138 > multiStats = append(multiStats, params.MetricsStatsHandler)
139 > }
140 > if len(multiStats) > 0 {
141 > grpcServerOptions = append(grpcServerOptions, grpc.StatsHandler(multiStats))
142 > }
143
144 > streamInterceptors := []grpc.StreamServerInterceptor{ fx.go
145 > params.TelemetryInterceptor.StreamIntercept,
146 > interceptor.CustomErrorStreamInterceptor,
147 > }
148 > if len(params.AdditionalStreamInterceptors) > 0 {
149 streamInterceptors = append(streamInterceptors, params.AdditionalStreamInterceptors...)
150 }
151
152 > return append( fx.go
153 > grpcServerOptions,
154 > grpc.ChainUnaryInterceptor(getUnaryInterceptors(params)...),
155 > grpc.ChainStreamInterceptor(streamInterceptors...),
156 > )
157 }
158
159 > func getUnaryInterceptors(params GrpcServerOptionsParams) []grpc.UnaryServerInterceptor { fx.go
160 > interceptors := []grpc.UnaryServerInterceptor{
161 > params.ServiceErrorInterceptor.Intercept,
162 > metrics.NewServerMetricsContextInjectorInterceptor(),
163 > metrics.NewServerMetricsTrailerPropagatorInterceptor(params.Logger),
164 > params.TelemetryInterceptor.UnaryIntercept,
165 > }
166 >
167 > interceptors = append(interceptors, params.AdditionalInterceptors...)
168 >
169 > if params.NamespaceRateLimitInterceptor != nil {
170 > interceptors = append(interceptors, params.NamespaceRateLimitInterceptor.Intercept)
171 > }
172
173 > interceptors = append(interceptors, params.RateLimitInterceptor.Intercept) fx.go
174 >
175 > if params.ContextMetadataInterceptor != nil {
176 > interceptors = append(interceptors, params.ContextMetadataInterceptor.Intercept)
177 > }
178
179 > return append(interceptors, params.RetryableInterceptor.Intercept) fx.go
180 }
go.temporal.io/server/client/client_bean.go 66 introduced LOC · 13 ranges

Open complete file

56
57 // NewClientBean provides a collection of clients
58 > func NewClientBean(factory Factory, clusterMetadata cluster.Metadata) (Bean, error) { client_bean.go
59 >
60 > historyClient, err := factory.NewHistoryClientWithTimeout(history.DefaultTimeout)
61 > if err != nil {
62 return nil, err
63 }
64
65 > adminClients := map[string]adminservice.AdminServiceClient{} client_bean.go
66 > frontendClients := map[string]frontendClient{}
67 >
68 > currentClusterName := clusterMetadata.GetCurrentClusterName()
69 > // Init local cluster client with membership info
70 > adminClient, err := factory.NewLocalAdminClientWithTimeout(
71 > admin.DefaultTimeout,
72 > admin.DefaultLargeTimeout,
73 > )
74 > if err != nil {
75 return nil, err
76 }
77 > conn, client, err := factory.NewLocalFrontendClientWithTimeout( client_bean.go
78 > frontend.DefaultTimeout,
79 > frontend.DefaultLongPollTimeout,
80 > )
81 > if err != nil {
82 return nil, err
83 }
84 > adminClients[currentClusterName] = adminClient client_bean.go
85 > frontendClients[currentClusterName] = frontendClient{
86 > connection: conn,
87 > WorkflowServiceClient: client,
88 > }
89 >
90 > bean := &clientBeanImpl{
91 > factory: factory,
92 > historyClient: historyClient,
93 > clusterMetadata: clusterMetadata,
94 > adminClients: adminClients,
95 > frontendClients: frontendClients,
96 > }
97 > bean.registerClientEviction()
98 > return bean, nil
99 }
100
101 > func (h *clientBeanImpl) registerClientEviction() { client_bean.go
102 > currentCluster := h.clusterMetadata.GetCurrentClusterName()
103 > h.clusterMetadata.RegisterMetadataChangeCallback(
104 > h,
105 > func(oldClusterMetadata map[string]*cluster.ClusterInformation, newClusterMetadata map[string]*cluster.ClusterInformation) {
106 > for clusterName := range newClusterMetadata {
107 > if clusterName == currentCluster {
108 > continue
109 }
110 h.adminClientsLock.Lock()
136 }
137
138 > func (h *clientBeanImpl) GetHistoryClient() historyservice.HistoryServiceClient { client_bean.go
139 > return h.historyClient
140 > }
141
142 > func (h *clientBeanImpl) GetMatchingClient(namespaceIDToName NamespaceIDToNameFunc) (matchingservice.MatchingServiceClient, error) { client_bean.go
143 > if client := h.matchingClient.Load(); client != nil {
144 return client.(matchingservice.MatchingServiceClient), nil
145 }
146 > return h.lazyInitMatchingClient(namespaceIDToName) client_bean.go
147 }
148
149 > func (h *clientBeanImpl) GetFrontendClient() workflowservice.WorkflowServiceClient { client_bean.go
150 > return h.frontendClients[h.clusterMetadata.GetCurrentClusterName()]
151 > }
152
153 > func (h *clientBeanImpl) GetRemoteAdminClient(cluster string) (adminservice.AdminServiceClient, error) { client_bean.go
154 > h.adminClientsLock.RLock()
155 > client, ok := h.adminClients[cluster]
156 > h.adminClientsLock.RUnlock()
157 > if ok {
158 > return client, nil
159 > }
160
161 clusterInfo, clusterFound := h.clusterMetadata.GetAllClusterInfo()[cluster]
234 }
235
236 > func (h *clientBeanImpl) lazyInitMatchingClient(namespaceIDToName NamespaceIDToNameFunc) (matchingservice.MatchingServiceClient, error) { client_bean.go
237 > h.Lock()
238 > defer h.Unlock()
239 > if cached := h.matchingClient.Load(); cached != nil {
240 return cached.(matchingservice.MatchingServiceClient), nil
241 }
242 > client, err := h.factory.NewMatchingClientWithTimeout(namespaceIDToName, matching.DefaultTimeout, matching.DefaultLongPollTimeout) client_bean.go
243 > if err != nil {
244 return nil, err
245 }
246 > h.matchingClient.Store(client) client_bean.go
247 > return client, nil
248 }
go.temporal.io/server/common/persistence/persistence_metric_clients.go 62 introduced LOC · 9 ranges

Open complete file

189 ctx context.Context,
190 request *UpdateShardRequest,
191 > ) (retErr error) { persistence_metric_clients.go
192 > caller := headers.GetCallerInfo(ctx).CallerName
193 > startTime := time.Now().UTC()
194 > defer func() {
195 > p.healthSignals.Record(request.ShardInfo.GetShardId(), time.Since(startTime), retErr)
196 > p.recordRequestMetrics(metrics.PersistenceUpdateShardScope, caller, time.Since(startTime), retErr)
197 > p.recordDataLossMetrics(metrics.PersistenceUpdateShardScope, caller, retErr, "", "")
198 > }()
199 > return p.persistence.UpdateShard(ctx, request)
200 }
201
218 }
219
220 > func (p *executionPersistenceClient) GetName() string { persistence_metric_clients.go
221 > return p.persistence.GetName()
222 > }
223
224 func (p *executionPersistenceClient) GetHistoryBranchUtil() HistoryBranchUtil {
419 ctx context.Context,
420 request *GetHistoryTasksRequest,
421 > ) (_ *GetHistoryTasksResponse, retErr error) { persistence_metric_clients.go
422 > var operation string
423 > switch request.TaskCategory.ID() {
424 > case tasks.CategoryIDTransfer:
425 > operation = metrics.PersistenceGetTransferTasksScope
426 > case tasks.CategoryIDTimer:
427 > operation = metrics.PersistenceGetTimerTasksScope
428 > case tasks.CategoryIDVisibility:
429 > operation = metrics.PersistenceGetVisibilityTasksScope
430 case tasks.CategoryIDReplication:
431 operation = metrics.PersistenceGetReplicationTasksScope
432 case tasks.CategoryIDArchival:
433 operation = metrics.PersistenceGetArchivalTasksScope
434 > case tasks.CategoryIDOutbound: persistence_metric_clients.go
435 > operation = metrics.PersistenceGetOutboundTasksScope
436 default:
437 return nil, serviceerror.NewInternalf("unknown task category type: %v", request.TaskCategory)
438 }
439
440 > caller := headers.GetCallerInfo(ctx).CallerName persistence_metric_clients.go
441 > startTime := time.Now().UTC()
442 > defer func() {
443 > p.healthSignals.Record(request.ShardID, time.Since(startTime), retErr)
444 > p.recordRequestMetrics(operation, caller, time.Since(startTime), retErr)
445 > p.recordDataLossMetrics(operation, caller, retErr, "", "")
446 > }()
447 > return p.persistence.GetHistoryTasks(ctx, request)
448 }
449
890 }
891
892 > func (p *metadataPersistenceClient) WatchNamespaces(ctx context.Context) (_ <-chan *NamespaceWatchEvent, retErr error) { persistence_metric_clients.go
893 > caller := headers.GetCallerInfo(ctx).CallerName
894 > startTime := time.Now().UTC()
895 > defer func() {
896 > metricErr := retErr
897 > // WatchNotSupported isn't really a persistence error. It's just a signal that persistence doesn't support watching.
898 > if errors.Is(metricErr, ErrWatchNotSupported) {
899 > metricErr = nil
900 > }
901 > p.healthSignals.Record(CallerSegmentMissing, time.Since(startTime), metricErr)
902 > p.recordRequestMetrics(metrics.PersistenceWatchNamespacesScope, caller, time.Since(startTime), metricErr)
903 > p.recordDataLossMetrics(metrics.PersistenceWatchNamespacesScope, caller, metricErr, "", "")
904 }()
905 > return p.persistence.WatchNamespaces(ctx) persistence_metric_clients.go
906 }
907
1240 func (p *clusterMetadataPersistenceClient) GetCurrentClusterMetadata(
1241 ctx context.Context,
1242 > ) (_ *GetClusterMetadataResponse, retErr error) { persistence_metric_clients.go
1243 > caller := headers.GetCallerInfo(ctx).CallerName
1244 > startTime := time.Now().UTC()
1245 > defer func() {
1246 > p.healthSignals.Record(CallerSegmentMissing, time.Since(startTime), retErr)
1247 > p.recordRequestMetrics(metrics.PersistenceGetCurrentClusterMetadataScope, caller, time.Since(startTime), retErr)
1248 > p.recordDataLossMetrics(metrics.PersistenceGetCurrentClusterMetadataScope, caller, retErr, "", "")
1249 > }()
1250 > return p.persistence.GetCurrentClusterMetadata(ctx)
1251 }
1252
1378 ctx context.Context,
1379 request *ListNexusEndpointsRequest,
1380 > ) (_ *ListNexusEndpointsResponse, retErr error) { persistence_metric_clients.go
1381 > caller := headers.GetCallerInfo(ctx).CallerName
1382 > startTime := time.Now().UTC()
1383 > defer func() {
1384 > p.healthSignals.Record(CallerSegmentMissing, time.Since(startTime), retErr)
1385 > p.recordRequestMetrics(metrics.PersistenceListNexusEndpointsScope, caller, time.Since(startTime), retErr)
1386 > p.recordDataLossMetrics(metrics.PersistenceListNexusEndpointsScope, caller, retErr, "", "")
1387 > }()
1388 > return p.persistence.ListNexusEndpoints(ctx, request)
1389 }
1390
go.temporal.io/server/common/metrics/runtime.go 59 introduced LOC · 8 ranges

Open complete file

37 logger log.Logger,
38 instanceID string,
39 > ) *RuntimeMetricsReporter { runtime.go
40 > if len(instanceID) > 0 {
41 handler = handler.WithTags(StringTag(instance, instanceID))
42 }
43 > var memstats runtime.MemStats runtime.go
44 > runtime.ReadMemStats(&memstats)
45 >
46 > return &RuntimeMetricsReporter{
47 > handler: handler,
48 > reportInterval: reportInterval,
49 > logger: logger,
50 > lastNumGC: memstats.NumGC,
51 > quit: make(chan struct{}),
52 > buildTime: build.InfoData.GitTime,
53 > buildInfoHandler: handler.WithTags(
54 > StringTag(gitRevisionTag, build.InfoData.GitRevision),
55 > StringTag(buildDateTag, build.InfoData.GitTime.Format(time.RFC3339)),
56 > StringTag(buildPlatformTag, build.InfoData.GoArch),
57 > StringTag(goVersionTag, build.InfoData.GoVersion),
58 > StringTag(buildVersionTag, headers.ServerVersion),
59 > ),
60 > }
61 }
62
63 // report Sends runtime metrics to the local metrics collector.
64 > func (r *RuntimeMetricsReporter) report() { runtime.go
65 > var memStats runtime.MemStats
66 > runtime.ReadMemStats(&memStats)
67 >
68 > NumGoRoutinesGauge.With(r.handler).Record(float64(runtime.NumGoroutine()))
69 > GoMaxProcsGauge.With(r.handler).Record(float64(runtime.GOMAXPROCS(0)))
70 > MemoryAllocatedGauge.With(r.handler).Record(float64(memStats.Alloc))
71 > MemoryHeapGauge.With(r.handler).Record(float64(memStats.HeapAlloc))
72 > MemoryHeapObjectsGauge.With(r.handler).Record(float64(memStats.HeapObjects))
73 > MemoryHeapIdleGauge.With(r.handler).Record(float64(memStats.HeapIdle))
74 > MemoryHeapInuseGauge.With(r.handler).Record(float64(memStats.HeapInuse))
75 > MemoryHeapReleasedGauge.With(r.handler).Record(float64(memStats.HeapReleased))
76 > MemoryStackGauge.With(r.handler).Record(float64(memStats.StackInuse))
77 > MemoryMallocsGauge.With(r.handler).Record(float64(memStats.Mallocs))
78 > MemoryFreesGauge.With(r.handler).Record(float64(memStats.Frees))
79 >
80 > NumGCGauge.With(r.handler).Record(float64(memStats.NumGC))
81 > GcPauseNsTotal.With(r.handler).Record(float64(memStats.PauseTotalNs))
82 >
83 > // memStats.NumGC is a perpetually incrementing counter (unless it wraps at 2^32)
84 > num := memStats.NumGC
85 > lastNum := atomic.SwapUint32(&r.lastNumGC, num) // reset for the next iteration
86 > if delta := num - lastNum; delta > 0 {
87 > NumGCCounter.With(r.handler).Record(int64(delta))
88 > if delta > 255 {
89 // too many GCs happened, the timestamps buffer got wrapped around. Report only the last 256
90 lastNum = num - 256
91 }
92 > for i := lastNum; i != num; i++ { runtime.go
93 > pause := memStats.PauseNs[i%256]
94 > GcPauseMsTimer.With(r.handler).Record(time.Duration(pause))
95 > }
96 }
97
98 // report build info
99 > r.buildInfoHandler.Gauge(buildInfoMetricName).Record(1.0) runtime.go
100 > r.buildInfoHandler.Gauge(buildAgeMetricName).Record(float64(time.Since(r.buildTime)))
101 }
102
103 // Start Starts the reporter thread that periodically emits metrics.
104 > func (r *RuntimeMetricsReporter) Start() { runtime.go
105 > if !atomic.CompareAndSwapInt32(&r.started, 0, 1) {
106 return
107 }
108 > r.report() runtime.go
109 > go func() {
110 > ticker := time.NewTicker(r.reportInterval)
111 > for {
112 > select {
113 case <-ticker.C:
114 r.report()
119 }
120 }()
121 > r.logger.Info("RuntimeMetricsReporter started") runtime.go
122 }
123