factory.go ×21

Frontier kind: Code frontier

unlabeled · c_38cc31521396

313 tests · 2232 LOC · 103 files · introduces 0 tests · 191 LOC · 11 files

Introduces — evidence that enters the hierarchy at this concept

Code
40 ranges191 lines · 11 files
Tests
0 tests

Contains — complete concept membership

All code (extent)
341 ranges2232 lines · 103 files · Browse complete extent
All tests (intent)
313 testsBrowse complete intent

Neighbourhood graph

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

Introduced files, introduced tests, and structurally relevant concept specialization

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

Graph controls are ready.

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

Native relationship evidence

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

Introduced tests

Every collected test enters the hierarchy at exactly one concept.

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

Introduced code

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

11 files ranked by introduced lines: 191 introduced LOC across 40 ranges. Expand a file to inspect source; the > gutter marks introduced lines.

go.temporal.io/server/common/persistence/client/factory.go 69 introduced LOC · 21 ranges

Open complete file

141
142 // NewShardManager returns a new shard manager
143 > func (f *factoryImpl) NewShardManager() (persistence.ShardManager, error) { factory.go
144 > shardStore, err := f.dataStoreFactory.NewShardStore()
145 > if err != nil {
146 return nil, err
147 }
148
149 > result := persistence.NewShardManager(shardStore, f.serializer) factory.go
150 > if f.systemRateLimiter != nil && f.namespaceRateLimiter != nil {
151 result = persistence.NewShardPersistenceRateLimitedClient(result, f.systemRateLimiter, f.namespaceRateLimiter, f.shardRateLimiter, f.logger)
152 }
153 > if f.metricsHandler != nil && f.healthSignals != nil { factory.go
154 > result = persistence.NewShardPersistenceMetricsClient(result, f.metricsHandler, f.healthSignals, f.logger, f.enableDataLossMetrics)
155 > }
156 > result = persistence.NewShardPersistenceRetryableClient(result, retryPolicy, IsPersistenceTransientError)
157 > return result, nil
158 }
159
160 // NewMetadataManager returns a new metadata manager
161 > func (f *factoryImpl) NewMetadataManager() (persistence.MetadataManager, error) { factory.go
162 > store, err := f.dataStoreFactory.NewMetadataStore()
163 > if err != nil {
164 return nil, err
165 }
166
167 > result := persistence.NewMetadataManagerImpl(store, f.serializer, f.logger, f.clusterName) factory.go
168 > if f.systemRateLimiter != nil && f.namespaceRateLimiter != nil {
169 result = persistence.NewMetadataPersistenceRateLimitedClient(result, f.systemRateLimiter, f.namespaceRateLimiter, f.shardRateLimiter, f.logger)
170 }
171 > if f.metricsHandler != nil && f.healthSignals != nil { factory.go
172 > result = persistence.NewMetadataPersistenceMetricsClient(result, f.metricsHandler, f.healthSignals, f.logger, f.enableDataLossMetrics)
173 > }
174 > result = persistence.NewMetadataPersistenceRetryableClient(result, retryPolicy, IsPersistenceTransientError)
175 > return result, nil
176 }
177
178 // NewClusterMetadataManager returns a new cluster metadata manager
179 > func (f *factoryImpl) NewClusterMetadataManager() (persistence.ClusterMetadataManager, error) { factory.go
180 > store, err := f.dataStoreFactory.NewClusterMetadataStore()
181 > if err != nil {
182 return nil, err
183 }
184
185 > result := persistence.NewClusterMetadataManagerImpl(store, f.serializer, f.clusterName, f.logger) factory.go
186 > if f.systemRateLimiter != nil && f.namespaceRateLimiter != nil {
187 result = persistence.NewClusterMetadataPersistenceRateLimitedClient(result, f.systemRateLimiter, f.namespaceRateLimiter, f.shardRateLimiter, f.logger)
188 }
189 > if f.metricsHandler != nil && f.healthSignals != nil { factory.go
190 > result = persistence.NewClusterMetadataPersistenceMetricsClient(result, f.metricsHandler, f.healthSignals, f.logger, f.enableDataLossMetrics)
191 > }
192 > result = persistence.NewClusterMetadataPersistenceRetryableClient(result, retryPolicy, IsPersistenceTransientError)
193 > return result, nil
194 }
195
196 // NewExecutionManager returns a new execution manager
197 > func (f *factoryImpl) NewExecutionManager() (persistence.ExecutionManager, error) { factory.go
198 > store, err := f.dataStoreFactory.NewExecutionStore()
199 > if err != nil {
200 return nil, err
201 }
202
203 > result := persistence.NewExecutionManager( factory.go
204 > store,
205 > f.serializer,
206 > f.eventBlobCache,
207 > f.logger,
208 > f.config.TransactionSizeLimit,
209 > f.enableBestEffortDeleteTasksOnWorkflowUpdate,
210 > )
211 > if f.systemRateLimiter != nil && f.namespaceRateLimiter != nil {
212 result = persistence.NewExecutionPersistenceRateLimitedClient(result, f.systemRateLimiter, f.namespaceRateLimiter, f.shardRateLimiter, f.logger)
213 }
214 > if f.metricsHandler != nil && f.healthSignals != nil { factory.go
215 > result = persistence.NewExecutionPersistenceMetricsClient(result, f.metricsHandler, f.healthSignals, f.logger, f.enableDataLossMetrics)
216 > }
217 > result = persistence.NewExecutionPersistenceRetryableClient(result, retryPolicy, IsPersistenceTransientError)
218 > return result, nil
219 }
220
221 > func (f *factoryImpl) NewNamespaceReplicationQueue() (persistence.NamespaceReplicationQueue, error) { factory.go
222 > result, err := f.dataStoreFactory.NewQueue(persistence.NamespaceReplicationQueueType)
223 > if err != nil {
224 return nil, err
225 }
226
227 > if f.systemRateLimiter != nil && f.namespaceRateLimiter != nil { factory.go
228 result = persistence.NewQueuePersistenceRateLimitedClient(result, f.systemRateLimiter, f.namespaceRateLimiter, f.shardRateLimiter, f.logger)
229 }
230 > if f.metricsHandler != nil && f.healthSignals != nil { factory.go
231 > result = persistence.NewQueuePersistenceMetricsClient(result, f.metricsHandler, f.healthSignals, f.logger, f.enableDataLossMetrics)
232 > }
233 > result = persistence.NewQueuePersistenceRetryableClient(result, namespaceQueueRetryPolicy, IsNamespaceQueueTransientError)
234 > return persistence.NewNamespaceReplicationQueue(result, f.serializer, f.clusterName, f.metricsHandler, f.logger)
235 }
236
243 }
244
245 > func (f *factoryImpl) NewNexusEndpointManager() (persistence.NexusEndpointManager, error) { factory.go
246 > store, err := f.dataStoreFactory.NewNexusEndpointStore()
247 > if err != nil {
248 return nil, err
249 }
250
251 > result := persistence.NewNexusEndpointManager(store, f.serializer, f.logger) factory.go
252 > if f.systemRateLimiter != nil && f.namespaceRateLimiter != nil {
253 result = persistence.NewNexusEndpointPersistenceRateLimitedClient(result, f.systemRateLimiter, f.namespaceRateLimiter, f.shardRateLimiter, f.logger)
254 }
255 > if f.metricsHandler != nil && f.healthSignals != nil { factory.go
256 > result = persistence.NewNexusEndpointPersistenceMetricsClient(result, f.metricsHandler, f.healthSignals, f.logger, f.enableDataLossMetrics)
257 > }
258 > result = persistence.NewNexusEndpointPersistenceRetryableClient(result, retryPolicy, IsPersistenceTransientError)
259 > return result, nil
260 }
261
292 }
293
294 > if f.metricsHandler == nil { factory.go
295 f.metricsHandler = metrics.NoopMetricsHandler
296 }
297 > if f.healthSignals == nil { factory.go
298 f.healthSignals = persistence.NoopHealthSignalAggregator
299 }
300 > f.healthSignals.Start() factory.go
301 }
go.temporal.io/server/common/persistence/persistence_metric_clients.go 58 introduced LOC · 6 ranges

Open complete file

77
78 // NewShardPersistenceMetricsClient creates a client to manage shards
79 > func NewShardPersistenceMetricsClient(persistence ShardManager, metricsHandler metrics.Handler, healthSignals HealthSignalAggregator, logger log.Logger, enableDataLossMetrics dynamicconfig.BoolPropertyFn) ShardManager { persistence_metric_clients.go
80 > return &shardPersistenceClient{
81 > metricEmitter: metricEmitter{
82 > metricsHandler: metricsHandler,
83 > logger: logger,
84 > enableDataLossMetrics: enableDataLossMetrics,
85 > },
86 > healthSignals: healthSignals,
87 > persistence: persistence,
88 > }
89 > }
90
91 // NewExecutionPersistenceMetricsClient creates a client to manage executions
116
117 // NewMetadataPersistenceMetricsClient creates a MetadataManager client to manage metadata
118 > func NewMetadataPersistenceMetricsClient(persistence MetadataManager, metricsHandler metrics.Handler, healthSignals HealthSignalAggregator, logger log.Logger, enableDataLossMetrics dynamicconfig.BoolPropertyFn) MetadataManager { persistence_metric_clients.go
119 > return &metadataPersistenceClient{
120 > metricEmitter: metricEmitter{
121 > metricsHandler: metricsHandler,
122 > logger: logger,
123 > enableDataLossMetrics: enableDataLossMetrics,
124 > },
125 > healthSignals: healthSignals,
126 > persistence: persistence,
127 > }
128 > }
129
130 // NewClusterMetadataPersistenceMetricsClient creates a ClusterMetadataManager client to manage cluster metadata
131 > func NewClusterMetadataPersistenceMetricsClient(persistence ClusterMetadataManager, metricsHandler metrics.Handler, healthSignals HealthSignalAggregator, logger log.Logger, enableDataLossMetrics dynamicconfig.BoolPropertyFn) ClusterMetadataManager { persistence_metric_clients.go
132 > return &clusterMetadataPersistenceClient{
133 > metricEmitter: metricEmitter{
134 > metricsHandler: metricsHandler,
135 > logger: logger,
136 > enableDataLossMetrics: enableDataLossMetrics,
137 > },
138 > healthSignals: healthSignals,
139 > persistence: persistence,
140 > }
141 > }
142
143 // NewQueuePersistenceMetricsClient creates a client to manage queue
144 > func NewQueuePersistenceMetricsClient(persistence Queue, metricsHandler metrics.Handler, healthSignals HealthSignalAggregator, logger log.Logger, enableDataLossMetrics dynamicconfig.BoolPropertyFn) Queue { persistence_metric_clients.go
145 > return &queuePersistenceClient{
146 > metricEmitter: metricEmitter{
147 > metricsHandler: metricsHandler,
148 > logger: logger,
149 > enableDataLossMetrics: enableDataLossMetrics,
150 > },
151 > healthSignals: healthSignals,
152 > persistence: persistence,
153 > }
154 > }
155
156 // NewNexusEndpointPersistenceMetricsClient creates a NexusEndpointManager to manage nexus endpoints
157 > func NewNexusEndpointPersistenceMetricsClient(persistence NexusEndpointManager, metricsHandler metrics.Handler, healthSignals HealthSignalAggregator, logger log.Logger, enableDataLossMetrics dynamicconfig.BoolPropertyFn) NexusEndpointManager { persistence_metric_clients.go
158 > return &nexusEndpointPersistenceClient{
159 > metricEmitter: metricEmitter{
160 > metricsHandler: metricsHandler,
161 > logger: logger,
162 > enableDataLossMetrics: enableDataLossMetrics,
163 > },
164 > healthSignals: healthSignals,
165 > persistence: persistence,
166 > }
167 > }
168
169 func (p *shardPersistenceClient) GetName() string {
1055 ctx context.Context,
1056 blob *commonpb.DataBlob,
1058 > return p.persistence.Init(ctx, blob)
1059 > }
1060
1061 func (p *queuePersistenceClient) EnqueueMessage(
go.temporal.io/server/common/persistence/namespace_replication_queue.go 16 introduced LOC · 3 ranges

Open complete file

28 metricsHandler metrics.Handler,
29 logger log.Logger,
30 > ) (NamespaceReplicationQueue, error) { namespace_replication_queue.go
31 >
32 > blob, err := serializer.QueueMetadataToBlob(
33 > &persistencespb.QueueMetadata{
34 > ClusterAckLevels: make(map[string]int64),
35 > })
36 > if err != nil {
37 return nil, err
38 }
39 > err = queue.Init(context.TODO(), blob) namespace_replication_queue.go
40 > if err != nil {
41 return nil, err
42 }
43
44 > return &namespaceReplicationQueueImpl{ namespace_replication_queue.go
45 > queue: queue,
46 > clusterName: clusterName,
47 > metricsHandler: metricsHandler,
48 > logger: logger,
49 > serializer: serializer,
50 > }, nil
51 }
52
go.temporal.io/server/api/persistence/v1/queue_metadata.pb.go 8 introduced LOC · 1 range

Open complete file

44 func (*QueueMetadata) ProtoMessage() {}
45
46 > func (x *QueueMetadata) ProtoReflect() protoreflect.Message { queue_metadata.pb.go
47 > mi := &file_temporal_server_api_persistence_v1_queue_metadata_proto_msgTypes[0]
48 > if x != nil {
49 > ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
50 > if ms.LoadMessageInfo() == nil {
51 > ms.StoreMessageInfo(mi)
52 > }
53 > return ms
54 }
55 return mi.MessageOf(x)
go.temporal.io/server/common/persistence/cluster_metadata_store.go 8 introduced LOC · 1 range

Open complete file

37 currentClusterName string,
38 logger log.Logger,
39 > ) ClusterMetadataManager { cluster_metadata_store.go
40 > return &clusterMetadataManagerImpl{
41 > serializer: serializer,
42 > persistence: persistence,
43 > currentClusterName: currentClusterName,
44 > logger: logger,
45 > }
46 > }
47
48 func (m *clusterMetadataManagerImpl) GetName() string {
go.temporal.io/server/common/persistence/metadata_manager.go 8 introduced LOC · 1 range

Open complete file

37 logger log.Logger,
38 clusterName string,
39 > ) MetadataManager { metadata_manager.go
40 > return &metadataManagerImpl{
41 > serializer: serializer,
42 > persistence: persistence,
43 > logger: logger,
44 > clusterName: clusterName,
45 > }
46 > }
47
48 func (m *metadataManagerImpl) GetName() string {
go.temporal.io/server/common/persistence/nexus_endpoint_manager.go 7 introduced LOC · 1 range

Open complete file

32 serializer serialization.Serializer,
33 logger log.Logger,
34 > ) NexusEndpointManager { nexus_endpoint_manager.go
35 > return &nexusEndpointManagerImpl{
36 > persistence: persistence,
37 > serializer: serializer,
38 > logger: logger,
39 > }
40 > }
41
42 func (m *nexusEndpointManagerImpl) GetName() string {
go.temporal.io/server/common/persistence/serialization/codec.go 7 introduced LOC · 2 ranges

Open complete file

79
80 switch encoding {
81 > case enumspb.ENCODING_TYPE_JSON: codec.go
82 > blob, err := codec.NewJSONPBEncoder().Encode(m)
83 > if err != nil {
84 return nil, err
85 }
86 > return &commonpb.DataBlob{ codec.go
87 > Data: blob,
88 > EncodingType: enumspb.ENCODING_TYPE_JSON,
89 > }, nil
90 case enumspb.ENCODING_TYPE_PROTO3:
91 data, err := proto.MarshalOptions{Deterministic: opts.deterministic}.Marshal(m)
go.temporal.io/server/common/persistence/persistence_retryable_clients.go 5 introduced LOC · 2 ranges

Open complete file

1046 ctx context.Context,
1047 blob *commonpb.DataBlob,
1049 > op := func(ctx context.Context) error {
1050 > return p.persistence.Init(ctx, blob)
1051 > }
1052
1053 > return backoff.ThrottleRetryContext(ctx, op, p.policy, p.isRetryable) persistence_retryable_clients.go
1054 }
1055
go.temporal.io/server/common/persistence/serialization/serializer.go 4 introduced LOC · 1 range

Open complete file

562 }
563
564 > func (t *serializerImpl) QueueMetadataToBlob(metadata *persistencespb.QueueMetadata) (*commonpb.DataBlob, error) { serializer.go
565 > // TODO change ENCODING_TYPE_JSON to ENCODING_TYPE_PROTO3
566 > return encodeBlob(metadata, enumspb.ENCODING_TYPE_JSON)
567 > }
568
569 func (t *serializerImpl) QueueMetadataFromBlob(data *commonpb.DataBlob) (*persistencespb.QueueMetadata, error) {
go.temporal.io/server/common/persistence/noop_health_signal_aggregator.go 1 introduced LOC · 1 range

Open complete file

13 func newNoopSignalAggregator() *noopSignalAggregator { return &noopSignalAggregator{} }
14
15 > func (a *noopSignalAggregator) Start() {} noop_health_signal_aggregator.go
16
17 func (a *noopSignalAggregator) Stop() {}