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.
package frontend
import (
"fmt"
"net"
"github.com/gorilla/mux"
"go.temporal.io/server/api/adminservice/v1"
"go.temporal.io/server/chasm"
"go.temporal.io/server/chasm/lib/activity"
"go.temporal.io/server/chasm/lib/callback"
chasmnexus "go.temporal.io/server/chasm/lib/nexusoperation"
nexusoperationpb "go.temporal.io/server/chasm/lib/nexusoperation/gen/nexusoperationpb/v1"
chasmscheduler "go.temporal.io/server/chasm/lib/scheduler"
"go.temporal.io/server/chasm/lib/scheduler/gen/schedulerpb/v1"
chasmtests "go.temporal.io/server/chasm/lib/tests"
chasmworkflow "go.temporal.io/server/chasm/lib/workflow"
"go.temporal.io/server/client"
"go.temporal.io/server/common"
"go.temporal.io/server/common/archiver"
"go.temporal.io/server/common/archiver/provider"
"go.temporal.io/server/common/authorization"
"go.temporal.io/server/common/clock"
"go.temporal.io/server/common/cluster"
"go.temporal.io/server/common/config"
"go.temporal.io/server/common/dynamicconfig"
"go.temporal.io/server/common/log"
"go.temporal.io/server/common/log/tag"
"go.temporal.io/server/common/membership"
"go.temporal.io/server/common/metrics"
"go.temporal.io/server/common/namespace"
"go.temporal.io/server/common/namespace/nsreplication"
"go.temporal.io/server/common/persistence"
"go.temporal.io/server/common/persistence/serialization"
"go.temporal.io/server/common/persistence/visibility"
"go.temporal.io/server/common/persistence/visibility/manager"
"go.temporal.io/server/common/primitives"
"go.temporal.io/server/common/quotas"
"go.temporal.io/server/common/quotas/calculator"
"go.temporal.io/server/common/resolver"
"go.temporal.io/server/common/resource"
"go.temporal.io/server/common/rpc"
"go.temporal.io/server/common/rpc/encryption"
"go.temporal.io/server/common/rpc/interceptor"
"go.temporal.io/server/common/sdk"
"go.temporal.io/server/common/searchattribute"
"go.temporal.io/server/common/telemetry"
"go.temporal.io/server/common/testing/testhooks"
"go.temporal.io/server/service"
"go.temporal.io/server/service/frontend/configs"
"go.temporal.io/server/service/history/tasks"
"go.temporal.io/server/service/worker/scheduler"
"go.temporal.io/server/service/worker/workerdeployment"
"go.uber.org/fx"
"google.golang.org/grpc"
"google.golang.org/grpc/health"
healthpb "google.golang.org/grpc/health/grpc_health_v1"
"google.golang.org/grpc/keepalive"
)
type (
FEReplicatorNamespaceReplicationQueue persistence.NamespaceReplicationQueue
namespaceChecker struct {
r namespace.Registry
}
)
var Module = fx.Options(
resource.Module,
chasmtests.Module,
scheduler.Module,
workerdeployment.Module,
// Note that with this approach routes may be registered in arbitrary order.
// This is okay because our routes don't have overlapping matches.
// The only important detail is that the PathPrefix("/") route registered in the HTTPAPIServerProvider comes last.
// Coincidentally, this is the case today, likely because it has more dependencies that the other dependencies.
// This approach isn't perfect but at it allows the router to be pluggable and we have enough functional test
// coverage to catch misconfiguration.
// A more robust approach would require using fx groups but we shouldn't overcomplicate until this becomes an issue.
fx.Provide(MuxRouterProvider),
fx.Provide(ConfigProvider),
fx.Provide(ServiceErrorInterceptorProvider),
fx.Provide(NamespaceLogInterceptorProvider),
fx.Provide(NamespaceHandoverInterceptorProvider),
fx.Provide(interceptor.NewRoutingKeyExtractor),
fx.Provide(BusinessIDInterceptorProvider),
fx.Provide(RedirectionInterceptorProvider),
fx.Provide(ErrorHandlerProvider),
fx.Provide(TelemetryInterceptorProvider),
fx.Provide(RetryableInterceptorProvider),
fx.Provide(RateLimitInterceptorProvider),
fx.Provide(interceptor.NewHealthInterceptor),
fx.Provide(NamespaceCountLimitInterceptorProvider),
fx.Provide(NamespaceValidatorInterceptorProvider),
fx.Provide(NamespaceRateLimitInterceptorProvider),
fx.Provide(SDKVersionInterceptorProvider),
fx.Provide(CallerInfoInterceptorProvider),
fx.Provide(SlowRequestLoggerInterceptorProvider),
fx.Provide(MaskInternalErrorDetailsInterceptorProvider),
fx.Provide(ContextMetadataInterceptorProvider),
fx.Provide(GrpcServerOptionsProvider),
fx.Provide(VisibilityManagerProvider),
fx.Provide(ThrottledLoggerRpsFnProvider),
fx.Provide(PersistenceRateLimitingParamsProvider),
service.PersistenceLazyLoadedServiceResolverModule,
fx.Provide(FEReplicatorNamespaceReplicationQueueProvider),
fx.Provide(nsreplication.NewNoopDataMerger),
fx.Provide(nsreplication.NewDefaultAdmitter),
fx.Provide(AuthorizationInterceptorProvider),
fx.Provide(NamespaceCheckerProvider),
fx.Provide(func(so GrpcServerOptions) *grpc.Server { return grpc.NewServer(so.Options...) }),
fx.go ×44
fx.Provide(callbackValidatorProvider),
fx.Provide(HandlerProvider),
fx.Provide(AdminHandlerProvider),
fx.Provide(NamespaceDLQHandlerProvider),
fx.Provide(OperatorHandlerProvider),
fx.Provide(NewVersionChecker),
fx.Provide(ServiceResolverProvider),
fx.Provide(newNexusCompletionHandler),
fx.Provide(NewNexusOperationHTTPHandler),
fx.Provide(newNexusCompletionHTTPHandler),
fx.Invoke(RegisterNexusOperationHTTPHandler),
fx.Invoke(RegisterNexusCompletionHTTPHandler),
fx.Invoke(RegisterOpenAPIHTTPHandler),
fx.Provide(HTTPAPIServerProvider),
fx.Provide(NewServiceProvider),
fx.Provide(NexusEndpointClientProvider),
fx.Invoke(ServiceLifetimeHooks),
fx.Provide(nexusoperationpb.NewNexusOperationServiceLayeredClient),
fx.Provide(schedulerpb.NewSchedulerServiceLayeredClient),
fx.Provide(chasmnexus.NewFrontendHandler),
chasmnexus.Module,
chasmscheduler.Module,
chasmworkflow.Module,
callback.Module,
activity.FrontendModule,
fx.Provide(visibility.ChasmVisibilityManagerProvider),
fx.Provide(chasm.ChasmVisibilityInterceptorProvider),
)
func NewServiceProvider(
serviceConfig *Config,
server *grpc.Server,
healthServer *health.Server,
httpAPIServer *HTTPAPIServer,
handler Handler,
adminHandler *AdminHandler,
operatorHandler *OperatorHandlerImpl,
versionChecker *VersionChecker,
visibilityMgr manager.VisibilityManager,
logger log.SnTaggedLogger,
grpcListener net.Listener,
metricsHandler metrics.Handler,
membershipMonitor membership.Monitor,
return NewService(
serviceConfig,
server,
healthServer,
httpAPIServer,
handler,
adminHandler,
operatorHandler,
versionChecker,
visibilityMgr,
logger,
grpcListener,
metricsHandler,
membershipMonitor,
)
}
// GrpcServerOptions are the options to build the frontend gRPC server along
// with the interceptors that are already set in the options.
type GrpcServerOptions struct {
Options []grpc.ServerOption
UnaryInterceptors []grpc.UnaryServerInterceptor
}
func AuthorizationInterceptorProvider(
cfg *config.Config,
serviceConfig *Config,
logger log.Logger,
namespaceChecker authorization.NamespaceChecker,
metricsHandler metrics.Handler,
authorizer authorization.Authorizer,
claimMapper authorization.ClaimMapper,
audienceGetter authorization.JWTAudienceMapper,
dc *dynamicconfig.Collection,
return authorization.NewInterceptor(
claimMapper,
authorizer,
metricsHandler,
logger,
namespaceChecker,
audienceGetter,
cfg.Global.Authorization.AuthHeaderName,
cfg.Global.Authorization.AuthExtraHeaderName,
serviceConfig.ExposeAuthorizerErrors,
dynamicconfig.EnableCrossNamespaceCommands.Get(dc),
dynamicconfig.EnablePrincipalPropagation.Get(dc),
dynamicconfig.DisableStreamingAuthorizer.Get(dc),
)
}
func NamespaceCheckerProvider(registry namespace.Registry) authorization.NamespaceChecker {
fx.go ×44
return &namespaceChecker{r: registry}
}
// This will get called before the namespace state validation interceptor. We want to
// disable readthrough to avoid polluting the negative lookup cache, e.g. if this call is
// for RegisterNamespace and the namespace doesn't exist yet.
opts := namespace.GetNamespaceOptions{DisableReadthrough: true}
_, err := n.r.GetNamespaceWithOptions(name, opts)
return err
}
func GrpcServerOptionsProvider(
logger log.Logger,
cfg *config.Config,
serviceConfig *Config,
serviceName primitives.ServiceName,
rpcFactory common.RPCFactory,
serviceErrorInterceptor *interceptor.ServiceErrorInterceptor,
namespaceLogInterceptor *interceptor.NamespaceLogInterceptor,
namespaceRateLimiterInterceptor interceptor.NamespaceRateLimitInterceptor,
namespaceCountLimiterInterceptor *interceptor.ConcurrentRequestLimitInterceptor,
namespaceValidatorInterceptor *interceptor.NamespaceValidatorInterceptor,
namespaceHandoverInterceptor *interceptor.NamespaceHandoverInterceptor,
businessIDInterceptor *interceptor.RoutingKeyInterceptor,
redirectionInterceptor *interceptor.Redirection,
telemetryInterceptor *interceptor.TelemetryInterceptor,
retryableInterceptor *interceptor.RetryableInterceptor,
healthInterceptor *interceptor.HealthInterceptor,
rateLimitInterceptor *interceptor.RateLimitInterceptor,
traceStatsHandler telemetry.ServerStatsHandler,
metricsStatsHandler metrics.ServerStatsHandler,
sdkVersionInterceptor *interceptor.SDKVersionInterceptor,
callerInfoInterceptor *interceptor.CallerInfoInterceptor,
authInterceptor *authorization.Interceptor,
maskInternalErrorDetailsInterceptor *interceptor.MaskInternalErrorDetailsInterceptor,
contextMetadataInterceptor *interceptor.ContextMetadataInterceptor,
slowRequestLoggerInterceptor *interceptor.SlowRequestLoggerInterceptor,
chasmRequestVisibilityInterceptor *chasm.ChasmVisibilityInterceptor,
customInterceptors []grpc.UnaryServerInterceptor,
customStreamInterceptors []grpc.StreamServerInterceptor,
metricsHandler metrics.Handler,
kep := keepalive.EnforcementPolicy{
MinTime: serviceConfig.KeepAliveMinTime(),
PermitWithoutStream: serviceConfig.KeepAlivePermitWithoutStream(),
}
kp := keepalive.ServerParameters{
MaxConnectionIdle: serviceConfig.KeepAliveMaxConnectionIdle(),
MaxConnectionAge: serviceConfig.KeepAliveMaxConnectionAge(),
MaxConnectionAgeGrace: serviceConfig.KeepAliveMaxConnectionAgeGrace(),
Time: serviceConfig.KeepAliveTime(),
Timeout: serviceConfig.KeepAliveTimeout(),
}
var grpcServerOptions []grpc.ServerOption
var err error
switch serviceName {
case primitives.FrontendService:
grpcServerOptions, err = rpcFactory.GetFrontendGRPCServerOptions()
case primitives.InternalFrontendService:
grpcServerOptions, err = rpcFactory.GetInternodeGRPCServerOptions()
default:
err = fmt.Errorf("unexpected frontend service name %q", serviceName)
}
logger.Fatal("creating gRPC server options failed", tag.Error(err))
}
// Order of interceptors is important
// Mask error interceptor should be the most outer interceptor since it handle the errors format
// Service Error Interceptor should be the next most outer interceptor on error handling
maskInternalErrorDetailsInterceptor.Intercept,
serviceErrorInterceptor.Intercept,
interceptor.NewFrontendServiceErrorInterceptor(logger),
// BusinessID interceptor extracts business ID and adds it to context for use, must be before any interceptor that touches namespaces (namespaceValidator, handoverInterceptor)
businessIDInterceptor.Intercept,
namespaceValidatorInterceptor.NamespaceValidateIntercept,
namespaceLogInterceptor.Intercept, // TODO: Deprecate this with a outer custom interceptor
metrics.NewServerMetricsContextInjectorInterceptor(),
authInterceptor.Intercept,
// Handover interceptor has to above redirection because the request will route to the correct cluster after handover completed.
// And retry cannot be performed before customInterceptors.
namespaceHandoverInterceptor.Intercept,
redirectionInterceptor.Intercept,
// Telemetry interceptor must be after redirection to ensure metrics are recorded in the correct cluster
telemetryInterceptor.UnaryIntercept,
healthInterceptor.Intercept,
namespaceValidatorInterceptor.StateValidationIntercept,
namespaceCountLimiterInterceptor.Intercept,
namespaceRateLimiterInterceptor.Intercept,
rateLimitInterceptor.Intercept,
sdkVersionInterceptor.Intercept,
callerInfoInterceptor.Intercept,
slowRequestLoggerInterceptor.Intercept,
chasmRequestVisibilityInterceptor.Intercept,
contextMetadataInterceptor.Intercept,
}
if len(customInterceptors) > 0 {
// TODO: Deprecate WithChainedFrontendGrpcInterceptors and provide a inner custom interceptor
http_api_server.go ×23
unaryInterceptors = append(unaryInterceptors, customInterceptors...)
}
// retry interceptor should be the most inner interceptor
streamInterceptor := []grpc.StreamServerInterceptor{
authInterceptor.InterceptStream,
telemetryInterceptor.StreamIntercept,
}
if len(customStreamInterceptors) > 0 {
}
grpcServerOptions,
grpc.KeepaliveParams(kp),
grpc.KeepaliveEnforcementPolicy(kep),
grpc.ChainUnaryInterceptor(unaryInterceptors...),
grpc.ChainStreamInterceptor(streamInterceptor...),
)
multiStats := rpc.MultiStatsHandler{}
if traceStatsHandler != nil {
}
multiStats = append(multiStats, metricsStatsHandler)
}
if len(multiStats) > 0 {
grpcServerOptions = append(grpcServerOptions, grpc.StatsHandler(multiStats))
}
return GrpcServerOptions{Options: grpcServerOptions, UnaryInterceptors: unaryInterceptors}
}
func ConfigProvider(
dc *dynamicconfig.Collection,
persistenceConfig config.Persistence,
return NewConfig(
dc,
persistenceConfig.NumHistoryShards,
)
}
func ServiceErrorInterceptorProvider(
dc *dynamicconfig.Collection,
return interceptor.NewServiceErrorInterceptor(
dynamicconfig.MaxServiceErrorMessageLength.Get(dc),
)
}
func ThrottledLoggerRpsFnProvider(serviceConfig *Config) resource.ThrottledLoggerRpsFn {
fx.go ×44
return func() float64 { return float64(serviceConfig.ThrottledLogRPS()) }
}
func NamespaceLogInterceptorProvider(
namespaceLogger resource.NamespaceLogger,
namespaceRegistry namespace.Registry,
return interceptor.NewNamespaceLogInterceptor(
namespaceRegistry,
namespaceLogger)
}
return interceptor.NewRetryableInterceptor(
common.CreateFrontendHandlerRetryPolicy(),
common.IsServiceHandlerRetryableError,
)
}
func RedirectionInterceptorProvider(
configuration *Config,
namespaceCache namespace.Registry,
policy config.DCRedirectionPolicy,
logger log.Logger,
clientBean client.Bean,
metricsHandler metrics.Handler,
timeSource clock.TimeSource,
clusterMetadata cluster.Metadata,
return interceptor.NewRedirection(
configuration.EnableNamespaceNotActiveAutoForwarding,
configuration.ForceNamespaceSelectedAPIAutoForwarding,
namespaceCache,
policy,
logger,
clientBean,
metricsHandler,
timeSource,
clusterMetadata,
)
}
func BusinessIDInterceptorProvider(
extractor interceptor.RoutingKeyExtractor,
logger log.Logger,
return interceptor.NewRoutingKeyInterceptor(
[]interceptor.RoutingKeyExtractorFunc{
interceptor.WorkflowServiceExtractor(extractor),
},
logger,
)
}
type NamespaceHandoverInterceptorParams struct {
fx.In
DynamicConfig *dynamicconfig.Collection
NamespaceRegistry namespace.Registry
Logger log.Logger
MetricsHandler metrics.Handler
TimeSource clock.TimeSource
RequestErrorHandler *interceptor.RequestErrorHandler
AdditionalAllowedMethodsDuringHandover []string `group:"additionalAllowedMethodsDuringHandover"`
}
func NamespaceHandoverInterceptorProvider(
params NamespaceHandoverInterceptorParams,
return interceptor.NewNamespaceHandoverInterceptor(
params.DynamicConfig,
params.NamespaceRegistry,
params.MetricsHandler,
params.Logger,
params.TimeSource,
params.RequestErrorHandler,
params.AdditionalAllowedMethodsDuringHandover,
)
}
func ErrorHandlerProvider(
logger log.Logger,
serviceConfig *Config,
return interceptor.NewRequestErrorHandler(
logger,
serviceConfig.LogAllReqErrors,
)
}
func TelemetryInterceptorProvider(
logger log.Logger,
metricsHandler metrics.Handler,
namespaceRegistry namespace.Registry,
serviceConfig *Config,
requestErrorHandler *interceptor.RequestErrorHandler,
return interceptor.NewTelemetryInterceptor(
namespaceRegistry,
metricsHandler,
logger,
serviceConfig.LogAllReqErrors,
requestErrorHandler,
)
}
func getRateFnWithMetrics(rateFn quotas.RateFn, handler metrics.Handler) quotas.RateFn {
fx.go ×3
return func() float64 {
rate := rateFn()
metrics.HostRPSLimit.With(handler).Record(rate)
return rate
}
}
func RateLimitInterceptorProvider(
serviceConfig *Config,
frontendServiceResolver membership.ServiceResolver,
handler metrics.Handler,
logger log.SnTaggedLogger,
rateFn := calculator.NewLoggedCalculator(
calculator.ClusterAwareQuotaCalculator{
MemberCounter: frontendServiceResolver,
PerInstanceQuota: serviceConfig.RPS,
GlobalQuota: serviceConfig.GlobalRPS,
},
log.With(logger, tag.ComponentRPCHandler, tag.ScopeHost),
).GetQuota
rateFnWithMetrics := getRateFnWithMetrics(rateFn, handler)
namespaceReplicationInducingRateFn := func() float64 {
return float64(serviceConfig.NamespaceReplicationInducingAPIsRPS())
}
configs.NewRequestToRateLimiter(
quotas.NewDefaultIncomingRateBurst(rateFnWithMetrics),
quotas.NewDefaultIncomingRateBurst(rateFn),
quotas.NewDefaultIncomingRateBurst(namespaceReplicationInducingRateFn),
serviceConfig.OperatorRPSRatio,
),
map[string]int{
healthpb.Health_Check_FullMethodName: 0, // exclude health check requests from rate limiting.
adminservice.AdminService_DeepHealthCheck_FullMethodName: 0, // exclude deep health check requests from rate limiting.
},
)
}
func ContextMetadataInterceptorProvider(
logger log.Logger,
dc *dynamicconfig.Collection,
setTrailer := dynamicconfig.FrontendContextMetadataSetTrailer.Get(dc)()
return interceptor.NewContextMetadataInterceptor(setTrailer, logger)
}
func MaskInternalErrorDetailsInterceptorProvider(
logger log.Logger,
serviceConfig *Config,
namespaceRegistry namespace.Registry,
return interceptor.NewMaskInternalErrorDetailsInterceptor(
serviceConfig.MaskInternalErrorDetails, namespaceRegistry, logger,
)
}
func NamespaceRateLimitInterceptorProvider(
serviceName primitives.ServiceName,
serviceConfig *Config,
namespaceRegistry namespace.Registry,
frontendServiceResolver membership.ServiceResolver,
metricsHandler metrics.Handler,
logger log.SnTaggedLogger,
var globalNamespaceRPS, globalNamespaceVisibilityRPS, globalNamespaceNamespaceReplicationInducingAPIsRPS dynamicconfig.IntPropertyFnWithNamespaceFilter
switch serviceName {
case primitives.FrontendService:
globalNamespaceRPS = serviceConfig.GlobalNamespaceRPS
globalNamespaceVisibilityRPS = serviceConfig.GlobalNamespaceVisibilityRPS
globalNamespaceNamespaceReplicationInducingAPIsRPS = serviceConfig.GlobalNamespaceNamespaceReplicationInducingAPIsRPS
case primitives.InternalFrontendService:
globalNamespaceRPS = serviceConfig.InternalFEGlobalNamespaceRPS
globalNamespaceVisibilityRPS = serviceConfig.InternalFEGlobalNamespaceVisibilityRPS
// Internal frontend has no special limit for this set of APIs
globalNamespaceNamespaceReplicationInducingAPIsRPS = serviceConfig.InternalFEGlobalNamespaceRPS
default:
panic("invalid service name")
}
calculator.ClusterAwareNamespaceQuotaCalculator{
MemberCounter: frontendServiceResolver,
PerInstanceQuota: serviceConfig.MaxNamespaceRPSPerInstance,
GlobalQuota: globalNamespaceRPS,
},
log.With(logger, tag.ComponentRPCHandler, tag.ScopeNamespace),
).GetQuota
visibilityRateFn := calculator.NewLoggedNamespaceCalculator(
calculator.ClusterAwareNamespaceQuotaCalculator{
MemberCounter: frontendServiceResolver,
PerInstanceQuota: serviceConfig.MaxNamespaceVisibilityRPSPerInstance,
GlobalQuota: globalNamespaceVisibilityRPS,
},
log.With(logger, tag.ComponentVisibilityHandler, tag.ScopeNamespace),
).GetQuota
namespaceReplicationInducingRateFn := calculator.NewLoggedNamespaceCalculator(
calculator.ClusterAwareNamespaceQuotaCalculator{
MemberCounter: frontendServiceResolver,
PerInstanceQuota: serviceConfig.MaxNamespaceNamespaceReplicationInducingAPIsRPSPerInstance,
GlobalQuota: globalNamespaceNamespaceReplicationInducingAPIsRPS,
},
log.With(logger, tag.ComponentNamespaceReplication, tag.ScopeNamespace),
).GetQuota
namespaceRateLimiter := quotas.NewNamespaceRequestRateLimiter(
func(req quotas.Request) quotas.RequestRateLimiter {
quotas.NewNamespaceRateBurst(
req.Caller,
namespaceRateFn,
quotas.NamespaceBurstRatioFn(serviceConfig.MaxNamespaceBurstRatioPerInstance),
),
quotas.NewNamespaceRateBurst(
req.Caller,
visibilityRateFn,
quotas.NamespaceBurstRatioFn(serviceConfig.MaxNamespaceVisibilityBurstRatioPerInstance),
),
quotas.NewNamespaceRateBurst(
req.Caller,
namespaceReplicationInducingRateFn,
quotas.NamespaceBurstRatioFn(serviceConfig.MaxNamespaceNamespaceReplicationInducingAPIsBurstRatioPerInstance),
),
serviceConfig.OperatorRPSRatio,
)
},
)
namespaceRegistry,
namespaceRateLimiter,
map[string]int{}, // no token overrides
configs.PollTaskAPISet,
serviceConfig.PollWaitForNamespaceRateLimitToken,
metricsHandler,
)
}
func NamespaceCountLimitInterceptorProvider(
serviceConfig *Config,
namespaceRegistry namespace.Registry,
serviceResolver membership.ServiceResolver,
logger log.SnTaggedLogger,
return interceptor.NewConcurrentRequestLimitInterceptor(
namespaceRegistry,
serviceResolver,
logger,
serviceConfig.MaxConcurrentLongRunningRequestsPerInstance,
serviceConfig.MaxGlobalConcurrentLongRunningRequests,
configs.ExecutionAPICountLimitOverride,
)
}
type NamespaceValidatorInterceptorParams struct {
fx.In
ServiceConfig *Config
NamespaceRegistry namespace.Registry
AdditionalAllowedMethodsDuringHandover []string `group:"additionalAllowedMethodsDuringHandover"`
}
func NamespaceValidatorInterceptorProvider(
params NamespaceValidatorInterceptorParams,
return interceptor.NewNamespaceValidatorInterceptor(
params.NamespaceRegistry,
params.ServiceConfig.EnableTokenNamespaceEnforcement,
params.ServiceConfig.MaxIDLengthLimit,
params.AdditionalAllowedMethodsDuringHandover,
)
}
return interceptor.NewSDKVersionInterceptor()
}
func CallerInfoInterceptorProvider(
namespaceRegistry namespace.Registry,
return interceptor.NewCallerInfoInterceptor(namespaceRegistry)
}
func SlowRequestLoggerInterceptorProvider(
logger log.Logger,
dc *dynamicconfig.Collection,
return interceptor.NewSlowRequestLoggerInterceptor(
logger,
dynamicconfig.SlowRequestLoggingThreshold.Get(dc),
)
}
func PersistenceRateLimitingParamsProvider(
serviceConfig *Config,
persistenceLazyLoadedServiceResolver service.PersistenceLazyLoadedServiceResolver,
logger log.SnTaggedLogger,
return service.NewPersistenceRateLimitingParams(
serviceConfig.PersistenceMaxQPS,
serviceConfig.PersistenceGlobalMaxQPS,
serviceConfig.PersistenceNamespaceMaxQPS,
serviceConfig.PersistenceGlobalNamespaceMaxQPS,
serviceConfig.PersistencePerShardNamespaceMaxQPS,
serviceConfig.OperatorRPSRatio,
serviceConfig.PersistenceQPSBurstRatio,
serviceConfig.PersistenceDynamicRateLimitingParams,
persistenceLazyLoadedServiceResolver,
logger,
)
}
func VisibilityManagerProvider(
logger log.Logger,
persistenceConfig *config.Persistence,
customVisibilityStoreFactory visibility.VisibilityStoreFactory,
metricsHandler metrics.Handler,
serviceConfig *Config,
persistenceServiceResolver resolver.ServiceResolver,
searchAttributesMapperProvider searchattribute.MapperProvider,
saProvider searchattribute.Provider,
namespaceRegistry namespace.Registry,
chasmRegistry *chasm.Registry,
serializer serialization.Serializer,
return visibility.NewManager(
*persistenceConfig,
persistenceServiceResolver,
customVisibilityStoreFactory,
nil, // frontend visibility never write
saProvider,
searchAttributesMapperProvider,
namespaceRegistry,
chasmRegistry,
serviceConfig.VisibilityPersistenceMaxReadQPS,
serviceConfig.VisibilityPersistenceMaxWriteQPS,
serviceConfig.OperatorRPSRatio,
serviceConfig.VisibilityPersistenceSlowQueryThreshold,
serviceConfig.EnableReadFromSecondaryVisibility,
serviceConfig.VisibilityEnableShadowReadMode,
dynamicconfig.GetStringPropertyFn(visibility.SecondaryVisibilityWritingModeOff), // frontend visibility never write
serviceConfig.VisibilityDisableOrderByClause,
serviceConfig.VisibilityEnableManualPagination,
serviceConfig.VisibilityEnableUnifiedQueryConverter,
metricsHandler,
logger,
serializer,
)
}
func FEReplicatorNamespaceReplicationQueueProvider(
namespaceReplicationQueue persistence.NamespaceReplicationQueue,
clusterMetadata cluster.Metadata,
var replicatorNamespaceReplicationQueue persistence.NamespaceReplicationQueue
if clusterMetadata.IsGlobalNamespaceEnabled() {
replicatorNamespaceReplicationQueue = namespaceReplicationQueue
}
}
func ServiceResolverProvider(
membershipMonitor membership.Monitor,
serviceName primitives.ServiceName,
return membershipMonitor.GetResolver(serviceName)
}
func AdminHandlerProvider(
persistenceConfig *config.Persistence,
configuration *Config,
replicatorNamespaceReplicationQueue FEReplicatorNamespaceReplicationQueue,
visibilityMgr manager.VisibilityManager,
logger log.SnTaggedLogger,
namespaceReplicationQueue persistence.NamespaceReplicationQueue,
taskManager persistence.TaskManager,
fairTaskManager persistence.FairTaskManager,
persistenceExecutionManager persistence.ExecutionManager,
clusterMetadataManager persistence.ClusterMetadataManager,
persistenceMetadataManager persistence.MetadataManager,
clientFactory client.Factory,
clientBean client.Bean,
historyClient resource.HistoryClient,
sdkClientFactory sdk.ClientFactory,
membershipMonitor membership.Monitor,
hostInfoProvider membership.HostInfoProvider,
metricsHandler metrics.Handler,
namespaceRegistry namespace.Registry,
saProvider searchattribute.Provider,
saManager searchattribute.Manager,
saMapperProvider searchattribute.MapperProvider,
clusterMetadata cluster.Metadata,
healthServer *health.Server,
eventSerializer serialization.Serializer,
timeSource clock.TimeSource,
taskCategoryRegistry tasks.TaskCategoryRegistry,
matchingClient resource.MatchingClient,
chasmRegistry *chasm.Registry,
namespaceDataMerger nsreplication.NamespaceDataMerger,
schedulerClient schedulerpb.SchedulerServiceClient,
namespaceDLQHandler nsreplication.DLQMessageHandler,
args := NewAdminHandlerArgs{
persistenceConfig,
configuration,
namespaceReplicationQueue,
replicatorNamespaceReplicationQueue,
visibilityMgr,
logger,
taskManager,
fairTaskManager,
persistenceExecutionManager,
clusterMetadataManager,
persistenceMetadataManager,
clientFactory,
clientBean,
historyClient,
sdkClientFactory,
membershipMonitor,
hostInfoProvider,
metricsHandler,
namespaceRegistry,
saProvider,
saManager,
saMapperProvider,
clusterMetadata,
healthServer,
eventSerializer,
timeSource,
chasmRegistry,
namespaceDataMerger,
schedulerClient,
taskCategoryRegistry,
matchingClient,
}
return NewAdminHandler(args, namespaceDLQHandler)
}
// NamespaceDLQHandlerProvider provides the default namespace DLQ message handler.
func NamespaceDLQHandlerProvider(
clusterMetadata cluster.Metadata,
persistenceMetadataManager persistence.MetadataManager,
namespaceDataMerger nsreplication.NamespaceDataMerger,
namespaceAdmitter nsreplication.NamespaceReplicationAdmitter,
namespaceReplicationQueue persistence.NamespaceReplicationQueue,
logger log.SnTaggedLogger,
testHooks testhooks.TestHooks,
taskExecutor := nsreplication.NewTaskExecutor(
clusterMetadata.GetCurrentClusterName(),
persistenceMetadataManager,
namespaceDataMerger,
namespaceAdmitter,
logger,
testHooks,
)
return nsreplication.NewDLQMessageHandler(
taskExecutor,
namespaceReplicationQueue,
logger,
)
}
func OperatorHandlerProvider(
configuration *Config,
logger log.SnTaggedLogger,
sdkClientFactory sdk.ClientFactory,
metricsHandler metrics.Handler,
visibilityMgr manager.VisibilityManager,
saManager searchattribute.Manager,
healthServer *health.Server,
historyClient resource.HistoryClient,
clusterMetadataManager persistence.ClusterMetadataManager,
clusterMetadata cluster.Metadata,
clientFactory client.Factory,
namespaceRegistry namespace.Registry,
nexusEndpointClient *NexusEndpointClient,
args := NewOperatorHandlerImplArgs{
configuration,
logger,
sdkClientFactory,
metricsHandler,
visibilityMgr,
saManager,
healthServer,
historyClient,
clusterMetadataManager,
clusterMetadata,
clientFactory,
namespaceRegistry,
nexusEndpointClient,
}
return NewOperatorHandlerImpl(args)
}
// callbackValidatorProvider creates a callback Validator using the production dynamic config keys
// so that existing operator configurations (callback.allowedAddresses) are honored.
return callback.NewValidator(
callback.MaxPerExecution.Get(dc),
dynamicconfig.FrontendCallbackURLMaxLength.Get(dc),
dynamicconfig.FrontendCallbackHeaderMaxSize.Get(dc),
callback.AllowedAddresses.Get(dc),
)
}
func HandlerProvider(
dc *dynamicconfig.Collection,
cfg *config.Config,
serviceName primitives.ServiceName,
dcRedirectionPolicy config.DCRedirectionPolicy,
serviceConfig *Config,
versionChecker *VersionChecker,
namespaceReplicationQueue FEReplicatorNamespaceReplicationQueue,
visibilityMgr manager.VisibilityManager,
chasmVisibilityMgr chasm.VisibilityManager,
logger log.SnTaggedLogger,
throttledLogger log.ThrottledLogger,
persistenceExecutionManager persistence.ExecutionManager,
clusterMetadataManager persistence.ClusterMetadataManager,
persistenceMetadataManager persistence.MetadataManager,
clientBean client.Bean,
historyClient resource.HistoryClient,
matchingClient resource.MatchingClient,
workerDeploymentStoreClient workerdeployment.Client,
schedulerClient schedulerpb.SchedulerServiceClient,
archiverProvider provider.ArchiverProvider,
metricsHandler metrics.Handler,
payloadSerializer serialization.Serializer,
timeSource clock.TimeSource,
namespaceRegistry namespace.Registry,
saMapperProvider searchattribute.MapperProvider,
saProvider searchattribute.Provider,
saValidator *searchattribute.Validator,
clusterMetadata cluster.Metadata,
archivalMetadata archiver.ArchivalMetadata,
healthServer *health.Server,
membershipMonitor membership.Monitor,
healthInterceptor *interceptor.HealthInterceptor,
scheduleSpecBuilder *scheduler.SpecBuilder,
activityHandler activity.FrontendHandler,
callbackValidator callback.Validator,
nexusOperationHandler chasmnexus.FrontendHandler,
registry *chasm.Registry,
frontendServiceResolver membership.ServiceResolver,
workerDeploymentReadRateLimiter := configs.NewGlobalNamespaceRateLimiter(
frontendServiceResolver,
serviceConfig.GlobalWorkerDeploymentReadRPS,
serviceConfig.GlobalWorkerDeploymentReadBurstRatio,
log.With(logger, tag.ComponentRPCHandler, tag.ScopeNamespace),
)
wfHandler := NewWorkflowHandler(
callbackValidator,
serviceConfig,
namespaceReplicationQueue,
visibilityMgr,
logger,
throttledLogger,
persistenceExecutionManager.GetName(),
clusterMetadataManager,
persistenceMetadataManager,
historyClient,
matchingClient,
workerDeploymentStoreClient,
schedulerClient,
archiverProvider,
payloadSerializer,
namespaceRegistry,
saMapperProvider,
saProvider,
saValidator,
clusterMetadata,
archivalMetadata,
healthServer,
timeSource,
membershipMonitor,
healthInterceptor,
scheduleSpecBuilder,
httpEnabled(cfg, serviceName),
activityHandler,
nexusOperationHandler,
registry,
workerDeploymentReadRateLimiter,
chasmworkflow.NewValidator(
chasmworkflow.NewConfig(dc),
saMapperProvider,
saValidator,
),
)
return wfHandler
}
func RegisterNexusOperationHTTPHandler(
h *NexusOperationHTTPHandler,
router *mux.Router,
h.RegisterRoutes(router)
}
func RegisterNexusCompletionHTTPHandler(
h *nexusCompletionHTTPHandler,
router *mux.Router,
h.RegisterRoutes(router)
}
func RegisterOpenAPIHTTPHandler(
rateLimitInterceptor *interceptor.RateLimitInterceptor,
logger log.Logger,
router *mux.Router,
h := NewOpenAPIHTTPHandler(
rateLimitInterceptor,
logger,
)
h.RegisterRoutes(router)
return h
}
// Instantiate a router to support additional route prefixes.
return mux.NewRouter().UseEncodedPath()
}
// If the service is not the frontend service, HTTP API is disabled
if serviceName != primitives.FrontendService && serviceName != primitives.InternalFrontendService {
return false
}
// If HTTP API port is 0, it is disabled
}
// HTTPAPIServerProvider provides an HTTP API server if enabled or nil
// otherwise.
func HTTPAPIServerProvider(
cfg *config.Config,
serviceName primitives.ServiceName,
serviceConfig *Config,
grpcListener net.Listener,
tlsConfigProvider encryption.TLSConfigProvider,
handler Handler,
operatorHandler *OperatorHandlerImpl,
grpcServerOptions GrpcServerOptions,
metricsHandler metrics.Handler,
namespaceRegistry namespace.Registry,
logger log.Logger,
router *mux.Router,
if !httpEnabled(cfg, serviceName) {
}
return NewHTTPAPIServer(
serviceConfig,
rpcConfig,
grpcListener,
tlsConfigProvider,
handler,
operatorHandler,
grpcServerOptions.UnaryInterceptors,
metricsHandler,
router,
namespaceRegistry,
logger,
)
}
func NexusEndpointClientProvider(
dc *dynamicconfig.Collection,
namespaceRegistry namespace.Registry,
matchingClient resource.MatchingClient,
nexusEndpointManager persistence.NexusEndpointManager,
logger log.Logger,
clientConfig := newNexusEndpointClientConfig(dc)
return newNexusEndpointClient(
clientConfig,
namespaceRegistry,
matchingClient,
nexusEndpointManager,
logger,
)
}
lc.Append(fx.StartStopHook(svc.Start, svc.Stop))
}