onebox.go ×75

Frontier kind: Joint frontier

unlabeled · c_a0575922eafe

1 test · 29268 LOC · 661 files · introduces 1 test · 1256 LOC · 37 files

Introduces — evidence that enters the hierarchy at this concept

Code
238 ranges1256 lines · 37 files
Tests
1 test

Contains — complete concept membership

All code (extent)
6162 ranges29268 lines · 661 files · Browse complete extent
All tests (intent)
1 testBrowse 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.

1 test introduced at this concept.

Introduced code

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

Showing the top 20 of 37 files by introduced lines: 1213 of 1256 introduced LOC and 218 of 238 ranges. Expand a file to inspect source; the > gutter marks introduced lines.

go.temporal.io/server/tests/testcore/onebox.go 406 introduced LOC · 75 ranges

Open complete file

179
180 // newTemporal returns an instance that hosts full temporal in one process
181 > func newTemporal(t *testing.T, params *temporalParams) *temporalImpl { onebox.go
182 > impl := &temporalImpl{
183 > logger: params.logger,
184 > clusterMetadataConfig: params.clusterMetadataConfig,
185 > persistenceConfig: params.persistenceConfig,
186 > metadataMgr: params.metadataMgr,
187 > clusterMetadataMgr: params.clusterMetadataManager,
188 > shardMgr: params.shardMgr,
189 > taskMgr: params.taskMgr,
190 > executionManager: params.executionManager,
191 > namespaceReplicationQueue: params.namespaceReplicationQueue,
192 > abstractDataStoreFactory: params.abstractDataStoreFactory,
193 > visibilityStoreFactory: params.visibilityStoreFactory,
194 > esConfig: params.esConfig,
195 > esClient: params.esClient,
196 > archiverMetadata: params.archiverMetadata,
197 > archiverProvider: params.archiverProvider,
198 > frontendConfig: params.frontendConfig,
199 > historyConfig: params.historyConfig,
200 > matchingConfig: params.matchingConfig,
201 > workerConfig: params.workerConfig,
202 > mockAdminClient: params.mockAdminClient,
203 > namespaceReplicationTaskExecutor: params.namespaceReplicationTaskExecutor,
204 > dcRedirectionPolicy: params.dcRedirectionPolicy,
205 > tlsConfigProvider: params.tlsConfigProvider,
206 > captureMetricsHandler: params.captureMetricsHandler,
207 > dcClient: dynamicconfig.NewMemoryClient(),
208 > testHooks: testhooks.NewTestHooks(),
209 > serviceFxOptions: params.serviceFxOptions,
210 > taskCategoryRegistry: params.taskCategoryRegistry,
211 > hostsByProtocolByService: params.hostsByProtocolByService,
212 > replicationStreamRecorder: NewReplicationStreamRecorder(),
213 > spanExporters: params.spanExporters,
214 > tokenProvider: params.tokenProvider,
215 > enableHistoryTaskRecorder: params.enableHistoryTaskRecorder,
216 > }
217 >
218 > // Configure output file path for on-demand logging (call WriteToLog() to write)
219 > clusterName := params.clusterMetadataConfig.CurrentClusterName
220 > outputFile := fmt.Sprintf("/tmp/replication_stream_messages_%s.txt", clusterName)
221 > impl.replicationStreamRecorder.SetOutputFile(outputFile)
222 > impl.clients = newClients(
223 > impl.logger,
224 > impl.hostsByProtocolByService[grpcProtocol],
225 > &impl.frontendMembershipAddress,
226 > impl.tlsConfigProvider,
227 > impl.GetMetricsHandler(),
228 > impl.dcClient,
229 > impl.testHooks,
230 > impl.historyConfig.NumHistoryShards,
231 > impl.metadataMgr,
232 > impl.tokenProvider,
233 > )
234 >
235 > // Global defaults: applied without cleanup so they persist across cluster reuse.
236 > for k, v := range defaultDynamicConfigOverrides {
237 > impl.overrideDynamicConfigForClusterLifetime(k, v)
238 > }
239 // Override Nexus callback URL. This is parameterized on the frontend's HTTP address,
240 // so it can't be overriden in the loop above.
241 > impl.setNexusCallbackURL() onebox.go
242 > // Per-test overrides: cleaned up when the creating test finishes.
243 > for k, v := range params.dynamicConfigOverrides {
244 impl.overrideDynamicConfigForTest(t, k, v)
245 }
246 > return impl onebox.go
247 }
248
249 > func (c *temporalImpl) Start() error { onebox.go
250 > // create temporal-system namespace, this must be created before starting
251 > // the services - so directly use the metadataManager to create this
252 > if err := c.createSystemNamespace(); err != nil {
253 return err
254 }
255 > c.startMatching() onebox.go
256 > c.startHistory()
257 > c.startFrontend()
258 > c.startWorker()
259 >
260 > return nil
261 }
262
263 > func (c *temporalImpl) Stop() error { onebox.go
264 > ctx, cancel := context.WithTimeout(context.Background(), 15*time.Second)
265 > defer cancel()
266 >
267 > var errs []error
268 > errs = append(errs, c.close()...)
269 >
270 > slices.Reverse(c.fxApps) // less log spam if we go backwards
271 > for _, app := range c.fxApps {
272 > errs = append(errs, app.Stop(ctx))
273 > }
274
275 > return multierr.Combine(errs...) onebox.go
276 }
277
278 > func (c *temporalImpl) makeHostMap(serviceName primitives.ServiceName, self string) map[primitives.ServiceName]static.Hosts { onebox.go
279 > hostMap := maps.Clone(c.hostsByProtocolByService[grpcProtocol])
280 > hosts := hostMap[serviceName]
281 > hosts.Self = self
282 > hostMap[serviceName] = hosts
283 > return hostMap
284 > }
285
286 // Use this to get an address for a remote cluster to connect to.
289 }
290
291 > func (c *temporalImpl) FrontendHTTPAddress() string { onebox.go
292 > // randomize like a load balancer would
293 > addrs := c.hostsByProtocolByService[httpProtocol][primitives.FrontendService].All
294 > return addrs[rand.Intn(len(addrs))]
295 > }
296
297 > func (c *temporalImpl) FrontendGRPCAddress() string { onebox.go
298 > return c.hostsByProtocolByService[grpcProtocol][primitives.FrontendService].All[0]
299 > }
300
301 > func (c *temporalImpl) WorkerGRPCAddress() string { onebox.go
302 > return c.hostsByProtocolByService[grpcProtocol][primitives.WorkerService].All[0]
303 > }
304
305 func (c *temporalImpl) DcClient() *dynamicconfig.MemoryClient {
319 }
320
321 > func (c *temporalImpl) copyPersistenceConfig() config.Persistence { onebox.go
322 > persistenceConfig := copyPersistenceConfig(c.persistenceConfig)
323 > if c.esConfig != nil {
324 esDataStoreName := "es-visibility"
325 persistenceConfig.VisibilityStore = esDataStoreName
328 }
329 }
330 > return persistenceConfig onebox.go
331 }
332
333 > func (c *temporalImpl) startFrontend() { onebox.go
334 > serviceName := primitives.FrontendService
335 >
336 > var grpcResolver *membership.GRPCResolver
337 >
338 > for _, host := range c.hostsByProtocolByService[grpcProtocol][serviceName].All {
339 > logger := log.With(c.logger, tag.Host(host))
340 > app := fx.New(
341 > fx.Supply(
342 > c.copyPersistenceConfig(),
343 > serviceName,
344 > c.mockAdminClient,
345 > ),
346 > fx.Provide(c.frontendConfigProvider),
347 > fx.Provide(func() listenHostPort { return listenHostPort(host) }),
348 > fx.Provide(func() httpPort { return mustPortFromAddress(c.FrontendHTTPAddress()) }),
349 > fx.Provide(func() config.DCRedirectionPolicy { return c.dcRedirectionPolicy }),
350 > fx.Provide(func() log.Logger { return logger }),
351 > fx.Provide(func() log.ThrottledLogger { return logger }),
352 > fx.Provide(func() resource.NamespaceLogger { return logger }),
353 fx.Provide(c.newRPCFactory),
354 static.MembershipModule(c.makeHostMap(serviceName, host)),
355 > fx.Provide(func() *cluster.Config { return c.clusterMetadataConfig }), onebox.go
356 > fx.Provide(func() carchiver.ArchivalMetadata { return c.archiverMetadata }),
357 > fx.Provide(func() provider.ArchiverProvider { return c.archiverProvider }),
358 fx.Provide(sdkClientFactoryProvider),
359 fx.Provide(c.GetMetricsHandler),
360 > fx.Provide(func() []grpc.UnaryServerInterceptor { onebox.go
361 > if c.replicationStreamRecorder != nil {
362 > return []grpc.UnaryServerInterceptor{
363 > c.replicationStreamRecorder.UnaryServerInterceptor(c.clusterMetadataConfig.CurrentClusterName),
364 > }
365 > }
366 return nil
367 }),
368 > fx.Provide(func() []grpc.StreamServerInterceptor { onebox.go
369 > if c.replicationStreamRecorder != nil {
370 > return []grpc.StreamServerInterceptor{
371 > c.replicationStreamRecorder.StreamServerInterceptor(c.clusterMetadataConfig.CurrentClusterName),
372 > }
373 > }
374 return nil
375 }),
376 > fx.Provide(func() authorization.Authorizer { return c }), onebox.go
377 > fx.Provide(func() authorization.ClaimMapper { return c }),
378 > fx.Provide(func() authorization.JWTAudienceMapper { return nil }),
379 fx.Provide(newClientFactoryProvider),
380 > fx.Provide(func() searchattribute.Mapper { return nil }), onebox.go
381 // Comment the line above and uncomment the line below to test with search attributes mapper.
382 // fx.Provide(func() searchattribute.Mapper { return NewSearchAttributeTestMapper() }),
383 > fx.Provide(func() resolver.ServiceResolver { return resolver.NewNoopResolver() }), onebox.go
384 fx.Provide(persistenceClient.FactoryProvider),
385 > fx.Provide(func() persistenceClient.AbstractDataStoreFactory { return c.abstractDataStoreFactory }), onebox.go
386 > fx.Provide(func() visibility.VisibilityStoreFactory { return c.visibilityStoreFactory }),
387 > fx.Provide(func() dynamicconfig.Client { return c.dcClient }),
388 > fx.Decorate(func() testhooks.TestHooks { return c.testHooks }),
389 fx.Provide(resource.DefaultSnTaggedLoggerProvider),
390 fx.Provide(func() esclient.Client { return c.esClient }),
399 chasm.Module,
400 )
401 > err := app.Err() onebox.go
402 > if err != nil {
403 logger.Fatal("unable to construct frontend service", tag.Error(err))
404 }
405
406 > c.fxApps = append(c.fxApps, app) onebox.go
407 >
408 > if err := app.Start(context.Background()); err != nil {
409 logger.Fatal("unable to start frontend service", tag.Error(err))
410 }
412
413 // Address for SDKs
414 > c.frontendMembershipAddress = grpcResolver.MakeURL(serviceName) onebox.go
415 }
416
417 > func (c *temporalImpl) startHistory() { onebox.go
418 > serviceName := primitives.HistoryService
419 >
420 > testhooks.NewHook(testhooks.HistoryChasmRuntimeProvider, func(
421 > chasmEngine chasm.Engine,
422 > chasmVisibilityManager chasm.VisibilityManager,
423 > _ *chasm.Registry,
424 > ) {
425 > c.chasmEngine = chasmEngine
426 > c.chasmVisibilityMgr = chasmVisibilityManager
427 > }).Apply(c.testHooks, testhooks.GlobalScope)
428
429 > persistenceFactoryProvider := persistenceClient.FactoryProvider onebox.go
430 > if c.enableHistoryTaskRecorder {
431 persistenceFactoryProvider = func(params persistenceClient.NewFactoryParams) persistenceClient.Factory {
432 return &historyTaskRecordingPersistenceFactory{
440 }
441
442 > for _, host := range c.hostsByProtocolByService[grpcProtocol][serviceName].All { onebox.go
443 > var namespaceRegistry namespace.Registry
444 > logger := log.With(c.logger, tag.Host(host))
445 > app := fx.New(
446 > fx.Supply(
447 > c.copyPersistenceConfig(),
448 > serviceName,
449 > c.mockAdminClient,
450 > ),
451 > fx.Provide(c.configProvider),
452 > fx.Provide(c.GetMetricsHandler),
453 > fx.Provide(func() otellog.Logger { return wideevents.NoopLogger() }),
454 > fx.Provide(func() listenHostPort { return listenHostPort(host) }),
455 > fx.Provide(func() httpPort { return mustPortFromAddress(c.FrontendHTTPAddress()) }),
456 fx.Provide(func() config.DCRedirectionPolicy { return config.DCRedirectionPolicy{} }),
457 > fx.Provide(func() log.Logger { return logger }), onebox.go
458 > fx.Provide(func() log.ThrottledLogger { return logger }),
459 fx.Provide(c.newRPCFactory),
460 > fx.Decorate(func(base []grpc.UnaryServerInterceptor) []grpc.UnaryServerInterceptor { onebox.go
461 > if c.replicationStreamRecorder != nil {
462 > return append(base, c.replicationStreamRecorder.UnaryServerInterceptor(c.clusterMetadataConfig.CurrentClusterName))
463 > }
464 return base
465 }),
466 > fx.Provide(func() []grpc.StreamServerInterceptor { onebox.go
467 > if c.replicationStreamRecorder != nil {
468 > return []grpc.StreamServerInterceptor{
469 > c.replicationStreamRecorder.StreamServerInterceptor(c.clusterMetadataConfig.CurrentClusterName),
470 > }
471 > }
472 return nil
473 }),
474 static.MembershipModule(c.makeHostMap(serviceName, host)),
475 > fx.Provide(func() *cluster.Config { return c.clusterMetadataConfig }), onebox.go
476 > fx.Provide(func() carchiver.ArchivalMetadata { return c.archiverMetadata }),
477 > fx.Provide(func() provider.ArchiverProvider { return c.archiverProvider }),
478 fx.Provide(sdkClientFactoryProvider),
479 fx.Provide(newClientFactoryProvider),
480 > fx.Provide(func() searchattribute.Mapper { return nil }), onebox.go
481 // Comment the line above and uncomment the line below to test with search attributes mapper.
482 // fx.Provide(func() searchattribute.Mapper { return NewSearchAttributeTestMapper() }),
483 > fx.Provide(func() resolver.ServiceResolver { return resolver.NewNoopResolver() }), onebox.go
484 fx.Provide(persistenceFactoryProvider),
485 > fx.Provide(func() persistenceClient.AbstractDataStoreFactory { return c.abstractDataStoreFactory }), onebox.go
486 > fx.Provide(func() visibility.VisibilityStoreFactory { return c.visibilityStoreFactory }),
487 > fx.Provide(func() dynamicconfig.Client { return c.dcClient }),
488 > fx.Decorate(func() testhooks.TestHooks { return c.testHooks }),
489 fx.Provide(resource.DefaultSnTaggedLoggerProvider),
490 fx.Provide(func() esclient.Client { return c.esClient }),
501 fx.Populate(&namespaceRegistry),
502 )
503 > err := app.Err() onebox.go
504 > if err != nil {
505 logger.Fatal("unable to construct history service", tag.Error(err))
506 }
507 > c.fxApps = append(c.fxApps, app) onebox.go
508 >
509 > if err := app.Start(context.Background()); err != nil {
510 logger.Fatal("unable to start history service", tag.Error(err))
511 }
513 }
514
515 > func (c *temporalImpl) startMatching() { onebox.go
516 > serviceName := primitives.MatchingService
517 >
518 > for _, host := range c.hostsByProtocolByService[grpcProtocol][serviceName].All {
519 > var namespaceRegistry namespace.Registry
520 > logger := log.With(c.logger, tag.Host(host))
521 > app := fx.New(
522 > fx.Supply(
523 > c.copyPersistenceConfig(),
524 > serviceName,
525 > c.mockAdminClient,
526 > ),
527 > fx.Provide(c.configProvider),
528 > fx.Provide(c.GetMetricsHandler),
529 > fx.Provide(func() listenHostPort { return listenHostPort(host) }),
530 > fx.Provide(func() httpPort { return mustPortFromAddress(c.FrontendHTTPAddress()) }),
531 > fx.Provide(func() log.Logger { return logger }),
532 > fx.Provide(func() log.ThrottledLogger { return logger }),
533 fx.Provide(c.newRPCFactory),
534 static.MembershipModule(c.makeHostMap(serviceName, host)),
535 > fx.Provide(func() *cluster.Config { return c.clusterMetadataConfig }), onebox.go
536 fx.Provide(func() carchiver.ArchivalMetadata { return c.archiverMetadata }),
537 fx.Provide(func() provider.ArchiverProvider { return c.archiverProvider }),
538 fx.Provide(newClientFactoryProvider),
539 > fx.Provide(func() searchattribute.Mapper { return nil }), onebox.go
540 > fx.Provide(func() resolver.ServiceResolver { return resolver.NewNoopResolver() }),
541 fx.Provide(persistenceClient.FactoryProvider),
542 > fx.Provide(func() persistenceClient.AbstractDataStoreFactory { return c.abstractDataStoreFactory }), onebox.go
543 > fx.Provide(func() visibility.VisibilityStoreFactory { return c.visibilityStoreFactory }),
544 > fx.Provide(func() dynamicconfig.Client { return c.dcClient }),
545 > fx.Decorate(func() testhooks.TestHooks { return c.testHooks }),
546 fx.Provide(func() esclient.Client { return c.esClient }),
547 fx.Provide(c.GetTLSConfigProvider),
556 fx.Populate(&namespaceRegistry),
557 )
558 > err := app.Err() onebox.go
559 > if err != nil {
560 logger.Fatal("unable to start matching service", tag.Error(err))
561 }
562 > c.fxApps = append(c.fxApps, app) onebox.go
563 > if err := app.Start(context.Background()); err != nil {
564 logger.Fatal("unable to start matching service", tag.Error(err))
565 }
567 }
568
569 > func (c *temporalImpl) startWorker() { onebox.go
570 > serviceName := primitives.WorkerService
571 >
572 > clusterConfigCopy := cluster.Config{
573 > EnableGlobalNamespace: c.clusterMetadataConfig.EnableGlobalNamespace,
574 > FailoverVersionIncrement: c.clusterMetadataConfig.FailoverVersionIncrement,
575 > MasterClusterName: c.clusterMetadataConfig.MasterClusterName,
576 > CurrentClusterName: c.clusterMetadataConfig.CurrentClusterName,
577 > ClusterInformation: maps.Clone(c.clusterMetadataConfig.ClusterInformation),
578 > }
579 >
580 > for _, host := range c.hostsByProtocolByService[grpcProtocol][serviceName].All {
581 > var namespaceRegistry namespace.Registry
582 > logger := log.With(c.logger, tag.Host(host))
583 > app := fx.New(
584 >
585 > fx.Supply(
586 > c.copyPersistenceConfig(),
587 > serviceName,
588 > c.mockAdminClient,
589 > ),
590 > fx.Provide(c.configProvider),
591 > fx.Provide(c.GetMetricsHandler),
592 > fx.Provide(func() listenHostPort { return listenHostPort(host) }),
593 > fx.Provide(func() httpPort { return mustPortFromAddress(c.FrontendHTTPAddress()) }),
594 fx.Provide(func() config.DCRedirectionPolicy { return config.DCRedirectionPolicy{} }),
595 > fx.Provide(func() log.Logger { return logger }), onebox.go
596 > fx.Provide(func() log.ThrottledLogger { return logger }),
597 fx.Provide(c.newRPCFactory),
598 static.MembershipModule(c.makeHostMap(serviceName, host)),
599 > fx.Provide(func() *cluster.Config { return &clusterConfigCopy }), onebox.go
600 fx.Provide(func() carchiver.ArchivalMetadata { return c.archiverMetadata }),
601 fx.Provide(func() provider.ArchiverProvider { return c.archiverProvider }),
602 fx.Provide(sdkClientFactoryProvider),
603 fx.Provide(newClientFactoryProvider),
604 > fx.Provide(func() searchattribute.Mapper { return nil }), onebox.go
605 > fx.Provide(func() resolver.ServiceResolver { return resolver.NewNoopResolver() }),
606 fx.Provide(persistenceClient.FactoryProvider),
607 > fx.Provide(func() persistenceClient.AbstractDataStoreFactory { return c.abstractDataStoreFactory }), onebox.go
608 > fx.Provide(func() visibility.VisibilityStoreFactory { return c.visibilityStoreFactory }),
609 > fx.Provide(func() dynamicconfig.Client { return c.dcClient }),
610 > fx.Decorate(func() testhooks.TestHooks { return c.testHooks }),
611 fx.Provide(resource.DefaultSnTaggedLoggerProvider),
612 fx.Provide(func() esclient.Client { return c.esClient }),
621 fx.Populate(&namespaceRegistry),
622 )
623 > err := app.Err() onebox.go
624 > if err != nil {
625 logger.Fatal("unable to start worker service", tag.Error(err))
626 }
627
628 > c.fxApps = append(c.fxApps, app) onebox.go
629 > if err := app.Start(context.Background()); err != nil {
630 logger.Fatal("unable to start worker service", tag.Error(err))
631 }
633 }
634
635 > func (c *temporalImpl) getFxOptionsForService(serviceName primitives.ServiceName) fx.Option { onebox.go
636 > return fx.Options(c.serviceFxOptions[serviceName]...)
637 > }
638
639 > func (c *temporalImpl) createSystemNamespace() error { onebox.go
640 > err := c.metadataMgr.InitializeSystemNamespaces(context.Background(), c.clusterMetadataConfig.CurrentClusterName)
641 > if err != nil {
642 return fmt.Errorf("failed to create temporal-system namespace: %v", err)
643 }
644 > return nil onebox.go
645 }
646
649 }
650
651 > func (c *temporalImpl) GetTLSConfigProvider() encryption.TLSConfigProvider { onebox.go
652 > // If we just return this directly, the interface will be non-nil but the
653 > // pointer will be nil
654 > if c.tlsConfigProvider != nil {
655 return c.tlsConfigProvider
656 }
657 > return nil onebox.go
658 }
659
660 > func (c *temporalImpl) GetTaskCategoryRegistry() tasks.TaskCategoryRegistry { onebox.go
661 > return c.taskCategoryRegistry
662 > }
663
664 func (c *temporalImpl) TLSConfigProvider() *encryption.FixedTLSConfigProvider {
672 }
673
674 > func (c *temporalImpl) GetMetricsHandler() metrics.Handler { onebox.go
675 > if c.captureMetricsHandler != nil {
676 > return c.captureMetricsHandler
677 > }
678 return metrics.NoopMetricsHandler
679 }
680
681 > func (c *temporalImpl) frontendConfigProvider() *config.Config { onebox.go
682 > // Set HTTP port and a test HTTP forwarded header
683 > return &config.Config{
684 > Services: map[string]config.Service{
685 > string(primitives.FrontendService): {
686 > RPC: config.RPC{
687 > HTTPPort: int(mustPortFromAddress(c.FrontendHTTPAddress())),
688 > HTTPAdditionalForwardedHeaders: []string{
689 > "this-header-forwarded",
690 > "this-header-prefix-forwarded-*",
691 > },
692 > },
693 > },
694 > },
695 > DCRedirectionPolicy: c.dcRedirectionPolicy,
696 > ExporterConfig: telemetry.ExportConfig{
697 > CustomExporters: c.spanExporters,
698 > },
699 > }
700 > }
701
702 > func (c *temporalImpl) configProvider(serviceName primitives.ServiceName) *config.Config { onebox.go
703 > return &config.Config{
704 > Services: map[string]config.Service{
705 > string(serviceName): {
706 > RPC: config.RPC{},
707 > },
708 > },
709 > DCRedirectionPolicy: config.DCRedirectionPolicy{},
710 > ExporterConfig: telemetry.ExportConfig{
711 > CustomExporters: c.spanExporters,
712 > },
713 > }
714 > }
715
716 func (c *temporalImpl) newRPCFactory(
724 httpPort httpPort,
725 metricsHandler metrics.Handler,
726 > ) (common.RPCFactory, error) { onebox.go
727 > host, portStr, err := net.SplitHostPort(string(grpcHostPort))
728 > if err != nil {
729 return nil, fmt.Errorf("failed parsing host:port: %w", err)
730 }
731 > port, err := strconv.Atoi(portStr) onebox.go
732 > if err != nil {
733 return nil, fmt.Errorf("invalid port: %w", err)
734 }
735 > var frontendTLSConfig *tls.Config onebox.go
736 > if tlsConfigProvider != nil {
737 if frontendTLSConfig, err = tlsConfigProvider.GetFrontendClientConfig(); err != nil {
738 return nil, fmt.Errorf("failed getting client TLS config: %w", err)
739 }
740 }
741 > var options []grpc.DialOption onebox.go
742 > if tracingStatsHandler != nil {
743 options = append(options, grpc.WithStatsHandler(tracingStatsHandler))
744 }
745 // Add replication stream recorder injector
746 > if c.replicationStreamRecorder != nil { onebox.go
747 > options = append(options,
748 > grpc.WithChainUnaryInterceptor(c.replicationStreamRecorder.UnaryInterceptor(c.clusterMetadataConfig.CurrentClusterName)),
749 > grpc.WithChainStreamInterceptor(c.replicationStreamRecorder.StreamInterceptor(c.clusterMetadataConfig.CurrentClusterName)),
750 > )
751 > }
752 > rpcConfig := config.RPC{BindOnIP: host, GRPCPort: port, HTTPPort: int(httpPort)}
753 > cfg := &config.Config{
754 > Services: map[string]config.Service{
755 > string(sn): {
756 > RPC: rpcConfig,
757 > },
758 > },
759 > }
760 > return rpc.NewFactory(
761 > cfg,
762 > sn,
763 > logger,
764 > metricsHandler,
765 > tlsConfigProvider,
766 > grpcResolver.MakeURL(primitives.FrontendService),
767 > grpcResolver.MakeURL(primitives.FrontendService),
768 > int(httpPort),
769 > frontendTLSConfig,
770 > options,
771 > resource.PerServiceDialOptionsProvider(logger),
772 > monitor,
773 > c.tokenProvider,
774 > ), nil
775 }
776
803 caller *authorization.Claims,
804 target *authorization.CallTarget,
805 > ) (authorization.Result, error) { onebox.go
806 > c.callbackLock.RLock()
807 > onAuthorize := c.onAuthorize
808 > c.callbackLock.RUnlock()
809 > if onAuthorize != nil {
810 return onAuthorize(ctx, caller, target)
811 }
812 > return authorization.Result{Decision: authorization.DecisionAllow}, nil onebox.go
813 }
814
817 // The race condition happens because all the services are using the same datastore map in the config.
818 // Also all services will retry to modify the maxQPS field in the datastore during start up and use the modified maxQPS value to create a persistence factory.
819 > func copyPersistenceConfig(cfg config.Persistence) config.Persistence { onebox.go
820 > var newCfg config.Persistence
821 > b, err := json.Marshal(cfg)
822 > if err != nil {
823 panic("copy persistence config: " + err.Error())
824 > } else if err = json.Unmarshal(b, &newCfg); err != nil { onebox.go
825 panic("copy persistence config: " + err.Error())
826 }
827
828 // Preserve fault injection injectors after the JSON copy.
829 > for name, dataStore := range cfg.DataStores { onebox.go
830 > if dataStore.FaultInjection == nil {
831 > continue
832 }
833 newDataStore := newCfg.DataStores[name]
839 }
840
841 > return newCfg onebox.go
842 }
843
844 > func (c *temporalImpl) setNexusCallbackURL() { onebox.go
845 > // Set Nexus callback URL with the cluster's HTTP address. This is a sensible default to avoid
846 > // users to need to manually set this.
847 > //nolint:revive // test callback endpoints are served by the local HTTP API in functional tests
848 > nexusCallbackTemplate := fmt.Sprintf(
849 > "http://%s/namespaces/{{.NamespaceName}}/nexus/callback",
850 > c.FrontendHTTPAddress(),
851 > )
852 > c.overrideDynamicConfigForClusterLifetime(nexusoperations.CallbackURLTemplate.Key(), nexusCallbackTemplate)
853 > c.overrideDynamicConfigForClusterLifetime(chasmnexus.CallbackURLTemplate.Key(), nexusCallbackTemplate)
854 > }
855
856 > func (c *temporalImpl) overrideDynamicConfigForClusterLifetime(name dynamicconfig.Key, value any) { onebox.go
857 > c.dcClient.PartialOverrideValue(name, value)
858 > }
859
860 // overrideDynamicConfigForTest overrides a dynamic config value for the duration of the test.
861 > func (c *temporalImpl) overrideDynamicConfigForTest(t *testing.T, name dynamicconfig.Key, value any) func() { onebox.go
862 > cleanup := c.dcClient.PartialOverrideValue(name, value)
863 > t.Cleanup(cleanup)
864 > return cleanup
865 > }
866
867 func (c *temporalImpl) injectHook(t *testing.T, hook testhooks.Hook, scope any) func() {
871 }
872
873 > func mustPortFromAddress(addr string) httpPort { onebox.go
874 > _, port, err := net.SplitHostPort(addr)
875 > if err != nil {
876 panic(fmt.Errorf("Invalid address: %w", err))
877 }
878 > portInt, err := strconv.Atoi(port) onebox.go
879 > if err != nil {
880 panic(fmt.Errorf("Cannot parse port: %w", err))
881 }
882 > return httpPort(portInt) onebox.go
883 }
go.temporal.io/server/tests/testcore/functional_test_base.go 236 introduced LOC · 34 ranges

Open complete file

193 }
194
195 > func (s *FunctionalTestBase) GetTestCluster() *TestCluster { functional_test_base.go
196 > return s.testCluster
197 > }
198
199 func (s *FunctionalTestBase) GetTestClusterConfig() *TestClusterConfig {
201 }
202
203 > func (s *FunctionalTestBase) FrontendClient() workflowservice.WorkflowServiceClient { functional_test_base.go
204 > return s.testCluster.FrontendClient()
205 > }
206
207 func (s *FunctionalTestBase) AdminClient() adminservice.AdminServiceClient {
209 }
210
211 > func (s *FunctionalTestBase) OperatorClient() operatorservice.OperatorServiceClient { functional_test_base.go
212 > return s.testCluster.OperatorClient()
213 > }
214
215 > func (s *FunctionalTestBase) HttpAPIAddress() string { functional_test_base.go
216 > return s.testCluster.Host().FrontendHTTPAddress()
217 > }
218
219 > func (s *FunctionalTestBase) Namespace() namespace.Name { functional_test_base.go
220 > return s.namespace
221 > }
222
223 func (s *FunctionalTestBase) NamespaceID() namespace.ID {
225 }
226
227 > func (s *FunctionalTestBase) ExternalNamespace() namespace.Name { functional_test_base.go
228 > return s.externalNamespace
229 > }
230
231 > func (s *FunctionalTestBase) FrontendGRPCAddress() string { functional_test_base.go
232 > return s.GetTestCluster().Host().FrontendGRPCAddress()
233 > }
234
235 > func (s *FunctionalTestBase) WorkerGRPCAddress() string { functional_test_base.go
236 > return s.GetTestCluster().WorkerGRPCAddress()
237 > }
238
239 func (s *FunctionalTestBase) SdkWorker() sdkworker.Worker {
268 }
269
270 > func (s *FunctionalTestBase) SetupSuiteWithCluster(options ...TestClusterOption) { functional_test_base.go
271 > // Reserve a slot from the dedicated test cluster pool.
272 > testClusterRouter.dedicated.reserveSlot(s.T())
273 > s.setupCluster(options...)
274 > clusterRequest{
275 > kind: clusterKindDedicated,
276 > dedicatedReason: "legacy-suite",
277 > needWorkerService: ApplyTestClusterOptions(options).EnableWorkerService,
278 > }.recordCreation(s.T())
279 > }
280
281 > func (s *FunctionalTestBase) setupCluster(options ...TestClusterOption) { functional_test_base.go
282 > params := ApplyTestClusterOptions(options)
283 >
284 > // A custom logger supplied via WithClusterLogger takes precedence.
285 > if params.Logger != nil {
286 s.Logger = params.Logger
287 }
288
289 // NOTE: A suite might set its own logger. Example: AcquireShardSuiteBase.
290 > if s.Logger == nil { functional_test_base.go
291 > // The cluster outlives any single test and s.T() changes as different tests use it,
292 > // but the proxy T's Name() must be stable, so this is never updated.
293 > s.t = &sharedClusterT{
294 > name: s.T().Name(),
295 > logFanout: os.Getenv("CI") != "",
296 > }
297 > tl := testlogger.NewTestLogger(s.t, testlogger.FailOnExpectedErrorOnly)
298 > // Fail tests when a soft assertion fires (see `softassert` package).
299 > tl.Expect(testlogger.Error, ".*", tag.FailedAssertion)
300 > s.Logger = tl
301 > }
302
303 > s.testClusterConfig = &TestClusterConfig{ functional_test_base.go
304 > FaultInjection: params.FaultInjectionConfig,
305 > HistoryConfig: HistoryConfig{
306 > NumHistoryShards: cmp.Or(params.NumHistoryShards, 4),
307 > },
308 > DCRedirectionPolicy: params.DCRedirectionPolicy,
309 > DynamicConfigOverrides: params.DynamicConfigOverrides,
310 > EnableMetricsCapture: true,
311 > EnableArchival: params.ArchivalEnabled,
312 > EnableMTLS: params.EnableMTLS,
313 > EnableHistoryTaskRecorder: params.EnableHistoryTaskRecorder,
314 > CustomHistoryArchiverFactory: params.CustomHistoryArchiverFactory,
315 > CustomVisibilityArchiverFactory: params.CustomVisibilityArchiverFactory,
316 > WorkerConfig: WorkerConfig{DisableWorker: !params.EnableWorkerService},
317 > }
318 >
319 > // Apply configuration for shared clusters.
320 > if params.SharedCluster {
321 s.testClusterConfig.Persistence = sharedClusterPersistence(GetPersistenceTestDefaults())
322 s.isShared = true
325 // Initialize the OTEL collector if OTEL is enabled.
326 // Must be done before the test cluster is created, so that the collector can be used by the test cluster.
327 > if otelOutputDir := os.Getenv("TEMPORAL_TEST_OTEL_OUTPUT"); otelOutputDir != "" { functional_test_base.go
328 // Create an OTEL exporter.
329 s.otelExporter = testtelemetry.NewFileExporter(otelOutputDir)
335 }
336
337 > var err error functional_test_base.go
338 > testClusterFactory := NewTestClusterFactory()
339 > s.testCluster, err = testClusterFactory.NewCluster(s.T(), s.testClusterConfig, s.Logger)
340 > s.Require().NoError(err)
341 >
342 > // Setup test cluster namespaces.
343 > s.namespace = namespace.Name(RandomizeStr("namespace"))
344 > s.namespaceID, err = s.RegisterNamespace(s.Namespace(), 1, enumspb.ARCHIVAL_STATE_DISABLED, "", "")
345 > s.Require().NoError(err)
346 >
347 > s.externalNamespace = namespace.Name(RandomizeStr("external-namespace"))
348 > _, err = s.RegisterNamespace(s.ExternalNamespace(), 1, enumspb.ARCHIVAL_STATE_DISABLED, "", "")
349 > s.Require().NoError(err)
350 }
351
362 // into partitions. Otherwise, the test suite will be executed multiple times
363 // in each partition.
364 > func (s *FunctionalTestBase) SetupTest() { functional_test_base.go
365 > s.checkTestShard()
366 > s.initAssertions()
367 > s.setupSdk()
368 > s.taskPoller = taskpoller.New(s.T(), s.FrontendClient(), s.Namespace().String())
369 > }
370
371 func (s *FunctionalTestBase) SetupSubTest() {
374
375 // TODO: remove once `parallelsuite` and testEnv is rolled out everywhere
376 > func (s *FunctionalTestBase) initAssertions() { functional_test_base.go
377 > // `s.Assertions` (as well as other test helpers which depends on `s.T()`) must be initialized on
378 > // both test and subtest levels (but not suite level, where `s.T()` is `nil`).
379 > //
380 > // If these helpers are not reinitialized on subtest level, any failed `assert` in
381 > // subtest will fail the entire test (not subtest) immediately without running other subtests.
382 >
383 > s.Assertions = require.New(s.T())
384 > s.ProtoAssertions = protorequire.New(s.T())
385 > s.HistoryRequire = historyrequire.New(s.T())
386 > s.UpdateUtils = updateutils.New(s.T())
387 > }
388
389 // checkTestShard supports test sharding based on environment variables.
390 > func (s *FunctionalTestBase) checkTestShard() { functional_test_base.go
391 > checkTestShard(s.T())
392 > }
393
394 > func ApplyTestClusterOptions(options []TestClusterOption) testClusterParams { functional_test_base.go
395 > params := testClusterParams{
396 > EnableWorkerService: true,
397 > }
398 > for _, opt := range options {
399 opt(&params)
400 }
401 > return params functional_test_base.go
402 }
403
404 > func (s *FunctionalTestBase) setupSdk() { functional_test_base.go
405 > // Set URL template after httpAPAddress is set, see commonnexus.RouteCompletionCallback
406 > s.OverrideDynamicConfig(
407 > nexusoperations.CallbackURLTemplate,
408 > "http://"+s.HttpAPIAddress()+"/namespaces/{{.NamespaceName}}/nexus/callback")
409 >
410 > clientOptions := sdkclient.Options{
411 > HostPort: s.FrontendGRPCAddress(),
412 > Namespace: s.Namespace().String(),
413 > Logger: log.NewSdkLogger(s.Logger),
414 > }
415 >
416 > if provider := s.testCluster.host.tlsConfigProvider; provider != nil {
417 clientOptions.ConnectionOptions.TLS = provider.FrontendClientConfig
418 }
419
420 > var err error functional_test_base.go
421 > s.sdkClient, err = sdkclient.Dial(clientOptions)
422 > s.NoError(err)
423 > // TODO(alex): move initialization to suite level?
424 > s.taskQueue = RandomizeStr("tq")
425 >
426 > workerOptions := sdkworker.Options{}
427 > s.sdkWorker = sdkworker.New(s.sdkClient, s.taskQueue, workerOptions)
428 > err = s.sdkWorker.Start()
429 > s.NoError(err)
430 }
431
432 > func (s *FunctionalTestBase) exportOTELTraces() { functional_test_base.go
433 > if s.otelExporter == nil {
434 > return
435 > }
436 if s.T().Failed() {
437 var validFilenameChars = regexp.MustCompile(`[^a-zA-Z0-9._-]+`)
448 }
449
450 > func (s *FunctionalTestBase) TearDownCluster() { functional_test_base.go
451 > s.Require().NoError(s.MarkNamespaceAsDeleted(s.Namespace()))
452 > s.Require().NoError(s.MarkNamespaceAsDeleted(s.ExternalNamespace()))
453 > s.Require().NoError(s.tearDownTestCluster())
454 > }
455
456 // tearDownTestCluster tears down the underlying TestCluster and runs the proxy
466 return nil
467 }
468 > err := s.testCluster.TearDownCluster() functional_test_base.go
469 > s.testCluster = nil
470 > return err
471 }
472
473 // **IMPORTANT**: When overridding this, make sure to invoke `s.FunctionalTestBase.TearDownTest()`.
474 > func (s *FunctionalTestBase) TearDownTest() { functional_test_base.go
475 > s.exportOTELTraces()
476 > s.tearDownSdk()
477 > }
478
479 // **IMPORTANT**: When overridding this, make sure to invoke `s.FunctionalTestBase.TearDownSubTest()`.
482 }
483
484 > func (s *FunctionalTestBase) tearDownSdk() { functional_test_base.go
485 > if s.sdkWorker != nil {
486 > s.sdkWorker.Stop()
487 > }
488 > if s.sdkClient != nil {
489 > s.sdkClient.Close()
490 > }
491 }
492
501 historyArchivalURI string,
502 visibilityArchivalURI string,
503 > ) (namespace.ID, error) { functional_test_base.go
504 > currentClusterName := s.testCluster.testBase.ClusterMetadata.GetCurrentClusterName()
505 > nsID := namespace.ID(uuid.NewString())
506 > expectedSearchAttributes := searchattribute.TestSearchAttributesToRegister()
507 > namespaceRequest := &persistence.CreateNamespaceRequest{
508 > Namespace: &persistencespb.NamespaceDetail{
509 > Info: &persistencespb.NamespaceInfo{
510 > Id: nsID.String(),
511 > Name: nsName.String(),
512 > State: enumspb.NAMESPACE_STATE_REGISTERED,
513 > Description: "namespace for functional tests",
514 > },
515 > Config: &persistencespb.NamespaceConfig{
516 > Retention: timestamp.DurationFromDays(retentionDays),
517 > HistoryArchivalState: archivalState,
518 > HistoryArchivalUri: historyArchivalURI,
519 > VisibilityArchivalState: archivalState,
520 > VisibilityArchivalUri: visibilityArchivalURI,
521 > BadBinaries: &namespacepb.BadBinaries{Binaries: map[string]*namespacepb.BadBinaryInfo{}},
522 > },
523 > ReplicationConfig: &persistencespb.NamespaceReplicationConfig{
524 > ActiveClusterName: currentClusterName,
525 > Clusters: []string{
526 > currentClusterName,
527 > },
528 > },
529 >
530 > FailoverVersion: common.EmptyVersion,
531 > },
532 > IsGlobalNamespace: false,
533 > }
534 > _, err := s.testCluster.testBase.MetadataManager.CreateNamespace(context.Background(), namespaceRequest)
535 >
536 > if err != nil {
537 return namespace.EmptyID, err
538 }
539
540 > namespaceCacheDeadline := time.Now().Add(5 * NamespaceCacheRefreshInterval) functional_test_base.go
541 > ticker := time.NewTicker(NamespaceCacheRefreshInterval / 2)
542 > defer ticker.Stop()
543 > for {
544 > _, describeErr := s.FrontendClient().DescribeNamespace(NewContext(), &workflowservice.DescribeNamespaceRequest{
545 > Namespace: nsName.String(),
546 > })
547 > if describeErr == nil {
548 > break
549 }
550 if time.Now().After(namespaceCacheDeadline) {
554 }
555
556 > _, err = s.OperatorClient().AddSearchAttributes(NewContext(), &operatorservice.AddSearchAttributesRequest{ functional_test_base.go
557 > Namespace: nsName.String(),
558 > SearchAttributes: expectedSearchAttributes,
559 > })
560 > if err != nil {
561 return namespace.EmptyID, err
562 }
563
564 > namespaceCacheDeadline = time.Now().Add(5 * NamespaceCacheRefreshInterval) functional_test_base.go
565 > for {
566 > listResp, listErr := s.OperatorClient().ListSearchAttributes(NewContext(), &operatorservice.ListSearchAttributesRequest{
567 > Namespace: nsName.String(),
568 > })
569 > if listErr == nil {
570 > customAttrs := listResp.GetCustomAttributes()
571 > allFound := true
572 > for saName := range expectedSearchAttributes {
573 > if _, ok := customAttrs[saName]; !ok {
574 allFound = false
575 break
576 }
577 }
578 > if allFound { functional_test_base.go
579 > break
580 }
581 }
586 }
587
588 > s.Logger.Info("Register namespace succeeded", functional_test_base.go
589 > tag.WorkflowNamespace(nsName.String()),
590 > tag.WorkflowNamespaceID(nsID.String()),
591 > )
592 > return nsID, nil
593 }
594
595 func (s *FunctionalTestBase) MarkNamespaceAsDeleted(
596 nsName namespace.Name,
597 > ) error { functional_test_base.go
598 > ctx, cancel := rpc.NewContextWithTimeoutAndVersionHeaders(10000 * time.Second)
599 > defer cancel()
600 > _, err := s.FrontendClient().UpdateNamespace(ctx, &workflowservice.UpdateNamespaceRequest{
601 > Namespace: nsName.String(),
602 > UpdateInfo: &namespacepb.UpdateNamespaceInfo{
603 > State: enumspb.NAMESPACE_STATE_DELETED,
604 > },
605 > })
606 >
607 > return err
608 > }
609
610 func (s *FunctionalTestBase) GetHistoryFunc(namespace string, execution *commonpb.WorkflowExecution) func() []*historypb.HistoryEvent {
643 }
644
645 > func (s *FunctionalTestBase) OverrideDynamicConfig(setting dynamicconfig.GenericSetting, value any) (cleanup func()) { functional_test_base.go
646 > return s.testCluster.host.overrideDynamicConfigForTest(s.T(), setting.Key(), value)
647 > }
648
649 // InjectHook sets a test hook inside the cluster.
go.temporal.io/server/tests/testcore/test_cluster.go 191 introduced LOC · 28 ranges

Open complete file

104 )
105
106 > func (f *defaultTestClusterFactory) NewCluster(t *testing.T, clusterConfig *TestClusterConfig, logger log.Logger) (*TestCluster, error) { test_cluster.go
107 > return newClusterWithPersistenceTestBaseFactory(t, clusterConfig, logger, f.tbFactory)
108 > }
109
110 > func NewTestClusterFactory() TestClusterFactory { test_cluster.go
111 > tbFactory := &defaultPersistenceTestBaseFactory{}
112 > return newTestClusterFactoryWithCustomTestBaseFactory(tbFactory)
113 > }
114
115 > func newTestClusterFactoryWithCustomTestBaseFactory(tbFactory persistenceTestBaseFactory) TestClusterFactory { test_cluster.go
116 > return &defaultTestClusterFactory{
117 > tbFactory: tbFactory,
118 > }
119 > }
120
121 type persistenceTestBaseFactory interface {
127 // GetPersistenceTestDefaults returns the default persistence options based on CLI flags.
128 // Use this when creating TestClusterConfig to ensure proper database configuration.
129 > func GetPersistenceTestDefaults() persistencetests.TestBaseOptions { test_cluster.go
130 > return *persistencetests.GetTestClusterOption(cliFlags.persistenceType, cliFlags.persistenceDriver)
131 > }
132
133 > func (f *defaultPersistenceTestBaseFactory) NewTestBase(options *persistencetests.TestBaseOptions) *persistencetests.TestBase { test_cluster.go
134 > defaults := GetPersistenceTestDefaults()
135 > options.ApplyDefaults(&defaults)
136 >
137 > if cliFlags.enableFaultInjection != "" && options.FaultInjection == nil {
138 // If -enableFaultInjection is passed to the test runner, then default fault injection config is added to the persistence options.
139 // If FaultInjectionConfig is already set by test, then it means that this test requires
142 }
143
144 > return persistencetests.NewTestBase(options) test_cluster.go
145 }
146
150 logger log.Logger,
151 tbFactory persistenceTestBaseFactory,
152 > ) (*TestCluster, error) { test_cluster.go
153 > // determine number of hosts per service
154 > const minNodes = 1
155 > clusterConfig.FrontendConfig.NumFrontendHosts = max(minNodes, clusterConfig.FrontendConfig.NumFrontendHosts)
156 > clusterConfig.HistoryConfig.NumHistoryHosts = max(minNodes, clusterConfig.HistoryConfig.NumHistoryHosts)
157 > clusterConfig.MatchingConfig.NumMatchingHosts = max(minNodes, clusterConfig.MatchingConfig.NumMatchingHosts)
158 > clusterConfig.WorkerConfig.NumWorkers = max(minNodes, clusterConfig.WorkerConfig.NumWorkers)
159 > if clusterConfig.WorkerConfig.DisableWorker {
160 clusterConfig.WorkerConfig.NumWorkers = 0
161 }
162
163 // allocate ports
164 > hostsByProtocolByService := map[transferProtocol]map[primitives.ServiceName]static.Hosts{ test_cluster.go
165 > grpcProtocol: {
166 > primitives.FrontendService: {All: makeAddresses(clusterConfig.FrontendConfig.NumFrontendHosts)},
167 > primitives.MatchingService: {All: makeAddresses(clusterConfig.MatchingConfig.NumMatchingHosts)},
168 > primitives.HistoryService: {All: makeAddresses(clusterConfig.HistoryConfig.NumHistoryHosts)},
169 > primitives.WorkerService: {All: makeAddresses(clusterConfig.WorkerConfig.NumWorkers)},
170 > },
171 > httpProtocol: {
172 > primitives.FrontendService: {All: makeAddresses(clusterConfig.FrontendConfig.NumFrontendHosts)},
173 > },
174 > }
175 >
176 > if len(clusterConfig.ClusterMetadata.ClusterInformation) > 0 {
177 // set self-address for current cluster
178 ci := clusterConfig.ClusterMetadata.ClusterInformation[clusterConfig.ClusterMetadata.CurrentClusterName]
182 }
183
184 > clusterMetadataConfig := cluster.NewTestClusterMetadataConfig( test_cluster.go
185 > clusterConfig.ClusterMetadata.EnableGlobalNamespace,
186 > clusterConfig.IsMasterCluster,
187 > )
188 > if !clusterConfig.IsMasterCluster && clusterConfig.ClusterMetadata.MasterClusterName != "" { // xdc cluster metadata setup
189 clusterMetadataConfig = &cluster.Config{
190 EnableGlobalNamespace: clusterConfig.ClusterMetadata.EnableGlobalNamespace,
195 }
196 }
197 > clusterConfig.Persistence.Logger = logger test_cluster.go
198 > clusterConfig.Persistence.FaultInjection = clusterConfig.FaultInjection
199 >
200 > testBase := tbFactory.NewTestBase(&clusterConfig.Persistence)
201 >
202 > testBase.Setup(clusterMetadataConfig)
203 > archiverMetadata, archiverProvider := newArchivalMetadataAndProvider(
204 > clusterConfig.EnableArchival,
205 > clusterConfig.CustomHistoryArchiverFactory,
206 > clusterConfig.CustomVisibilityArchiverFactory,
207 > testBase.ExecutionManager,
208 > logger,
209 > )
210 > var err error
211 >
212 > pConfig := testBase.DefaultTestCluster.Config()
213 > pConfig.NumHistoryShards = clusterConfig.HistoryConfig.NumHistoryShards
214 >
215 > var (
216 > esClient esclient.Client
217 > )
218 > if !UseSQLVisibility() {
219 clusterConfig.ESConfig = &esclient.Config{
220 Indices: map[string]string{
242 return nil, err
243 }
244 > } else { test_cluster.go
245 > clusterConfig.ESConfig = nil
246 > }
247
248 > clusterInfoMap := make(map[string]cluster.ClusterInformation) test_cluster.go
249 > for clusterName, clusterInfo := range clusterMetadataConfig.ClusterInformation {
250 > clusterInfo.ShardCount = clusterConfig.HistoryConfig.NumHistoryShards
251 > clusterInfo.ClusterID = uuid.NewString()
252 > clusterInfoMap[clusterName] = clusterInfo
253 > _, err := testBase.ClusterMetadataManager.SaveClusterMetadata(context.Background(), &persistence.SaveClusterMetadataRequest{
254 > ClusterMetadata: &persistencespb.ClusterMetadata{
255 > HistoryShardCount: clusterConfig.HistoryConfig.NumHistoryShards,
256 > ClusterName: clusterName,
257 > ClusterId: clusterInfo.ClusterID,
258 > IsConnectionEnabled: clusterInfo.Enabled,
259 > IsReplicationEnabled: clusterInfo.ReplicationEnabled,
260 > IsGlobalNamespaceEnabled: clusterMetadataConfig.EnableGlobalNamespace,
261 > FailoverVersionIncrement: clusterMetadataConfig.FailoverVersionIncrement,
262 > ClusterAddress: clusterInfo.RPCAddress,
263 > HttpAddress: clusterInfo.HTTPAddress,
264 > InitialFailoverVersion: clusterInfo.InitialFailoverVersion,
265 > }},
266 > )
267 > if err != nil {
268 return nil, err
269 }
270 }
271 > clusterMetadataConfig.ClusterInformation = clusterInfoMap test_cluster.go
272 >
273 > cfg := &config.Config{
274 > Persistence: pConfig,
275 > ClusterMetadata: clusterMetadataConfig,
276 > Visibility: config.Visibility{},
277 > }
278 > clusterMetadataConfig, pConfig, err = temporal.ApplyClusterMetadataConfigProvider(
279 > logger,
280 > cfg,
281 > resolver.NewNoopResolver(),
282 > persistenceclient.FactoryProvider,
283 > testBase.AbstractDataStoreFactory,
284 > testBase.VisibilityStoreFactory,
285 > metrics.NoopMetricsHandler,
286 > serialization.NewSerializer(),
287 > )
288 > if err != nil {
289 return nil, err
290 }
291
292 > var tlsConfigProvider *encryption.FixedTLSConfigProvider test_cluster.go
293 > if clusterConfig.TLSConfigProvider != nil {
294 tlsConfigProvider = clusterConfig.TLSConfigProvider
295 > } else if clusterConfig.EnableMTLS { test_cluster.go
296 if tlsConfigProvider, err = createFixedTLSConfigProvider(); err != nil {
297 return nil, err
299 }
300
301 > temporalParams := &temporalParams{ test_cluster.go
302 > clusterMetadataConfig: clusterMetadataConfig,
303 > persistenceConfig: pConfig,
304 > metadataMgr: testBase.MetadataManager,
305 > clusterMetadataManager: testBase.ClusterMetadataManager,
306 > shardMgr: testBase.ShardMgr,
307 > executionManager: testBase.ExecutionManager,
308 > namespaceReplicationQueue: testBase.NamespaceReplicationQueue,
309 > abstractDataStoreFactory: testBase.AbstractDataStoreFactory,
310 > visibilityStoreFactory: testBase.VisibilityStoreFactory,
311 > taskMgr: testBase.TaskMgr,
312 > logger: logger,
313 > esConfig: clusterConfig.ESConfig,
314 > esClient: esClient,
315 > archiverMetadata: archiverMetadata,
316 > archiverProvider: archiverProvider,
317 > frontendConfig: clusterConfig.FrontendConfig,
318 > historyConfig: clusterConfig.HistoryConfig,
319 > matchingConfig: clusterConfig.MatchingConfig,
320 > workerConfig: clusterConfig.WorkerConfig,
321 > mockAdminClient: clusterConfig.MockAdminClient,
322 > namespaceReplicationTaskExecutor: nsreplication.NewTaskExecutor(clusterConfig.ClusterMetadata.CurrentClusterName, testBase.MetadataManager, nsreplication.NewNoopDataMerger(), nsreplication.NewDefaultAdmitter(), logger, testhooks.TestHooks{}),
323 > dcRedirectionPolicy: clusterConfig.DCRedirectionPolicy,
324 > dynamicConfigOverrides: clusterConfig.DynamicConfigOverrides,
325 > tlsConfigProvider: tlsConfigProvider,
326 > serviceFxOptions: clusterConfig.ServiceFxOptions,
327 > taskCategoryRegistry: temporal.TaskCategoryRegistryProvider(archiverMetadata),
328 > hostsByProtocolByService: hostsByProtocolByService,
329 > spanExporters: clusterConfig.SpanExporters,
330 > tokenProvider: clusterConfig.TokenProvider,
331 > enableHistoryTaskRecorder: clusterConfig.EnableHistoryTaskRecorder,
332 > }
333 >
334 > if clusterConfig.EnableMetricsCapture {
335 > temporalParams.captureMetricsHandler = metricstest.NewCaptureHandler()
336 > }
337
338 > err = newPProfInitializerImpl(logger, PprofTestPort).Start() test_cluster.go
339 > if err != nil {
340 logger.Fatal("Failed to start pprof", tag.Error(err))
341 }
342
343 > cluster := newTemporal(t, temporalParams) test_cluster.go
344 > if err = cluster.Start(); err != nil {
345 return nil, err
346 }
347
348 > return &TestCluster{testBase: testBase, host: cluster}, nil test_cluster.go
349 }
350
452 }
453
454 > func newPProfInitializerImpl(logger log.Logger, port int) *pprof.PProfInitializerImpl { test_cluster.go
455 > return &pprof.PProfInitializerImpl{
456 > PProf: &config.PProf{
457 > Port: port,
458 > },
459 > Logger: logger,
460 > }
461 > }
462
463 // TODO: remove when onebox uses temporal.NewServerFx and archival wiring comes from production config.
468 executionManager persistence.ExecutionManager,
469 logger log.Logger,
470 > ) (archiver.ArchivalMetadata, provider.ArchiverProvider) { test_cluster.go
471 > dcCollection := dynamicconfig.NewNoopCollection()
472 > if !enabled {
473 > return archiver.NewArchivalMetadata(dcCollection, "", false, "", false, &config.ArchivalNamespaceDefaults{}),
474 > provider.NewArchiverProvider(nil, nil, nil, nil, nil, logger, metrics.NoopMetricsHandler)
475 > }
476
477 cfg := &config.FilestoreArchiver{
505
506 // TearDownCluster tears down the test cluster
507 > func (tc *TestCluster) TearDownCluster() error { test_cluster.go
508 > errs := tc.host.Stop()
509 > tc.testBase.TearDownWorkflowStore()
510 > if !UseSQLVisibility() && tc.host.esConfig != nil {
511 if err := deleteIndex(tc.host.esConfig, tc.host.logger); err != nil {
512 errs = multierr.Combine(errs, err)
513 }
514 }
515 > return errs test_cluster.go
516 }
517
521 }
522
523 > func (tc *TestCluster) FrontendClient() workflowservice.WorkflowServiceClient { test_cluster.go
524 > return tc.host.FrontendClient()
525 > }
526
527 func (tc *TestCluster) AdminClient() adminservice.AdminServiceClient {
529 }
530
531 > func (tc *TestCluster) OperatorClient() operatorservice.OperatorServiceClient { test_cluster.go
532 > return tc.host.OperatorClient()
533 > }
534
535 // HistoryClient returns a history client from the test cluster
554
555 // TODO (alex): expose only needed objects from TemporalImpl.
556 > func (tc *TestCluster) Host() *temporalImpl { test_cluster.go
557 > return tc.host
558 > }
559
560 func (tc *TestCluster) InjectHook(t *testing.T, hook testhooks.Hook, scope any) func() {
562 }
563
564 > func (tc *TestCluster) WorkerGRPCAddress() string { test_cluster.go
565 > return tc.host.WorkerGRPCAddress()
566 > }
567
568 func (tc *TestCluster) ClusterName() string {
631 }
632
633 > func makeAddresses(count int) []string { test_cluster.go
634 > hosts := make([]string, count)
635 > for i := range hosts {
636 > hosts[i] = fmt.Sprintf("127.0.0.1:%d", freeport.MustGetFreePort())
637 > }
638 > return hosts
639 }
go.temporal.io/server/tests/testcore/clients.go 89 introduced LOC · 17 ranges

Open complete file

83 metadataMgr persistence.MetadataManager,
84 tokenProvider auth.TokenProvider,
85 > ) clients { clients.go
86 > return clients{
87 > logger: logger,
88 > hostsByService: hostsByService,
89 > frontendMembershipAddress: frontendMembershipAddress,
90 > tlsConfigProvider: tlsConfigProvider,
91 > metricsHandler: metricsHandler,
92 > dcClient: dcClient,
93 > testHooks: testHooks,
94 > numHistoryShards: numHistoryShards,
95 > metadataMgr: metadataMgr,
96 > tokenProvider: tokenProvider,
97 > }
98 > }
99
100 func (c *clients) AdminClient() adminservice.AdminServiceClient {
103 }
104
105 > func (c *clients) OperatorClient() operatorservice.OperatorServiceClient { clients.go
106 > c.ensureFrontend()
107 > return c.frontend.operator
108 > }
109
110 > func (c *clients) FrontendClient() workflowservice.WorkflowServiceClient { clients.go
111 > c.ensureFrontend()
112 > return c.frontend.frontend
113 > }
114
115 > func (c *clients) ensureFrontend() { clients.go
116 > c.frontend.once.Do(func() {
117 > conn, err := c.newConn(primitives.FrontendService)
118 > if err != nil {
119 c.logger.Fatal("unable to create frontend test client", tag.Error(err))
120 }
121 > c.frontend.conn = conn clients.go
122 > c.frontend.admin = adminservice.NewAdminServiceClient(conn)
123 > c.frontend.frontend = workflowservice.NewWorkflowServiceClient(conn)
124 > c.frontend.operator = operatorservice.NewOperatorServiceClient(conn)
125 })
126 }
212 }
213
214 > func (c *clients) close() []error { clients.go
215 > var errs []error
216 > for _, conn := range []*grpc.ClientConn{
217 > c.frontend.conn,
218 > c.history.conn,
219 > } {
220 > if conn != nil {
221 > errs = append(errs, conn.Close())
222 > }
223 }
224 > c.frontend.conn = nil clients.go
225 > c.history.conn = nil
226 > return errs
227 }
228
229 > func (c *clients) newConn(serviceName primitives.ServiceName) (*grpc.ClientConn, error) { clients.go
230 > address, err := c.grpcAddress(serviceName)
231 > if err != nil {
232 return nil, err
233 }
234 > tlsConfig, err := c.tlsConfig(serviceName) clients.go
235 > if err != nil {
236 return nil, err
237 }
238
239 > return rpc.Dial(address, tlsConfig, c.logger, metrics.NoopMetricsHandler) clients.go
240 }
241
242 > func (c *clients) grpcAddress(serviceName primitives.ServiceName) (string, error) { clients.go
243 > hosts := c.hostsByService[serviceName].All
244 > if len(hosts) == 0 {
245 return "", fmt.Errorf("no %s gRPC hosts configured", serviceName)
246 }
247 > return hosts[0], nil clients.go
248 }
249
250 > func (c *clients) tlsConfig(serviceName primitives.ServiceName) (*tls.Config, error) { clients.go
251 > if c.tlsConfigProvider == nil {
252 > return nil, nil
253 > }
254 if serviceName == primitives.FrontendService {
255 return c.tlsConfigProvider.GetFrontendClientConfig()
261 clusterConfig *cluster.Config,
262 mockAdminClient map[string]adminservice.AdminServiceClient,
263 > ) client.FactoryProvider { clients.go
264 > return &clientFactoryProvider{
265 > config: clusterConfig,
266 > mockAdminClient: mockAdminClient,
267 > }
268 > }
269
270 type clientFactoryProvider struct {
282 logger log.Logger,
283 throttledLogger log.Logger,
284 > ) client.Factory { clients.go
285 > f := client.NewFactoryProvider().NewFactory(
286 > rpcFactory,
287 > monitor,
288 > metricsHandler,
289 > dc,
290 > testHooks,
291 > numberOfHistoryShards,
292 > logger,
293 > throttledLogger,
294 > )
295 > return &clientFactory{
296 > Factory: f,
297 > config: p.config,
298 > mockAdminClient: p.mockAdminClient,
299 > }
300 > }
301
302 type clientFactory struct {
326 dc *dynamicconfig.Collection,
327 tlsConfigProvider encryption.TLSConfigProvider,
328 > ) sdk.ClientFactory { clients.go
329 > var tlsConfig *tls.Config
330 > if tlsConfigProvider != nil {
331 var err error
332 if tlsConfig, err = tlsConfigProvider.GetFrontendClientConfig(); err != nil {
334 }
335 }
336 > return sdk.NewClientFactory( clients.go
337 > grpcResolver.MakeURL(primitives.FrontendService),
338 > tlsConfig,
339 > metricsHandler,
340 > logger,
341 > dynamicconfig.WorkerStickyCacheSize.Get(dc),
342 > )
343 }
go.temporal.io/server/common/membership/static/service_resolver.go 65 introduced LOC · 13 ranges

Open complete file

17 }
18
19 > func newStaticResolver(hosts []string) *staticResolver { service_resolver.go
20 > hostInfos := make([]membership.HostInfo, 0, len(hosts))
21 > for _, host := range hosts {
22 > hostInfos = append(hostInfos, membership.NewHostInfoFromAddress(host))
23 > }
24 > return &staticResolver{
25 > hostInfos: hostInfos,
26 > hashfunc: farm.Fingerprint32,
27 > listeners: make(map[string]chan<- *membership.ChangedEvent),
28 > }
29 }
30
31 > func (s *staticResolver) start(hosts []string) { service_resolver.go
32 > hostInfos := make([]membership.HostInfo, 0, len(hosts))
33 > for _, host := range hosts {
34 > hostInfos = append(hostInfos, membership.NewHostInfoFromAddress(host))
35 > }
36 > event := &membership.ChangedEvent{
37 > HostsAdded: hostInfos,
38 > }
39 >
40 > s.mu.Lock()
41 > defer s.mu.Unlock()
42 >
43 > s.hostInfos = hostInfos
44 >
45 > for _, ch := range s.listeners {
46 > select {
47 > case ch <- event:
48 default:
49 }
51 }
52
53 > func (s *staticResolver) Lookup(key string) (membership.HostInfo, error) { service_resolver.go
54 > s.mu.Lock()
55 > defer s.mu.Unlock()
56 > if len(s.hostInfos) == 0 {
57 return nil, membership.ErrInsufficientHosts
58 }
59 > hash := int(s.hashfunc([]byte(key))) service_resolver.go
60 > idx := hash % len(s.hostInfos)
61 > return s.hostInfos[idx], nil
62 }
63
64 > func (s *staticResolver) LookupN(key string, _ int) []membership.HostInfo { service_resolver.go
65 > info, err := s.Lookup(key)
66 > if err != nil {
67 return []membership.HostInfo{}
68 }
69 > return []membership.HostInfo{info} service_resolver.go
70 }
71
72 > func (s *staticResolver) AddListener(name string, notifyChannel chan<- *membership.ChangedEvent) error { service_resolver.go
73 > s.mu.Lock()
74 > defer s.mu.Unlock()
75 > _, ok := s.listeners[name]
76 > if ok {
77 return membership.ErrListenerAlreadyExist
78 }
79 > s.listeners[name] = notifyChannel service_resolver.go
80 > return nil
81 }
82
83 > func (s *staticResolver) RemoveListener(name string) error { service_resolver.go
84 > s.mu.Lock()
85 > defer s.mu.Unlock()
86 > _, ok := s.listeners[name]
87 > if !ok {
88 return nil
89 }
90 > delete(s.listeners, name) service_resolver.go
91 > return nil
92 }
93
98 }
99
100 > func (s *staticResolver) AvailableMemberCount() int { service_resolver.go
101 > s.mu.Lock()
102 > defer s.mu.Unlock()
103 > return len(s.hostInfos)
104 > }
105
106 > func (s *staticResolver) Members() []membership.HostInfo { service_resolver.go
107 > s.mu.Lock()
108 > defer s.mu.Unlock()
109 > return s.hostInfos
110 > }
111
112 > func (s *staticResolver) AvailableMembers() []membership.HostInfo { service_resolver.go
113 > return s.Members()
114 > }
115
116 func (s *staticResolver) RequestRefresh() {
go.temporal.io/server/tests/testcore/replication_stream_recorder.go 65 introduced LOC · 11 ranges

Open complete file

45 }
46
47 > func NewReplicationStreamRecorder() *ReplicationStreamRecorder { replication_stream_recorder.go
48 > return &ReplicationStreamRecorder{
49 > capturedMessages: make([]CapturedReplicationMessage, 0),
50 > }
51 > }
52
53 // SetOutputFile sets the file path for writing captured messages on-demand
54 > func (r *ReplicationStreamRecorder) SetOutputFile(filePath string) { replication_stream_recorder.go
55 > r.mu.Lock()
56 > defer r.mu.Unlock()
57 > r.outputFilePath = filePath
58 > }
59
60 // WriteToLog writes all captured messages to the configured output file
157
158 // UnaryInterceptor returns a gRPC unary client interceptor that captures messages
159 > func (r *ReplicationStreamRecorder) UnaryInterceptor(clusterName string) grpc.UnaryClientInterceptor { replication_stream_recorder.go
160 > return func(
161 > ctx context.Context,
162 > method string,
163 > req, reply any,
164 > cc *grpc.ClientConn,
165 > invoker grpc.UnaryInvoker,
166 > opts ...grpc.CallOption,
167 > ) error {
168 > target := cc.Target()
169 >
170 > // Capture outgoing request if it's a replication-related call
171 > if isReplicationMethod(method) {
172 if protoReq, ok := req.(proto.Message); ok {
173 r.recordMessage(method, protoReq, DirectionSend, clusterName, target, false)
175 }
176
177 > err := invoker(ctx, method, req, reply, cc, opts...) replication_stream_recorder.go
178 >
179 > // Capture incoming response if successful
180 > if err == nil && isReplicationMethod(method) {
181 if protoReply, ok := reply.(proto.Message); ok {
182 r.recordMessage(method, protoReply, DirectionRecv, clusterName, target, false)
184 }
185
187 }
188 }
189
190 // StreamInterceptor returns a gRPC stream client interceptor that captures stream messages
191 > func (r *ReplicationStreamRecorder) StreamInterceptor(clusterName string) grpc.StreamClientInterceptor { replication_stream_recorder.go
192 > return func(
193 > ctx context.Context,
194 > desc *grpc.StreamDesc,
195 > cc *grpc.ClientConn,
196 > method string,
197 > streamer grpc.Streamer,
198 > opts ...grpc.CallOption,
199 > ) (grpc.ClientStream, error) {
200 stream, err := streamer(ctx, desc, cc, method, opts...)
201 if err != nil {
246
247 // UnaryServerInterceptor returns a gRPC unary server interceptor that captures messages
248 > func (r *ReplicationStreamRecorder) UnaryServerInterceptor(clusterName string) grpc.UnaryServerInterceptor { replication_stream_recorder.go
249 > return func(
250 > ctx context.Context,
251 > req any,
252 > info *grpc.UnaryServerInfo,
253 > handler grpc.UnaryHandler,
254 > ) (any, error) {
255 > // Capture incoming request if it's a replication-related call
256 > if isReplicationMethod(info.FullMethod) {
257 if protoReq, ok := req.(proto.Message); ok {
258 r.recordMessage(info.FullMethod, protoReq, DirectionServerRecv, clusterName, "server", false)
260 }
261
262 > resp, err := handler(ctx, req) replication_stream_recorder.go
263 >
264 > // Capture outgoing response if successful
265 > if err == nil && isReplicationMethod(info.FullMethod) {
266 if protoResp, ok := resp.(proto.Message); ok {
267 r.recordMessage(info.FullMethod, protoResp, DirectionServerSend, clusterName, "server", false)
269 }
270
271 > return resp, err replication_stream_recorder.go
272 }
273 }
274
275 // StreamServerInterceptor returns a gRPC stream server interceptor that captures stream messages
276 > func (r *ReplicationStreamRecorder) StreamServerInterceptor(clusterName string) grpc.StreamServerInterceptor { replication_stream_recorder.go
277 > return func(
278 > srv any,
279 > ss grpc.ServerStream,
280 > info *grpc.StreamServerInfo,
281 > handler grpc.StreamHandler,
282 > ) error {
283 if isReplicationMethod(info.FullMethod) {
284 wrappedStream := &recordingServerStream{
322 }
323
324 > func isReplicationMethod(method string) bool { replication_stream_recorder.go
325 > // Capture StreamWorkflowReplicationMessages from both history and admin services
326 > // - Sender (active) uses history service to respond to receiver
327 > // - Receiver (standby) uses admin service to call sender
328 > return method == "/temporal.server.api.historyservice.v1.HistoryService/StreamWorkflowReplicationMessages" ||
329 > method == "/temporal.server.api.adminservice.v1.AdminService/StreamWorkflowReplicationMessages"
330 > }
go.temporal.io/server/common/membership/static/monitor.go 26 introduced LOC · 8 ranges

Open complete file

28
29 // NewMonitor creates a new Monitor with the given host mappings.
30 > func NewMonitor(hosts map[primitives.ServiceName]Hosts) membership.Monitor { monitor.go
31 > resolvers := make(map[primitives.ServiceName]*staticResolver, len(hosts))
32 > for service, hostList := range hosts {
33 > resolvers[service] = newStaticResolver(hostList.All)
34 > }
35
36 > return &staticMonitor{ monitor.go
37 > hosts: hosts,
38 > resolvers: resolvers,
39 > }
40 }
41
42 > func (s *staticMonitor) Start() { monitor.go
43 > for service, r := range s.resolvers {
44 > r.start(s.hosts[service].All)
45 > }
46 }
47
48 > func (s *staticMonitor) EvictSelf() error { monitor.go
49 > return nil
50 > }
51
52 func (s *staticMonitor) EvictSelfAt(asOf time.Time) (time.Duration, error) {
54 }
55
56 > func (s *staticMonitor) GetResolver(service primitives.ServiceName) (membership.ServiceResolver, error) { monitor.go
57 > resolver, ok := s.resolvers[service]
58 > if !ok {
59 return nil, membership.ErrUnknownService
60 }
61 > return resolver, nil monitor.go
62 }
63
66 }
67
68 > func (s *staticMonitor) WaitUntilInitialized(_ context.Context) error { monitor.go
69 > return nil
70 > }
71
72 > func (s *staticMonitor) SetDraining(draining bool) error { monitor.go
73 > return nil
74 > }
75
76 func (s *staticMonitor) ApproximateMaxPropagationTime() time.Duration {
go.temporal.io/server/tests/testcore/test_cluster_pool.go 20 introduced LOC · 6 ranges

Open complete file

107 func (p *clusterPool) reserveSlot(t *testing.T) *clusterPoolSlot {
108 if p.availableSlots != nil {
109 > slot := <-p.availableSlots test_cluster_pool.go
110 > t.Cleanup(func() { p.availableSlots <- slot })
111 > return slot
112 }
113 return p.nextSlot()
249 // reason explains why the cluster was created, for analytics. It falls back to a
250 // generic reason when the caller did not provide one.
251 > func (r clusterRequest) reason() string { test_cluster_pool.go
252 > switch r.kind {
253 case clusterKindShared:
254 return "shared pool"
256 return "suite-scoped"
257 }
258 > switch { test_cluster_pool.go
259 > case r.dedicatedReason != "":
260 > return r.dedicatedReason
261 case r.mustBeFresh():
262 return "custom config"
269 // run can be queried for which suite created how many clusters of each kind, and
270 // why. Events fall back to the test log when no events file is configured.
271 > func (r clusterRequest) recordCreation(t *testing.T) { test_cluster_pool.go
272 > suite, _, _ := strings.Cut(t.Name(), "/")
273 > line, err := json.Marshal(map[string]any{
274 > "suite": suite,
275 > "test": t.Name(),
276 > "kind": r.kind,
277 > "reason": r.reason(),
278 > "worker": r.needWorkerService,
279 > })
280 > if err != nil {
281 return
282 }
283
284 > if testClusterRouter.eventsFile == nil { test_cluster_pool.go
285 log.Printf("CLUSTEREVENT %s", line)
286 return
288 // O_APPEND makes each write land atomically at EOF and os.File serializes
289 // concurrent writes, so lines from parallel tests don't interleave.
290 > _, _ = testClusterRouter.eventsFile.Write(append(line, '\n')) test_cluster_pool.go
291 }
292
go.temporal.io/server/common/persistence/persistence-tests/persistence_test_base.go 16 introduced LOC · 2 ranges

Open complete file

66
67 // ApplyDefaults copies database configuration from src, preserving any non-zero values already set.
68 > func (o *TestBaseOptions) ApplyDefaults(src *TestBaseOptions) { persistence_test_base.go
69 > o.StoreType = cmp.Or(o.StoreType, src.StoreType)
70 > o.SQLDBPluginName = cmp.Or(o.SQLDBPluginName, src.SQLDBPluginName)
71 > o.DBName = cmp.Or(o.DBName, src.DBName)
72 > o.DBUsername = cmp.Or(o.DBUsername, src.DBUsername)
73 > o.DBPassword = cmp.Or(o.DBPassword, src.DBPassword)
74 > o.DBHost = cmp.Or(o.DBHost, src.DBHost)
75 > o.DBPort = cmp.Or(o.DBPort, src.DBPort)
76 > o.SchemaDir = cmp.Or(o.SchemaDir, src.SchemaDir)
77 > if o.ConnectAttributes == nil {
78 > o.ConnectAttributes = src.ConnectAttributes
79 > }
80 }
81
172
173 // NewTestBase returns a persistence test base backed by either cassandra or sql
174 > func NewTestBase(options *TestBaseOptions) *TestBase { persistence_test_base.go
175 > switch options.StoreType {
176 > case config.StoreTypeSQL:
177 > return NewTestBaseWithSQL(options)
178 case config.StoreTypeNoSQL:
179 return NewTestBaseWithCassandra(options)
go.temporal.io/server/temporal/fx.go 16 introduced LOC · 3 ranges

Open complete file

711 )
712 switch err.(type) {
713 > case nil: fx.go
714 > // Update current record
715 > if updateErr := updateCurrentClusterMetadataRecord(
716 > ctx,
717 > clusterMetadataManager,
718 > svc,
719 > indexSearchAttributes,
720 > resp,
721 > ); updateErr != nil {
722 return svc.ClusterMetadata, svc.Persistence, updateErr
723 }
724 // Ignore invalid cluster metadata
725 > overwriteCurrentClusterMetadataWithDBRecord( fx.go
726 > svc,
727 > resp,
728 > logger,
729 > )
730 case *serviceerror.NotFound:
731 // Initialize current cluster record
831
832 if updateIndexSearchAttributes(initialIndexSearchAttributes, currentClusterDBRecord) {
833 > updateDBRecord = true fx.go
834 > }
835
836 if !updateDBRecord {
go.temporal.io/server/service/frontend/workflow_handler.go 15 introduced LOC · 3 ranges

Open complete file

1124
1125 // These errors are expected from some versioning situations. We should not log them, it'd be too noisy.
1126 > var newerBuild *serviceerror.NewerBuildExists // expected when versioned poller is superceded workflow_handler.go
1127 > var failedPrecond *serviceerror.FailedPrecondition // expected when user data is disabled
1128 > if errors.As(err, &newerBuild) || errors.As(err, &failedPrecond) {
1129 return nil, err
1130 }
1131
1132 // For all other errors log an error and return it back to client.
1133 > ctxTimeout := "not-set" workflow_handler.go
1134 > ctxDeadline, ok := childCtx.Deadline()
1135 > if ok {
1136 > ctxTimeout = ctxDeadline.Sub(callTime).String()
1137 > }
1138 > wh.logger.Error("Unable to call matching.PollWorkflowTaskQueue.",
1139 > tag.WorkflowTaskQueueName(request.GetTaskQueue().GetName()),
1140 > tag.Timeout(ctxTimeout),
1141 > tag.Error(err))
1142 > return nil, err
1143 }
1144
6711 // First check if this err is due to context cancellation. This means client connection to frontend is closed.
6712 if !errors.Is(ctx.Err(), context.Canceled) {
6713 > return false workflow_handler.go
6714 > }
6715 // Our rpc stack does not propagates context cancellation to the other service. Lets make an explicit
6716 // call to matching to notify this poller is gone to prevent any tasks being dispatched to zombie pollers.
go.temporal.io/server/common/membership/static/fx.go 14 introduced LOC · 3 ranges

Open complete file

11 func MembershipModule(
12 hostsByService map[primitives.ServiceName]Hosts,
13 > ) fx.Option { fx.go
14 > return fx.Options(
15 > fx.Provide(func() membership.Monitor {
16 > return NewMonitor(hostsByService)
17 > }),
18 > fx.Provide(func(serviceName primitives.ServiceName) membership.HostInfoProvider {
19 > hosts := hostsByService[serviceName]
20 > if len(hosts.All) == 0 {
21 panic(fmt.Sprintf("hosts for %v service are missing in static hosts", serviceName))
22 }
23 > if len(hosts.Self) == 0 { fx.go
24 panic(fmt.Sprintf("self host for %v service is missing in static hosts", serviceName))
25 }
26
27 > for _, serviceHost := range hosts.All { fx.go
28 > if serviceHost == hosts.Self {
29 > hostInfo := membership.NewHostInfoFromAddress(serviceHost)
30 > return membership.NewHostInfoProvider(hostInfo)
31 > }
32 }
33 panic(fmt.Sprintf("self host %v for %v service is defined, but missing from the list of all static hosts", hosts.Self, serviceName))
go.temporal.io/server/common/rpc/interceptor/namespace_logger.go 12 introduced LOC · 2 ranges

Open complete file

39
40 if nli.logger != nil {
41 > methodName := api.MethodName(info.FullMethod) namespace_logger.go
42 > namespace := MustGetNamespaceName(nli.namespaceRegistry, req)
43 > tlsInfo := authorization.TLSInfoFromContext(ctx)
44 > var serverName string
45 > var certThumbprint string
46 > if tlsInfo != nil {
47 serverName = tlsInfo.State.ServerName
48 cert := authorization.PeerCert(tlsInfo)
51 }
52 }
53 > nli.logger.Debug( namespace_logger.go
54 > "Frontend method invoked.",
55 > tag.WorkflowNamespace(namespace.String()),
56 > tag.Operation(methodName),
57 > tag.ServerName(serverName),
58 > tag.CertThumbprint(certThumbprint))
59 }
60 return handler(ctx, req)
go.temporal.io/server/common/testing/taskpoller/taskpoller.go 7 introduced LOC · 1 range

Open complete file

96 client workflowservice.WorkflowServiceClient,
97 namespace string,
98 > ) *TaskPoller { taskpoller.go
99 > return &TaskPoller{
100 > t: t,
101 > client: client,
102 > namespace: namespace,
103 > }
104 > }
105
106 // PollWorkflowTask creates a workflow task poller that uses the given PollWorkflowTaskQueueRequest.
go.temporal.io/server/common/testing/testhooks/test_impl.go 7 introduced LOC · 3 ranges

Open complete file

65
66 func (h Hook) Scope() ScopeType { return h.scopeType }
67 > func (h Hook) Apply(th TestHooks, scope any) func() { return h.apply(th, scope) } test_impl.go
68
69 > func NewHook[T any, S any](key Key[T, S], value T) Hook { test_impl.go
70 > return Hook{
71 > scopeType: key.scopeType,
72 > apply: func(th TestHooks, scope any) func() {
73 > if _, ok := scope.(S); !ok {
74 panic("testhooks: scope type mismatch")
75 }
76 > return Set(th, key, value, scope) test_impl.go
77 },
78 }
go.temporal.io/server/common/log/tag/tags.go 6 introduced LOC · 2 ranges

Open complete file

539
540 // CertThumbprint returns tag for CertThumbprint
541 > func CertThumbprint(thumbprint string) ZapTag { tags.go
542 > return NewStringTag("cert-thumbprint", thumbprint)
543 > }
544
545 func WorkerComponent(v any) ZapTag {
974
975 // Timeout returns tag for timeout
976 > func Timeout(timeoutValue string) ZapTag { tags.go
977 > return NewStringTag("timeout", timeoutValue)
978 > }
979
980 func DeletedExecutionsCount(count int) ZapTag {
go.temporal.io/server/common/persistence/persistence-tests/setup.go 6 introduced LOC · 2 ranges

Open complete file

32
33 // GetTestClusterOption returns test options for the given store type and driver.
34 > func GetTestClusterOption(storeType, driver string) *TestBaseOptions { setup.go
35 > switch storeType {
36 > case config.StoreTypeSQL:
37 > switch driver {
38 case mysql.PluginName:
39 return GetMySQLTestClusterOption()
40 case postgresql.PluginName, postgresql.PluginNamePGX:
41 return GetPostgreSQLTestClusterOption(driver, nil)
42 > case sqlite.PluginName: setup.go
43 > return GetSQLiteMemoryTestClusterOption()
44 default:
45 panic(fmt.Sprintf("unknown sql driver: %v", driver))
go.temporal.io/server/service/frontend/http_api_server.go 6 introduced LOC · 2 ranges

Open complete file

133 }
134 for _, v := range rpcConfig.HTTPAdditionalForwardedHeaders {
135 > if before, ok := strings.CutSuffix(v, "*"); ok { http_api_server.go
136 > h.matchAdditionalHeaderPrefixes = append(h.matchAdditionalHeaderPrefixes, http.CanonicalHeaderKey(before))
137 > } else {
138 > h.matchAdditionalHeaders[http.CanonicalHeaderKey(v)] = true
139 > }
140 }
141
324 }
325 for _, prefix := range h.matchAdditionalHeaderPrefixes {
326 > if strings.HasPrefix(headerName, prefix) { http_api_server.go
327 return headerName, true
328 }
go.temporal.io/server/tests/testcore/test_env.go 6 introduced LOC · 1 range

Open complete file

622 // checkTestShard supports test sharding based on environment variables.
623 // This distributes tests across multiple CI shards for parallel execution.
624 > func checkTestShard(t *testing.T) { test_env.go
625 > totalStr := os.Getenv("TEST_TOTAL_SHARDS")
626 > indexStr := os.Getenv("TEST_SHARD_INDEX")
627 > if totalStr == "" || indexStr == "" {
628 > return
629 > }
630 total, err := strconv.Atoi(totalStr)
631 if err != nil || total < 1 {
go.temporal.io/server/tests/testcore/context.go 4 introduced LOC · 2 ranges

Open complete file

16 // If a parent context is provided, the returned context will be canceled when
17 // either the timeout expires OR the parent is canceled.
18 > func NewContext(parent ...context.Context) context.Context { context.go
19 > if len(parent) > 0 && parent[0] != nil {
20 // Create RPC context derived from parent
21 ctx, _ := rpc.NewContextFromParentWithTimeoutAndVersionHeaders(
27
28 // Create standalone RPC context
29 > ctx, _ := rpc.NewContextWithTimeoutAndVersionHeaders(testcontext.DefaultTimeout()) context.go
30 > return ctx
31 }