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

1057 LOC · 552 covered · 505 uncovered · 59 ranges · 245 concepts · 11 introducers · 138 tests

File neighbourhood

The centred file is linked to every concept that introduces one of its ranges, every test that runs code from the file, and the gray connector concepts standing between those tests and the file's own introducer concepts. Undirected links join concepts to every file where they introduce source and concepts to the tests they introduce; arrows show specialization between the displayed concepts and bridge only concepts omitted from this view. Concept colors match the source ranges below; connector concepts have no source color and are shown in gray.

Focused file, its introducer and connector concepts, their introduced files, and tests that run code from the file

In the embedded map, ordinary wheel input scrolls the page; use the visible controls to zoom and drag to pan. Open the full-screen map for canvas navigation: wheel pans, Ctrl/Command plus wheel zooms, and arrow keys pan when this region is focused. On touch screens, open the full-screen map to pan or pinch. If JavaScript or WebGL is unavailable, use the related-file, concept, and source links on this page.

Graph controls are ready.

Interactive rendering requires JavaScript and WebGL. Use the related-file, concept, and source links on this page while the interactive map is unavailable.

1 package frontend
2
3 import (
4 "fmt"
5 "net"
6
7 "github.com/gorilla/mux"
8 "go.temporal.io/server/api/adminservice/v1"
9 "go.temporal.io/server/chasm"
10 "go.temporal.io/server/chasm/lib/activity"
11 "go.temporal.io/server/chasm/lib/callback"
12 chasmnexus "go.temporal.io/server/chasm/lib/nexusoperation"
13 nexusoperationpb "go.temporal.io/server/chasm/lib/nexusoperation/gen/nexusoperationpb/v1"
14 chasmscheduler "go.temporal.io/server/chasm/lib/scheduler"
15 "go.temporal.io/server/chasm/lib/scheduler/gen/schedulerpb/v1"
16 chasmtests "go.temporal.io/server/chasm/lib/tests"
17 chasmworkflow "go.temporal.io/server/chasm/lib/workflow"
18 "go.temporal.io/server/client"
19 "go.temporal.io/server/common"
20 "go.temporal.io/server/common/archiver"
21 "go.temporal.io/server/common/archiver/provider"
22 "go.temporal.io/server/common/authorization"
23 "go.temporal.io/server/common/clock"
24 "go.temporal.io/server/common/cluster"
25 "go.temporal.io/server/common/config"
26 "go.temporal.io/server/common/dynamicconfig"
27 "go.temporal.io/server/common/log"
28 "go.temporal.io/server/common/log/tag"
29 "go.temporal.io/server/common/membership"
30 "go.temporal.io/server/common/metrics"
31 "go.temporal.io/server/common/namespace"
32 "go.temporal.io/server/common/namespace/nsreplication"
33 "go.temporal.io/server/common/persistence"
34 "go.temporal.io/server/common/persistence/serialization"
35 "go.temporal.io/server/common/persistence/visibility"
36 "go.temporal.io/server/common/persistence/visibility/manager"
37 "go.temporal.io/server/common/primitives"
38 "go.temporal.io/server/common/quotas"
39 "go.temporal.io/server/common/quotas/calculator"
40 "go.temporal.io/server/common/resolver"
41 "go.temporal.io/server/common/resource"
42 "go.temporal.io/server/common/rpc"
43 "go.temporal.io/server/common/rpc/encryption"
44 "go.temporal.io/server/common/rpc/interceptor"
45 "go.temporal.io/server/common/sdk"
46 "go.temporal.io/server/common/searchattribute"
47 "go.temporal.io/server/common/telemetry"
48 "go.temporal.io/server/common/testing/testhooks"
49 "go.temporal.io/server/service"
50 "go.temporal.io/server/service/frontend/configs"
51 "go.temporal.io/server/service/history/tasks"
52 "go.temporal.io/server/service/worker/scheduler"
53 "go.temporal.io/server/service/worker/workerdeployment"
54 "go.uber.org/fx"
55 "google.golang.org/grpc"
56 "google.golang.org/grpc/health"
57 healthpb "google.golang.org/grpc/health/grpc_health_v1"
58 "google.golang.org/grpc/keepalive"
59 )
60
61 type (
62 FEReplicatorNamespaceReplicationQueue persistence.NamespaceReplicationQueue
63
64 namespaceChecker struct {
65 r namespace.Registry
66 }
67 )
68
69 var Module = fx.Options(
70 resource.Module,
71 chasmtests.Module,
72 scheduler.Module,
73 workerdeployment.Module,
74 // Note that with this approach routes may be registered in arbitrary order.
75 // This is okay because our routes don't have overlapping matches.
76 // The only important detail is that the PathPrefix("/") route registered in the HTTPAPIServerProvider comes last.
77 // Coincidentally, this is the case today, likely because it has more dependencies that the other dependencies.
78 // This approach isn't perfect but at it allows the router to be pluggable and we have enough functional test
79 // coverage to catch misconfiguration.
80 // A more robust approach would require using fx groups but we shouldn't overcomplicate until this becomes an issue.
81 fx.Provide(MuxRouterProvider),
82 fx.Provide(ConfigProvider),
83 fx.Provide(ServiceErrorInterceptorProvider),
84 fx.Provide(NamespaceLogInterceptorProvider),
85 fx.Provide(NamespaceHandoverInterceptorProvider),
86 fx.Provide(interceptor.NewRoutingKeyExtractor),
87 fx.Provide(BusinessIDInterceptorProvider),
88 fx.Provide(RedirectionInterceptorProvider),
89 fx.Provide(ErrorHandlerProvider),
90 fx.Provide(TelemetryInterceptorProvider),
91 fx.Provide(RetryableInterceptorProvider),
92 fx.Provide(RateLimitInterceptorProvider),
93 fx.Provide(interceptor.NewHealthInterceptor),
94 fx.Provide(NamespaceCountLimitInterceptorProvider),
95 fx.Provide(NamespaceValidatorInterceptorProvider),
96 fx.Provide(NamespaceRateLimitInterceptorProvider),
97 fx.Provide(SDKVersionInterceptorProvider),
98 fx.Provide(CallerInfoInterceptorProvider),
99 fx.Provide(SlowRequestLoggerInterceptorProvider),
100 fx.Provide(MaskInternalErrorDetailsInterceptorProvider),
101 fx.Provide(ContextMetadataInterceptorProvider),
102 fx.Provide(GrpcServerOptionsProvider),
103 fx.Provide(VisibilityManagerProvider),
104 fx.Provide(ThrottledLoggerRpsFnProvider),
105 fx.Provide(PersistenceRateLimitingParamsProvider),
106 service.PersistenceLazyLoadedServiceResolverModule,
107 fx.Provide(FEReplicatorNamespaceReplicationQueueProvider),
108 fx.Provide(nsreplication.NewNoopDataMerger),
109 fx.Provide(nsreplication.NewDefaultAdmitter),
110 fx.Provide(AuthorizationInterceptorProvider),
111 fx.Provide(NamespaceCheckerProvider),
112 > fx.Provide(func(so GrpcServerOptions) *grpc.Server { return grpc.NewServer(so.Options...) }), fx.go ×44
113 fx.Provide(callbackValidatorProvider),
114 fx.Provide(HandlerProvider),
115 fx.Provide(AdminHandlerProvider),
116 fx.Provide(NamespaceDLQHandlerProvider),
117 fx.Provide(OperatorHandlerProvider),
118 fx.Provide(NewVersionChecker),
119 fx.Provide(ServiceResolverProvider),
120 fx.Provide(newNexusCompletionHandler),
121 fx.Provide(NewNexusOperationHTTPHandler),
122 fx.Provide(newNexusCompletionHTTPHandler),
123 fx.Invoke(RegisterNexusOperationHTTPHandler),
124 fx.Invoke(RegisterNexusCompletionHTTPHandler),
125 fx.Invoke(RegisterOpenAPIHTTPHandler),
126 fx.Provide(HTTPAPIServerProvider),
127 fx.Provide(NewServiceProvider),
128 fx.Provide(NexusEndpointClientProvider),
129 fx.Invoke(ServiceLifetimeHooks),
130 fx.Provide(nexusoperationpb.NewNexusOperationServiceLayeredClient),
131 fx.Provide(schedulerpb.NewSchedulerServiceLayeredClient),
132 fx.Provide(chasmnexus.NewFrontendHandler),
133 chasmnexus.Module,
134 chasmscheduler.Module,
135 chasmworkflow.Module,
136 callback.Module,
137 activity.FrontendModule,
138 fx.Provide(visibility.ChasmVisibilityManagerProvider),
139 fx.Provide(chasm.ChasmVisibilityInterceptorProvider),
140 )
141
142 func NewServiceProvider(
143 serviceConfig *Config,
144 server *grpc.Server,
145 healthServer *health.Server,
146 httpAPIServer *HTTPAPIServer,
147 handler Handler,
148 adminHandler *AdminHandler,
149 operatorHandler *OperatorHandlerImpl,
150 versionChecker *VersionChecker,
151 visibilityMgr manager.VisibilityManager,
152 logger log.SnTaggedLogger,
153 grpcListener net.Listener,
154 metricsHandler metrics.Handler,
155 membershipMonitor membership.Monitor,
156 > ) *Service { fx.go ×44
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
175 // with the interceptors that are already set in the options.
176 type GrpcServerOptions struct {
177 Options []grpc.ServerOption
178 UnaryInterceptors []grpc.UnaryServerInterceptor
179 }
180
181 func AuthorizationInterceptorProvider(
182 cfg *config.Config,
183 serviceConfig *Config,
184 logger log.Logger,
185 namespaceChecker authorization.NamespaceChecker,
186 metricsHandler metrics.Handler,
187 authorizer authorization.Authorizer,
188 claimMapper authorization.ClaimMapper,
189 audienceGetter authorization.JWTAudienceMapper,
190 dc *dynamicconfig.Collection,
191 > ) *authorization.Interceptor { fx.go ×44
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 ×44
209 > return &namespaceChecker{r: registry}
210 > }
211
212 > func (n *namespaceChecker) Exists(name namespace.Name) error { rpc.go ×1
213 > // This will get called before the namespace state validation interceptor. We want to
214 > // disable readthrough to avoid polluting the negative lookup cache, e.g. if this call is
215 > // for RegisterNamespace and the namespace doesn't exist yet.
216 > opts := namespace.GetNamespaceOptions{DisableReadthrough: true}
217 > _, err := n.r.GetNamespaceWithOptions(name, opts)
218 > return err
219 > }
220
221 func GrpcServerOptionsProvider(
222 logger log.Logger,
223 cfg *config.Config,
224 serviceConfig *Config,
225 serviceName primitives.ServiceName,
226 rpcFactory common.RPCFactory,
227 serviceErrorInterceptor *interceptor.ServiceErrorInterceptor,
228 namespaceLogInterceptor *interceptor.NamespaceLogInterceptor,
229 namespaceRateLimiterInterceptor interceptor.NamespaceRateLimitInterceptor,
230 namespaceCountLimiterInterceptor *interceptor.ConcurrentRequestLimitInterceptor,
231 namespaceValidatorInterceptor *interceptor.NamespaceValidatorInterceptor,
232 namespaceHandoverInterceptor *interceptor.NamespaceHandoverInterceptor,
233 businessIDInterceptor *interceptor.RoutingKeyInterceptor,
234 redirectionInterceptor *interceptor.Redirection,
235 telemetryInterceptor *interceptor.TelemetryInterceptor,
236 retryableInterceptor *interceptor.RetryableInterceptor,
237 healthInterceptor *interceptor.HealthInterceptor,
238 rateLimitInterceptor *interceptor.RateLimitInterceptor,
239 traceStatsHandler telemetry.ServerStatsHandler,
240 metricsStatsHandler metrics.ServerStatsHandler,
241 sdkVersionInterceptor *interceptor.SDKVersionInterceptor,
242 callerInfoInterceptor *interceptor.CallerInfoInterceptor,
243 authInterceptor *authorization.Interceptor,
244 maskInternalErrorDetailsInterceptor *interceptor.MaskInternalErrorDetailsInterceptor,
245 contextMetadataInterceptor *interceptor.ContextMetadataInterceptor,
246 slowRequestLoggerInterceptor *interceptor.SlowRequestLoggerInterceptor,
247 chasmRequestVisibilityInterceptor *chasm.ChasmVisibilityInterceptor,
248 customInterceptors []grpc.UnaryServerInterceptor,
249 customStreamInterceptors []grpc.StreamServerInterceptor,
250 metricsHandler metrics.Handler,
251 > ) GrpcServerOptions { fx.go ×44
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()
270 default:
271 err = fmt.Errorf("unexpected frontend service name %q", serviceName)
272 }
273 > if err != nil { fx.go ×44
274 logger.Fatal("creating gRPC server options failed", tag.Error(err))
275 }
276 > unaryInterceptors := []grpc.UnaryServerInterceptor{ fx.go ×44
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 http_api_server.go ×23
308 > unaryInterceptors = append(unaryInterceptors, customInterceptors...)
309 > }
310 // retry interceptor should be the most inner interceptor
311 > unaryInterceptors = append(unaryInterceptors, retryableInterceptor.Intercept) fx.go ×44
312 >
313 > streamInterceptor := []grpc.StreamServerInterceptor{
314 > authInterceptor.InterceptStream,
315 > telemetryInterceptor.StreamIntercept,
316 > }
317 > if len(customStreamInterceptors) > 0 {
318 > streamInterceptor = append(streamInterceptor, customStreamInterceptors...) onebox.go ×75
319 > }
320
321 > grpcServerOptions = append( fx.go ×44
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) data_store_factory.go ×29
332 > }
333 > if metricsStatsHandler != nil { fx.go ×44
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
342 func ConfigProvider(
343 dc *dynamicconfig.Collection,
344 persistenceConfig config.Persistence,
345 > ) *Config { fx.go ×44
346 > return NewConfig(
347 > dc,
348 > persistenceConfig.NumHistoryShards,
349 > )
350 > }
351
352 func ServiceErrorInterceptorProvider(
353 dc *dynamicconfig.Collection,
354 > ) *interceptor.ServiceErrorInterceptor { fx.go ×44
355 > return interceptor.NewServiceErrorInterceptor(
356 > dynamicconfig.MaxServiceErrorMessageLength.Get(dc),
357 > )
358 > }
359
360 > func ThrottledLoggerRpsFnProvider(serviceConfig *Config) resource.ThrottledLoggerRpsFn { fx.go ×44
361 > return func() float64 { return float64(serviceConfig.ThrottledLogRPS()) }
362 }
363
364 func NamespaceLogInterceptorProvider(
365 namespaceLogger resource.NamespaceLogger,
366 namespaceRegistry namespace.Registry,
367 > ) *interceptor.NamespaceLogInterceptor { fx.go ×44
368 > return interceptor.NewNamespaceLogInterceptor(
369 > namespaceRegistry,
370 > namespaceLogger)
371 > }
372
373 > func RetryableInterceptorProvider() *interceptor.RetryableInterceptor { fx.go ×44
374 > return interceptor.NewRetryableInterceptor(
375 > common.CreateFrontendHandlerRetryPolicy(),
376 > common.IsServiceHandlerRetryableError,
377 > )
378 > }
379
380 func RedirectionInterceptorProvider(
381 configuration *Config,
382 namespaceCache namespace.Registry,
383 policy config.DCRedirectionPolicy,
384 logger log.Logger,
385 clientBean client.Bean,
386 metricsHandler metrics.Handler,
387 timeSource clock.TimeSource,
388 clusterMetadata cluster.Metadata,
389 > ) *interceptor.Redirection { fx.go ×44
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 ×44
407 > return interceptor.NewRoutingKeyInterceptor(
408 > []interceptor.RoutingKeyExtractorFunc{
409 > interceptor.WorkflowServiceExtractor(extractor),
410 > },
411 > logger,
412 > )
413 > }
414
415 type NamespaceHandoverInterceptorParams struct {
416 fx.In
417 DynamicConfig *dynamicconfig.Collection
418 NamespaceRegistry namespace.Registry
419 Logger log.Logger
420 MetricsHandler metrics.Handler
421 TimeSource clock.TimeSource
422 RequestErrorHandler *interceptor.RequestErrorHandler
423 AdditionalAllowedMethodsDuringHandover []string `group:"additionalAllowedMethodsDuringHandover"`
424 }
425
426 func NamespaceHandoverInterceptorProvider(
427 params NamespaceHandoverInterceptorParams,
428 > ) *interceptor.NamespaceHandoverInterceptor { fx.go ×44
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 ×44
444 > return interceptor.NewRequestErrorHandler(
445 > logger,
446 > serviceConfig.LogAllReqErrors,
447 > )
448 > }
449
450 func TelemetryInterceptorProvider(
451 logger log.Logger,
452 metricsHandler metrics.Handler,
453 namespaceRegistry namespace.Registry,
454 serviceConfig *Config,
455 requestErrorHandler *interceptor.RequestErrorHandler,
456 > ) *interceptor.TelemetryInterceptor { fx.go ×44
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 { fx.go ×3
467 > return func() float64 {
468 > rate := rateFn()
469 > metrics.HostRPSLimit.With(handler).Record(rate)
470 > return rate
471 > }
472 }
473
474 func RateLimitInterceptorProvider(
475 serviceConfig *Config,
476 frontendServiceResolver membership.ServiceResolver,
477 handler metrics.Handler,
478 logger log.SnTaggedLogger,
479 > ) *interceptor.RateLimitInterceptor { fx.go ×3
480 > rateFn := calculator.NewLoggedCalculator(
481 > calculator.ClusterAwareQuotaCalculator{
482 > MemberCounter: frontendServiceResolver,
483 > PerInstanceQuota: serviceConfig.RPS,
484 > GlobalQuota: serviceConfig.GlobalRPS,
485 > },
486 > log.With(logger, tag.ComponentRPCHandler, tag.ScopeHost),
487 > ).GetQuota
488 > rateFnWithMetrics := getRateFnWithMetrics(rateFn, handler)
489 >
490 > namespaceReplicationInducingRateFn := func() float64 {
491 > return float64(serviceConfig.NamespaceReplicationInducingAPIsRPS())
492 > }
493
494 > return interceptor.NewRateLimitInterceptor( fx.go ×3
495 > configs.NewRequestToRateLimiter(
496 > quotas.NewDefaultIncomingRateBurst(rateFnWithMetrics),
497 > quotas.NewDefaultIncomingRateBurst(rateFn),
498 > quotas.NewDefaultIncomingRateBurst(namespaceReplicationInducingRateFn),
499 > serviceConfig.OperatorRPSRatio,
500 > ),
501 > map[string]int{
502 > healthpb.Health_Check_FullMethodName: 0, // exclude health check requests from rate limiting.
503 > adminservice.AdminService_DeepHealthCheck_FullMethodName: 0, // exclude deep health check requests from rate limiting.
504 > },
505 > )
506 }
507
508 func ContextMetadataInterceptorProvider(
509 logger log.Logger,
510 dc *dynamicconfig.Collection,
511 > ) *interceptor.ContextMetadataInterceptor { fx.go ×44
512 > setTrailer := dynamicconfig.FrontendContextMetadataSetTrailer.Get(dc)()
513 > return interceptor.NewContextMetadataInterceptor(setTrailer, logger)
514 > }
515
516 func MaskInternalErrorDetailsInterceptorProvider(
517 logger log.Logger,
518 serviceConfig *Config,
519 namespaceRegistry namespace.Registry,
520 > ) *interceptor.MaskInternalErrorDetailsInterceptor { fx.go ×44
521 > return interceptor.NewMaskInternalErrorDetailsInterceptor(
522 > serviceConfig.MaskInternalErrorDetails, namespaceRegistry, logger,
523 > )
524 > }
525
526 func NamespaceRateLimitInterceptorProvider(
527 serviceName primitives.ServiceName,
528 serviceConfig *Config,
529 namespaceRegistry namespace.Registry,
530 frontendServiceResolver membership.ServiceResolver,
531 metricsHandler metrics.Handler,
532 logger log.SnTaggedLogger,
533 > ) interceptor.NamespaceRateLimitInterceptor { fx.go ×3
534 > var globalNamespaceRPS, globalNamespaceVisibilityRPS, globalNamespaceNamespaceReplicationInducingAPIsRPS dynamicconfig.IntPropertyFnWithNamespaceFilter
535 >
536 > switch serviceName {
537 > case primitives.FrontendService:
538 > globalNamespaceRPS = serviceConfig.GlobalNamespaceRPS
539 > globalNamespaceVisibilityRPS = serviceConfig.GlobalNamespaceVisibilityRPS
540 > globalNamespaceNamespaceReplicationInducingAPIsRPS = serviceConfig.GlobalNamespaceNamespaceReplicationInducingAPIsRPS
541 case primitives.InternalFrontendService:
542 globalNamespaceRPS = serviceConfig.InternalFEGlobalNamespaceRPS
543 globalNamespaceVisibilityRPS = serviceConfig.InternalFEGlobalNamespaceVisibilityRPS
544 // Internal frontend has no special limit for this set of APIs
545 globalNamespaceNamespaceReplicationInducingAPIsRPS = serviceConfig.InternalFEGlobalNamespaceRPS
546 default:
547 panic("invalid service name")
548 }
549
550 > namespaceRateFn := calculator.NewLoggedNamespaceCalculator( fx.go ×3
551 > calculator.ClusterAwareNamespaceQuotaCalculator{
552 > MemberCounter: frontendServiceResolver,
553 > PerInstanceQuota: serviceConfig.MaxNamespaceRPSPerInstance,
554 > GlobalQuota: globalNamespaceRPS,
555 > },
556 > log.With(logger, tag.ComponentRPCHandler, tag.ScopeNamespace),
557 > ).GetQuota
558 > visibilityRateFn := calculator.NewLoggedNamespaceCalculator(
559 > calculator.ClusterAwareNamespaceQuotaCalculator{
560 > MemberCounter: frontendServiceResolver,
561 > PerInstanceQuota: serviceConfig.MaxNamespaceVisibilityRPSPerInstance,
562 > GlobalQuota: globalNamespaceVisibilityRPS,
563 > },
564 > log.With(logger, tag.ComponentVisibilityHandler, tag.ScopeNamespace),
565 > ).GetQuota
566 > namespaceReplicationInducingRateFn := calculator.NewLoggedNamespaceCalculator(
567 > calculator.ClusterAwareNamespaceQuotaCalculator{
568 > MemberCounter: frontendServiceResolver,
569 > PerInstanceQuota: serviceConfig.MaxNamespaceNamespaceReplicationInducingAPIsRPSPerInstance,
570 > GlobalQuota: globalNamespaceNamespaceReplicationInducingAPIsRPS,
571 > },
572 > log.With(logger, tag.ComponentNamespaceReplication, tag.ScopeNamespace),
573 > ).GetQuota
574 > namespaceRateLimiter := quotas.NewNamespaceRequestRateLimiter(
575 > func(req quotas.Request) quotas.RequestRateLimiter {
576 > return configs.NewRequestToRateLimiter( rate_burst.go ×3
577 > quotas.NewNamespaceRateBurst(
578 > req.Caller,
579 > namespaceRateFn,
580 > quotas.NamespaceBurstRatioFn(serviceConfig.MaxNamespaceBurstRatioPerInstance),
581 > ),
582 > quotas.NewNamespaceRateBurst(
583 > req.Caller,
584 > visibilityRateFn,
585 > quotas.NamespaceBurstRatioFn(serviceConfig.MaxNamespaceVisibilityBurstRatioPerInstance),
586 > ),
587 > quotas.NewNamespaceRateBurst(
588 > req.Caller,
589 > namespaceReplicationInducingRateFn,
590 > quotas.NamespaceBurstRatioFn(serviceConfig.MaxNamespaceNamespaceReplicationInducingAPIsBurstRatioPerInstance),
591 > ),
592 > serviceConfig.OperatorRPSRatio,
593 > )
594 > },
595 )
596 > return interceptor.NewNamespaceRateLimitInterceptor( fx.go ×3
597 > namespaceRegistry,
598 > namespaceRateLimiter,
599 > map[string]int{}, // no token overrides
600 > configs.PollTaskAPISet,
601 > serviceConfig.PollWaitForNamespaceRateLimitToken,
602 > metricsHandler,
603 > )
604 }
605
606 func NamespaceCountLimitInterceptorProvider(
607 serviceConfig *Config,
608 namespaceRegistry namespace.Registry,
609 serviceResolver membership.ServiceResolver,
610 logger log.SnTaggedLogger,
611 > ) *interceptor.ConcurrentRequestLimitInterceptor { fx.go ×44
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 {
623 fx.In
624 ServiceConfig *Config
625 NamespaceRegistry namespace.Registry
626 AdditionalAllowedMethodsDuringHandover []string `group:"additionalAllowedMethodsDuringHandover"`
627 }
628
629 func NamespaceValidatorInterceptorProvider(
630 params NamespaceValidatorInterceptorParams,
631 > ) *interceptor.NamespaceValidatorInterceptor { fx.go ×44
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 ×44
641 > return interceptor.NewSDKVersionInterceptor()
642 > }
643
644 func CallerInfoInterceptorProvider(
645 namespaceRegistry namespace.Registry,
646 > ) *interceptor.CallerInfoInterceptor { fx.go ×44
647 > return interceptor.NewCallerInfoInterceptor(namespaceRegistry)
648 > }
649
650 func SlowRequestLoggerInterceptorProvider(
651 logger log.Logger,
652 dc *dynamicconfig.Collection,
653 > ) *interceptor.SlowRequestLoggerInterceptor { fx.go ×44
654 > return interceptor.NewSlowRequestLoggerInterceptor(
655 > logger,
656 > dynamicconfig.SlowRequestLoggingThreshold.Get(dc),
657 > )
658 > }
659
660 func PersistenceRateLimitingParamsProvider(
661 serviceConfig *Config,
662 persistenceLazyLoadedServiceResolver service.PersistenceLazyLoadedServiceResolver,
663 logger log.SnTaggedLogger,
664 > ) service.PersistenceRateLimitingParams { fx.go ×44
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(
680 logger log.Logger,
681 persistenceConfig *config.Persistence,
682 customVisibilityStoreFactory visibility.VisibilityStoreFactory,
683 metricsHandler metrics.Handler,
684 serviceConfig *Config,
685 persistenceServiceResolver resolver.ServiceResolver,
686 searchAttributesMapperProvider searchattribute.MapperProvider,
687 saProvider searchattribute.Provider,
688 namespaceRegistry namespace.Registry,
689 chasmRegistry *chasm.Registry,
690 serializer serialization.Serializer,
691 > ) (manager.VisibilityManager, error) { fx.go ×44
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 ×44
721 > var replicatorNamespaceReplicationQueue persistence.NamespaceReplicationQueue
722 > if clusterMetadata.IsGlobalNamespaceEnabled() {
723 replicatorNamespaceReplicationQueue = namespaceReplicationQueue
724 }
725 > return replicatorNamespaceReplicationQueue fx.go ×44
726 }
727
728 func ServiceResolverProvider(
729 membershipMonitor membership.Monitor,
730 serviceName primitives.ServiceName,
731 > ) (membership.ServiceResolver, error) { fx.go ×44
732 > return membershipMonitor.GetResolver(serviceName)
733 > }
734
735 func AdminHandlerProvider(
736 persistenceConfig *config.Persistence,
737 configuration *Config,
738 replicatorNamespaceReplicationQueue FEReplicatorNamespaceReplicationQueue,
739 visibilityMgr manager.VisibilityManager,
740 logger log.SnTaggedLogger,
741 namespaceReplicationQueue persistence.NamespaceReplicationQueue,
742 taskManager persistence.TaskManager,
743 fairTaskManager persistence.FairTaskManager,
744 persistenceExecutionManager persistence.ExecutionManager,
745 clusterMetadataManager persistence.ClusterMetadataManager,
746 persistenceMetadataManager persistence.MetadataManager,
747 clientFactory client.Factory,
748 clientBean client.Bean,
749 historyClient resource.HistoryClient,
750 sdkClientFactory sdk.ClientFactory,
751 membershipMonitor membership.Monitor,
752 hostInfoProvider membership.HostInfoProvider,
753 metricsHandler metrics.Handler,
754 namespaceRegistry namespace.Registry,
755 saProvider searchattribute.Provider,
756 saManager searchattribute.Manager,
757 saMapperProvider searchattribute.MapperProvider,
758 clusterMetadata cluster.Metadata,
759 healthServer *health.Server,
760 eventSerializer serialization.Serializer,
761 timeSource clock.TimeSource,
762 taskCategoryRegistry tasks.TaskCategoryRegistry,
763 matchingClient resource.MatchingClient,
764 chasmRegistry *chasm.Registry,
765 namespaceDataMerger nsreplication.NamespaceDataMerger,
766 schedulerClient schedulerpb.SchedulerServiceClient,
767 namespaceDLQHandler nsreplication.DLQMessageHandler,
768 > ) *AdminHandler { fx.go ×44
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.
806 func NamespaceDLQHandlerProvider(
807 clusterMetadata cluster.Metadata,
808 persistenceMetadataManager persistence.MetadataManager,
809 namespaceDataMerger nsreplication.NamespaceDataMerger,
810 namespaceAdmitter nsreplication.NamespaceReplicationAdmitter,
811 namespaceReplicationQueue persistence.NamespaceReplicationQueue,
812 logger log.SnTaggedLogger,
813 testHooks testhooks.TestHooks,
814 > ) nsreplication.DLQMessageHandler { admin_handler.go ×3
815 > taskExecutor := nsreplication.NewTaskExecutor(
816 > clusterMetadata.GetCurrentClusterName(),
817 > persistenceMetadataManager,
818 > namespaceDataMerger,
819 > namespaceAdmitter,
820 > logger,
821 > testHooks,
822 > )
823 > return nsreplication.NewDLQMessageHandler(
824 > taskExecutor,
825 > namespaceReplicationQueue,
826 > logger,
827 > )
828 > }
829
830 func OperatorHandlerProvider(
831 configuration *Config,
832 logger log.SnTaggedLogger,
833 sdkClientFactory sdk.ClientFactory,
834 metricsHandler metrics.Handler,
835 visibilityMgr manager.VisibilityManager,
836 saManager searchattribute.Manager,
837 healthServer *health.Server,
838 historyClient resource.HistoryClient,
839 clusterMetadataManager persistence.ClusterMetadataManager,
840 clusterMetadata cluster.Metadata,
841 clientFactory client.Factory,
842 namespaceRegistry namespace.Registry,
843 nexusEndpointClient *NexusEndpointClient,
844 > ) *OperatorHandlerImpl { fx.go ×44
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 ×44
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(
875 dc *dynamicconfig.Collection,
876 cfg *config.Config,
877 serviceName primitives.ServiceName,
878 dcRedirectionPolicy config.DCRedirectionPolicy,
879 serviceConfig *Config,
880 versionChecker *VersionChecker,
881 namespaceReplicationQueue FEReplicatorNamespaceReplicationQueue,
882 visibilityMgr manager.VisibilityManager,
883 chasmVisibilityMgr chasm.VisibilityManager,
884 logger log.SnTaggedLogger,
885 throttledLogger log.ThrottledLogger,
886 persistenceExecutionManager persistence.ExecutionManager,
887 clusterMetadataManager persistence.ClusterMetadataManager,
888 persistenceMetadataManager persistence.MetadataManager,
889 clientBean client.Bean,
890 historyClient resource.HistoryClient,
891 matchingClient resource.MatchingClient,
892 workerDeploymentStoreClient workerdeployment.Client,
893 schedulerClient schedulerpb.SchedulerServiceClient,
894 archiverProvider provider.ArchiverProvider,
895 metricsHandler metrics.Handler,
896 payloadSerializer serialization.Serializer,
897 timeSource clock.TimeSource,
898 namespaceRegistry namespace.Registry,
899 saMapperProvider searchattribute.MapperProvider,
900 saProvider searchattribute.Provider,
901 saValidator *searchattribute.Validator,
902 clusterMetadata cluster.Metadata,
903 archivalMetadata archiver.ArchivalMetadata,
904 healthServer *health.Server,
905 membershipMonitor membership.Monitor,
906 healthInterceptor *interceptor.HealthInterceptor,
907 scheduleSpecBuilder *scheduler.SpecBuilder,
908 activityHandler activity.FrontendHandler,
909 callbackValidator callback.Validator,
910 nexusOperationHandler chasmnexus.FrontendHandler,
911 registry *chasm.Registry,
912 frontendServiceResolver membership.ServiceResolver,
913 > ) Handler { fx.go ×44
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 ×44
966 > h.RegisterRoutes(router)
967 > }
968
969 func RegisterNexusCompletionHTTPHandler(
970 h *nexusCompletionHTTPHandler,
971 router *mux.Router,
972 > ) { fx.go ×44
973 > h.RegisterRoutes(router)
974 > }
975
976 func RegisterOpenAPIHTTPHandler(
977 rateLimitInterceptor *interceptor.RateLimitInterceptor,
978 logger log.Logger,
979 router *mux.Router,
980 > ) *OpenAPIHTTPHandler { fx.go ×44
981 > h := NewOpenAPIHTTPHandler(
982 > rateLimitInterceptor,
983 > logger,
984 > )
985 > h.RegisterRoutes(router)
986 > return h
987 > }
988
989 > func MuxRouterProvider() *mux.Router { fx.go ×44
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 ×44
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 ×44
1001 }
1002
1003 // HTTPAPIServerProvider provides an HTTP API server if enabled or nil
1004 // otherwise.
1005 func HTTPAPIServerProvider(
1006 cfg *config.Config,
1007 serviceName primitives.ServiceName,
1008 serviceConfig *Config,
1009 grpcListener net.Listener,
1010 tlsConfigProvider encryption.TLSConfigProvider,
1011 handler Handler,
1012 operatorHandler *OperatorHandlerImpl,
1013 grpcServerOptions GrpcServerOptions,
1014 metricsHandler metrics.Handler,
1015 namespaceRegistry namespace.Registry,
1016 logger log.Logger,
1017 router *mux.Router,
1018 > ) (*HTTPAPIServer, error) { fx.go ×44
1019 > if !httpEnabled(cfg, serviceName) {
1020 > return nil, nil lite_server.go ×25
1021 > }
1022 > rpcConfig := cfg.Services[string(serviceName)].RPC http_api_server.go ×23
1023 > return NewHTTPAPIServer(
1024 > serviceConfig,
1025 > rpcConfig,
1026 > grpcListener,
1027 > tlsConfigProvider,
1028 > handler,
1029 > operatorHandler,
1030 > grpcServerOptions.UnaryInterceptors,
1031 > metricsHandler,
1032 > router,
1033 > namespaceRegistry,
1034 > logger,
1035 > )
1036 }
1037
1038 func NexusEndpointClientProvider(
1039 dc *dynamicconfig.Collection,
1040 namespaceRegistry namespace.Registry,
1041 matchingClient resource.MatchingClient,
1042 nexusEndpointManager persistence.NexusEndpointManager,
1043 logger log.Logger,
1044 > ) *NexusEndpointClient { fx.go ×44
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 ×44
1056 > lc.Append(fx.StartStopHook(svc.Start, svc.Stop))
1057 > }