data_store_factory.go ×29

Frontier kind: Code frontier

unlabeled · c_d15dcb68cae1

2 tests · 44881 LOC · 771 files · introduces 0 tests · 703 LOC · 21 files

Introduces — evidence that enters the hierarchy at this concept

Code
177 ranges703 lines · 21 files
Tests
0 tests

Contains — complete concept membership

All code (extent)
10615 ranges44881 lines · 771 files · Browse complete extent
All tests (intent)
2 testsBrowse complete intent

Neighbourhood graph

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

Introduced files, introduced tests, and structurally relevant concept specialization

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

Graph controls are ready.

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

Native relationship evidence

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

Introduced tests

Every collected test enters the hierarchy at exactly one concept.

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

Introduced code

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

Showing the top 20 of 21 files by introduced lines: 701 of 703 introduced LOC and 176 of 177 ranges. Expand a file to inspect source; the > gutter marks introduced lines.

go.temporal.io/server/common/persistence/telemetry/task_store_gen.go 110 introduced LOC · 25 ranges

Open complete file

33 logger log.Logger,
34 tracer trace.Tracer,
35 > ) telemetryTaskStore { task_store_gen.go
36 > return telemetryTaskStore{
37 > TaskStore: base,
38 > tracer: tracer,
39 > debugMode: telemetry.DebugMode(),
40 > }
41 > }
42
43 // CompleteTasksLessThan wraps TaskStore.CompleteTasksLessThan.
126
127 // CreateTaskQueue wraps TaskStore.CreateTaskQueue.
128 > func (d telemetryTaskStore) CreateTaskQueue(ctx context.Context, request *_sourcePersistence.InternalCreateTaskQueueRequest) (err error) { task_store_gen.go
129 > ctx, span := d.tracer.Start(
130 > ctx,
131 > "persistence.TaskStore/CreateTaskQueue",
132 > trace.WithAttributes(
133 > attribute.Key("persistence.store").String("TaskStore"),
134 > attribute.Key("persistence.method").String("CreateTaskQueue"),
135 > ))
136 > defer span.End()
137 >
138 > if deadline, ok := ctx.Deadline(); ok {
139 span.SetAttributes(attribute.String("deadline", deadline.Format(time.RFC3339Nano)))
140 span.SetAttributes(attribute.String("timeout", time.Until(deadline).String()))
141 }
142
143 > err = d.TaskStore.CreateTaskQueue(ctx, request) task_store_gen.go
144 > if err != nil {
145 span.RecordError(err)
146 }
147
148 > if d.debugMode { task_store_gen.go
149
150 requestPayload, err := json.MarshalIndent(request, "", " ")
157 }
158
159 > return task_store_gen.go
160 }
161
162 // CreateTasks wraps TaskStore.CreateTasks.
163 > func (d telemetryTaskStore) CreateTasks(ctx context.Context, request *_sourcePersistence.InternalCreateTasksRequest) (cp1 *_sourcePersistence.CreateTasksResponse, err error) { task_store_gen.go
164 > ctx, span := d.tracer.Start(
165 > ctx,
166 > "persistence.TaskStore/CreateTasks",
167 > trace.WithAttributes(
168 > attribute.Key("persistence.store").String("TaskStore"),
169 > attribute.Key("persistence.method").String("CreateTasks"),
170 > ))
171 > defer span.End()
172 >
173 > if deadline, ok := ctx.Deadline(); ok {
174 span.SetAttributes(attribute.String("deadline", deadline.Format(time.RFC3339Nano)))
175 span.SetAttributes(attribute.String("timeout", time.Until(deadline).String()))
176 }
177
178 > cp1, err = d.TaskStore.CreateTasks(ctx, request) task_store_gen.go
179 > if err != nil {
180 span.RecordError(err)
181 }
182
183 > if d.debugMode { task_store_gen.go
184
185 requestPayload, err := json.MarshalIndent(request, "", " ")
238
239 // GetTaskQueue wraps TaskStore.GetTaskQueue.
240 > func (d telemetryTaskStore) GetTaskQueue(ctx context.Context, request *_sourcePersistence.InternalGetTaskQueueRequest) (ip1 *_sourcePersistence.InternalGetTaskQueueResponse, err error) { task_store_gen.go
241 > ctx, span := d.tracer.Start(
242 > ctx,
243 > "persistence.TaskStore/GetTaskQueue",
244 > trace.WithAttributes(
245 > attribute.Key("persistence.store").String("TaskStore"),
246 > attribute.Key("persistence.method").String("GetTaskQueue"),
247 > ))
248 > defer span.End()
249 >
250 > if deadline, ok := ctx.Deadline(); ok {
251 > span.SetAttributes(attribute.String("deadline", deadline.Format(time.RFC3339Nano)))
252 > span.SetAttributes(attribute.String("timeout", time.Until(deadline).String()))
253 > }
254
255 > ip1, err = d.TaskStore.GetTaskQueue(ctx, request) task_store_gen.go
256 > if err != nil {
257 > span.RecordError(err)
258 > }
259
260 > if d.debugMode { task_store_gen.go
261
262 requestPayload, err := json.MarshalIndent(request, "", " ")
276 }
277
278 > return task_store_gen.go
279 }
280
281 // GetTaskQueueUserData wraps TaskStore.GetTaskQueueUserData.
282 > func (d telemetryTaskStore) GetTaskQueueUserData(ctx context.Context, request *_sourcePersistence.GetTaskQueueUserDataRequest) (ip1 *_sourcePersistence.InternalGetTaskQueueUserDataResponse, err error) { task_store_gen.go
283 > ctx, span := d.tracer.Start(
284 > ctx,
285 > "persistence.TaskStore/GetTaskQueueUserData",
286 > trace.WithAttributes(
287 > attribute.Key("persistence.store").String("TaskStore"),
288 > attribute.Key("persistence.method").String("GetTaskQueueUserData"),
289 > ))
290 > defer span.End()
291 >
292 > if deadline, ok := ctx.Deadline(); ok {
293 span.SetAttributes(attribute.String("deadline", deadline.Format(time.RFC3339Nano)))
294 span.SetAttributes(attribute.String("timeout", time.Until(deadline).String()))
295 }
296
297 > ip1, err = d.TaskStore.GetTaskQueueUserData(ctx, request) task_store_gen.go
298 > if err != nil {
299 > span.RecordError(err)
300 > }
301
302 > if d.debugMode { task_store_gen.go
303
304 requestPayload, err := json.MarshalIndent(request, "", " ")
364
365 // GetTasks wraps TaskStore.GetTasks.
366 > func (d telemetryTaskStore) GetTasks(ctx context.Context, request *_sourcePersistence.GetTasksRequest) (ip1 *_sourcePersistence.InternalGetTasksResponse, err error) { task_store_gen.go
367 > ctx, span := d.tracer.Start(
368 > ctx,
369 > "persistence.TaskStore/GetTasks",
370 > trace.WithAttributes(
371 > attribute.Key("persistence.store").String("TaskStore"),
372 > attribute.Key("persistence.method").String("GetTasks"),
373 > ))
374 > defer span.End()
375 >
376 > if deadline, ok := ctx.Deadline(); ok {
377 > span.SetAttributes(attribute.String("deadline", deadline.Format(time.RFC3339Nano)))
378 > span.SetAttributes(attribute.String("timeout", time.Until(deadline).String()))
379 > }
380
381 > ip1, err = d.TaskStore.GetTasks(ctx, request) task_store_gen.go
382 > if err != nil {
383 span.RecordError(err)
384 }
385
386 > if d.debugMode { task_store_gen.go
387
388 requestPayload, err := json.MarshalIndent(request, "", " ")
490
491 // UpdateTaskQueue wraps TaskStore.UpdateTaskQueue.
492 > func (d telemetryTaskStore) UpdateTaskQueue(ctx context.Context, request *_sourcePersistence.InternalUpdateTaskQueueRequest) (up1 *_sourcePersistence.UpdateTaskQueueResponse, err error) { task_store_gen.go
493 > ctx, span := d.tracer.Start(
494 > ctx,
495 > "persistence.TaskStore/UpdateTaskQueue",
496 > trace.WithAttributes(
497 > attribute.Key("persistence.store").String("TaskStore"),
498 > attribute.Key("persistence.method").String("UpdateTaskQueue"),
499 > ))
500 > defer span.End()
501 >
502 > if deadline, ok := ctx.Deadline(); ok {
503 > span.SetAttributes(attribute.String("deadline", deadline.Format(time.RFC3339Nano)))
504 > span.SetAttributes(attribute.String("timeout", time.Until(deadline).String()))
505 > }
506
507 > up1, err = d.TaskStore.UpdateTaskQueue(ctx, request) task_store_gen.go
508 > if err != nil {
509 span.RecordError(err)
510 }
511
512 > if d.debugMode { task_store_gen.go
513
514 requestPayload, err := json.MarshalIndent(request, "", " ")
go.temporal.io/server/common/persistence/telemetry/execution_store_gen.go 97 introduced LOC · 21 ranges

Open complete file

33 logger log.Logger,
34 tracer trace.Tracer,
35 > ) telemetryExecutionStore { execution_store_gen.go
36 > return telemetryExecutionStore{
37 > ExecutionStore: base,
38 > tracer: tracer,
39 > debugMode: telemetry.DebugMode(),
40 > }
41 > }
42
43 // AddHistoryTasks wraps ExecutionStore.AddHistoryTasks.
182
183 // CreateWorkflowExecution wraps ExecutionStore.CreateWorkflowExecution.
184 > func (d telemetryExecutionStore) CreateWorkflowExecution(ctx context.Context, request *_sourcePersistence.InternalCreateWorkflowExecutionRequest) (ip1 *_sourcePersistence.InternalCreateWorkflowExecutionResponse, err error) { execution_store_gen.go
185 > ctx, span := d.tracer.Start(
186 > ctx,
187 > "persistence.ExecutionStore/CreateWorkflowExecution",
188 > trace.WithAttributes(
189 > attribute.Key("persistence.store").String("ExecutionStore"),
190 > attribute.Key("persistence.method").String("CreateWorkflowExecution"),
191 > ))
192 > defer span.End()
193 >
194 > if deadline, ok := ctx.Deadline(); ok {
195 > span.SetAttributes(attribute.String("deadline", deadline.Format(time.RFC3339Nano)))
196 > span.SetAttributes(attribute.String("timeout", time.Until(deadline).String()))
197 > }
198
199 > ip1, err = d.ExecutionStore.CreateWorkflowExecution(ctx, request) execution_store_gen.go
200 > if err != nil {
201 span.RecordError(err)
202 }
203
204 > if d.debugMode { execution_store_gen.go
205
206 requestPayload, err := json.MarshalIndent(request, "", " ")
518
519 // GetHistoryTasks wraps ExecutionStore.GetHistoryTasks.
520 > func (d telemetryExecutionStore) GetHistoryTasks(ctx context.Context, request *_sourcePersistence.GetHistoryTasksRequest) (ip1 *_sourcePersistence.InternalGetHistoryTasksResponse, err error) { execution_store_gen.go
521 > ctx, span := d.tracer.Start(
522 > ctx,
523 > "persistence.ExecutionStore/GetHistoryTasks",
524 > trace.WithAttributes(
525 > attribute.Key("persistence.store").String("ExecutionStore"),
526 > attribute.Key("persistence.method").String("GetHistoryTasks"),
527 > ))
528 > defer span.End()
529 >
530 > if deadline, ok := ctx.Deadline(); ok {
531 > span.SetAttributes(attribute.String("deadline", deadline.Format(time.RFC3339Nano)))
532 > span.SetAttributes(attribute.String("timeout", time.Until(deadline).String()))
533 > }
534
535 > ip1, err = d.ExecutionStore.GetHistoryTasks(ctx, request) execution_store_gen.go
536 > if err != nil {
537 span.RecordError(err)
538 }
539
540 > if d.debugMode { execution_store_gen.go
541
542 requestPayload, err := json.MarshalIndent(request, "", " ")
644
645 // GetWorkflowExecution wraps ExecutionStore.GetWorkflowExecution.
646 > func (d telemetryExecutionStore) GetWorkflowExecution(ctx context.Context, request *_sourcePersistence.GetWorkflowExecutionRequest) (ip1 *_sourcePersistence.InternalGetWorkflowExecutionResponse, err error) { execution_store_gen.go
647 > ctx, span := d.tracer.Start(
648 > ctx,
649 > "persistence.ExecutionStore/GetWorkflowExecution",
650 > trace.WithAttributes(
651 > attribute.Key("persistence.store").String("ExecutionStore"),
652 > attribute.Key("persistence.method").String("GetWorkflowExecution"),
653 > ))
654 > defer span.End()
655 >
656 > if deadline, ok := ctx.Deadline(); ok {
657 > span.SetAttributes(attribute.String("deadline", deadline.Format(time.RFC3339Nano)))
658 > span.SetAttributes(attribute.String("timeout", time.Until(deadline).String()))
659 > }
660
661 > ip1, err = d.ExecutionStore.GetWorkflowExecution(ctx, request) execution_store_gen.go
662 > if err != nil {
663 span.RecordError(err)
664 }
665
666 > if d.debugMode { execution_store_gen.go
667
668 requestPayload, err := json.MarshalIndent(request, "", " ")
875
876 // ReadHistoryBranch wraps ExecutionStore.ReadHistoryBranch.
877 > func (d telemetryExecutionStore) ReadHistoryBranch(ctx context.Context, request *_sourcePersistence.InternalReadHistoryBranchRequest) (ip1 *_sourcePersistence.InternalReadHistoryBranchResponse, err error) { execution_store_gen.go
878 > ctx, span := d.tracer.Start(
879 > ctx,
880 > "persistence.ExecutionStore/ReadHistoryBranch",
881 > trace.WithAttributes(
882 > attribute.Key("persistence.store").String("ExecutionStore"),
883 > attribute.Key("persistence.method").String("ReadHistoryBranch"),
884 > ))
885 > defer span.End()
886 >
887 > if deadline, ok := ctx.Deadline(); ok {
888 > span.SetAttributes(attribute.String("deadline", deadline.Format(time.RFC3339Nano)))
889 > span.SetAttributes(attribute.String("timeout", time.Until(deadline).String()))
890 > }
891
892 > ip1, err = d.ExecutionStore.ReadHistoryBranch(ctx, request) execution_store_gen.go
893 > if err != nil {
894 span.RecordError(err)
895 }
896
897 > if d.debugMode { execution_store_gen.go
898
899 requestPayload, err := json.MarshalIndent(request, "", " ")
952
953 // UpdateWorkflowExecution wraps ExecutionStore.UpdateWorkflowExecution.
954 > func (d telemetryExecutionStore) UpdateWorkflowExecution(ctx context.Context, request *_sourcePersistence.InternalUpdateWorkflowExecutionRequest) (err error) { execution_store_gen.go
955 > ctx, span := d.tracer.Start(
956 > ctx,
957 > "persistence.ExecutionStore/UpdateWorkflowExecution",
958 > trace.WithAttributes(
959 > attribute.Key("persistence.store").String("ExecutionStore"),
960 > attribute.Key("persistence.method").String("UpdateWorkflowExecution"),
961 > ))
962 > defer span.End()
963 >
964 > if deadline, ok := ctx.Deadline(); ok {
965 > span.SetAttributes(attribute.String("deadline", deadline.Format(time.RFC3339Nano)))
966 > span.SetAttributes(attribute.String("timeout", time.Until(deadline).String()))
967 > }
968
969 > err = d.ExecutionStore.UpdateWorkflowExecution(ctx, request) execution_store_gen.go
970 > if err != nil {
971 span.RecordError(err)
972 }
973
974 > if d.debugMode { execution_store_gen.go
975
976 requestPayload, err := json.MarshalIndent(request, "", " ")
go.temporal.io/server/common/persistence/telemetry/cluster_metadata_store_gen.go 85 introduced LOC · 21 ranges

Open complete file

33 logger log.Logger,
34 tracer trace.Tracer,
35 > ) telemetryClusterMetadataStore { cluster_metadata_store_gen.go
36 > return telemetryClusterMetadataStore{
37 > ClusterMetadataStore: base,
38 > tracer: tracer,
39 > debugMode: telemetry.DebugMode(),
40 > }
41 > }
42
43 // DeleteClusterMetadata wraps ClusterMetadataStore.DeleteClusterMetadata.
77
78 // GetClusterMembers wraps ClusterMetadataStore.GetClusterMembers.
79 > func (d telemetryClusterMetadataStore) GetClusterMembers(ctx context.Context, request *_sourcePersistence.GetClusterMembersRequest) (gp1 *_sourcePersistence.GetClusterMembersResponse, err error) { cluster_metadata_store_gen.go
80 > ctx, span := d.tracer.Start(
81 > ctx,
82 > "persistence.ClusterMetadataStore/GetClusterMembers",
83 > trace.WithAttributes(
84 > attribute.Key("persistence.store").String("ClusterMetadataStore"),
85 > attribute.Key("persistence.method").String("GetClusterMembers"),
86 > ))
87 > defer span.End()
88 >
89 > if deadline, ok := ctx.Deadline(); ok {
90 span.SetAttributes(attribute.String("deadline", deadline.Format(time.RFC3339Nano)))
91 span.SetAttributes(attribute.String("timeout", time.Until(deadline).String()))
92 }
93
94 > gp1, err = d.ClusterMetadataStore.GetClusterMembers(ctx, request) cluster_metadata_store_gen.go
95 > if err != nil {
96 span.RecordError(err)
97 }
98
99 > if d.debugMode { cluster_metadata_store_gen.go
100
101 requestPayload, err := json.MarshalIndent(request, "", " ")
115 }
116
118 }
119
120 // GetClusterMetadata wraps ClusterMetadataStore.GetClusterMetadata.
121 > func (d telemetryClusterMetadataStore) GetClusterMetadata(ctx context.Context, request *_sourcePersistence.InternalGetClusterMetadataRequest) (ip1 *_sourcePersistence.InternalGetClusterMetadataResponse, err error) { cluster_metadata_store_gen.go
122 > ctx, span := d.tracer.Start(
123 > ctx,
124 > "persistence.ClusterMetadataStore/GetClusterMetadata",
125 > trace.WithAttributes(
126 > attribute.Key("persistence.store").String("ClusterMetadataStore"),
127 > attribute.Key("persistence.method").String("GetClusterMetadata"),
128 > ))
129 > defer span.End()
130 >
131 > if deadline, ok := ctx.Deadline(); ok {
132 > span.SetAttributes(attribute.String("deadline", deadline.Format(time.RFC3339Nano)))
133 > span.SetAttributes(attribute.String("timeout", time.Until(deadline).String()))
134 > }
135
136 > ip1, err = d.ClusterMetadataStore.GetClusterMetadata(ctx, request) cluster_metadata_store_gen.go
137 > if err != nil {
138 span.RecordError(err)
139 }
140
141 > if d.debugMode { cluster_metadata_store_gen.go
142
143 requestPayload, err := json.MarshalIndent(request, "", " ")
157 }
158
160 }
161
162 // ListClusterMetadata wraps ClusterMetadataStore.ListClusterMetadata.
163 > func (d telemetryClusterMetadataStore) ListClusterMetadata(ctx context.Context, request *_sourcePersistence.InternalListClusterMetadataRequest) (ip1 *_sourcePersistence.InternalListClusterMetadataResponse, err error) { cluster_metadata_store_gen.go
164 > ctx, span := d.tracer.Start(
165 > ctx,
166 > "persistence.ClusterMetadataStore/ListClusterMetadata",
167 > trace.WithAttributes(
168 > attribute.Key("persistence.store").String("ClusterMetadataStore"),
169 > attribute.Key("persistence.method").String("ListClusterMetadata"),
170 > ))
171 > defer span.End()
172 >
173 > if deadline, ok := ctx.Deadline(); ok {
174 span.SetAttributes(attribute.String("deadline", deadline.Format(time.RFC3339Nano)))
175 span.SetAttributes(attribute.String("timeout", time.Until(deadline).String()))
176 }
177
178 > ip1, err = d.ClusterMetadataStore.ListClusterMetadata(ctx, request) cluster_metadata_store_gen.go
179 > if err != nil {
180 span.RecordError(err)
181 }
182
183 > if d.debugMode { cluster_metadata_store_gen.go
184
185 requestPayload, err := json.MarshalIndent(request, "", " ")
199 }
200
202 }
203
204 // PruneClusterMembership wraps ClusterMetadataStore.PruneClusterMembership.
205 > func (d telemetryClusterMetadataStore) PruneClusterMembership(ctx context.Context, request *_sourcePersistence.PruneClusterMembershipRequest) (err error) { cluster_metadata_store_gen.go
206 > ctx, span := d.tracer.Start(
207 > ctx,
208 > "persistence.ClusterMetadataStore/PruneClusterMembership",
209 > trace.WithAttributes(
210 > attribute.Key("persistence.store").String("ClusterMetadataStore"),
211 > attribute.Key("persistence.method").String("PruneClusterMembership"),
212 > ))
213 > defer span.End()
214 >
215 > if deadline, ok := ctx.Deadline(); ok {
216 span.SetAttributes(attribute.String("deadline", deadline.Format(time.RFC3339Nano)))
217 span.SetAttributes(attribute.String("timeout", time.Until(deadline).String()))
218 }
219
220 > err = d.ClusterMetadataStore.PruneClusterMembership(ctx, request) cluster_metadata_store_gen.go
221 > if err != nil {
222 span.RecordError(err)
223 }
224
225 > if d.debugMode { cluster_metadata_store_gen.go
226
227 requestPayload, err := json.MarshalIndent(request, "", " ")
280
281 // UpsertClusterMembership wraps ClusterMetadataStore.UpsertClusterMembership.
282 > func (d telemetryClusterMetadataStore) UpsertClusterMembership(ctx context.Context, request *_sourcePersistence.UpsertClusterMembershipRequest) (err error) { cluster_metadata_store_gen.go
283 > ctx, span := d.tracer.Start(
284 > ctx,
285 > "persistence.ClusterMetadataStore/UpsertClusterMembership",
286 > trace.WithAttributes(
287 > attribute.Key("persistence.store").String("ClusterMetadataStore"),
288 > attribute.Key("persistence.method").String("UpsertClusterMembership"),
289 > ))
290 > defer span.End()
291 >
292 > if deadline, ok := ctx.Deadline(); ok {
293 span.SetAttributes(attribute.String("deadline", deadline.Format(time.RFC3339Nano)))
294 span.SetAttributes(attribute.String("timeout", time.Until(deadline).String()))
295 }
296
297 > err = d.ClusterMetadataStore.UpsertClusterMembership(ctx, request) cluster_metadata_store_gen.go
298 > if err != nil {
299 span.RecordError(err)
300 }
301
302 > if d.debugMode { cluster_metadata_store_gen.go
303
304 requestPayload, err := json.MarshalIndent(request, "", " ")
go.temporal.io/server/common/persistence/telemetry/data_store_factory.go 64 introduced LOC · 29 ranges

Open complete file

29 logger log.Logger,
30 tracer trace.Tracer,
31 > ) *TelemetryDataStoreFactory { data_store_factory.go
32 > return &TelemetryDataStoreFactory{
33 > baseFactory: baseFactory,
34 > logger: logger,
35 > tracer: tracer,
36 > }
37 > }
38
39 > func (d *TelemetryDataStoreFactory) Close() { data_store_factory.go
40 > d.baseFactory.Close()
41 > }
42
43 > func (d *TelemetryDataStoreFactory) NewTaskStore() (persistence.TaskStore, error) { data_store_factory.go
44 > if d.taskStore == nil {
45 > baseStore, err := d.baseFactory.NewTaskStore()
46 > if err != nil {
47 return nil, err
48 }
49 > d.taskStore = newTelemetryTaskStore(baseStore, d.logger, d.tracer) data_store_factory.go
50 }
51 > return d.taskStore, nil data_store_factory.go
52 }
53
54 > func (d *TelemetryDataStoreFactory) NewFairTaskStore() (persistence.TaskStore, error) { data_store_factory.go
55 > if d.fairTaskStore == nil {
56 > baseStore, err := d.baseFactory.NewFairTaskStore()
57 > if err != nil {
58 return nil, err
59 }
60 > d.fairTaskStore = newTelemetryTaskStore(baseStore, d.logger, d.tracer) data_store_factory.go
61 }
62 > return d.fairTaskStore, nil data_store_factory.go
63 }
64
65 > func (d *TelemetryDataStoreFactory) NewShardStore() (persistence.ShardStore, error) { data_store_factory.go
66 > if d.shardStore == nil {
67 > baseStore, err := d.baseFactory.NewShardStore()
68 > if err != nil {
69 return nil, err
70 }
71 > d.shardStore = newTelemetryShardStore(baseStore, d.logger, d.tracer) data_store_factory.go
72 }
73 > return d.shardStore, nil data_store_factory.go
74 }
75
76 > func (d *TelemetryDataStoreFactory) NewMetadataStore() (persistence.MetadataStore, error) { data_store_factory.go
77 > if d.metadataStore == nil {
78 > baseStore, err := d.baseFactory.NewMetadataStore()
79 > if err != nil {
80 return nil, err
81 }
82 > d.metadataStore = newTelemetryMetadataStore(baseStore, d.logger, d.tracer) data_store_factory.go
83 }
84 > return d.metadataStore, nil data_store_factory.go
85 }
86
87 > func (d *TelemetryDataStoreFactory) NewExecutionStore() (persistence.ExecutionStore, error) { data_store_factory.go
88 > if d.executionStore == nil {
89 > baseStore, err := d.baseFactory.NewExecutionStore()
90 > if err != nil {
91 return nil, err
92 }
93 > d.executionStore = newTelemetryExecutionStore(baseStore, d.logger, d.tracer) data_store_factory.go
94 }
95 > return d.executionStore, nil data_store_factory.go
96 }
97
98 > func (d *TelemetryDataStoreFactory) NewQueue(queueType persistence.QueueType) (persistence.Queue, error) { data_store_factory.go
99 > if d.queue == nil {
100 > baseQueue, err := d.baseFactory.NewQueue(queueType)
101 > if err != nil {
102 return baseQueue, err
103 }
104 > d.queue = newTelemetryQueue(baseQueue, d.logger, d.tracer) data_store_factory.go
105 }
106 > return d.queue, nil data_store_factory.go
107 }
108
109 > func (d *TelemetryDataStoreFactory) NewQueueV2() (persistence.QueueV2, error) { data_store_factory.go
110 > if d.queueV2 == nil {
111 > baseQueue, err := d.baseFactory.NewQueueV2()
112 > if err != nil {
113 return baseQueue, err
114 }
115 > d.queueV2 = newTelemetryQueueV2(baseQueue, d.logger, d.tracer) data_store_factory.go
116 }
117 > return d.queueV2, nil data_store_factory.go
118 }
119
120 > func (d *TelemetryDataStoreFactory) NewClusterMetadataStore() (persistence.ClusterMetadataStore, error) { data_store_factory.go
121 > if d.clusterMDStore == nil {
122 > baseStore, err := d.baseFactory.NewClusterMetadataStore()
123 > if err != nil {
124 return nil, err
125 }
126 > d.clusterMDStore = newTelemetryClusterMetadataStore(baseStore, d.logger, d.tracer) data_store_factory.go
127 }
128 > return d.clusterMDStore, nil data_store_factory.go
129 }
130
131 > func (d *TelemetryDataStoreFactory) NewNexusEndpointStore() (persistence.NexusEndpointStore, error) { data_store_factory.go
132 > if d.nexusEndpointStore == nil {
133 > baseStore, err := d.baseFactory.NewNexusEndpointStore()
134 > if err != nil {
135 return nil, err
136 }
137 > d.nexusEndpointStore = newTelemetryNexusEndpointStore(baseStore, d.logger, d.tracer) data_store_factory.go
138 }
139 > return d.nexusEndpointStore, nil data_store_factory.go
140 }
go.temporal.io/server/common/persistence/telemetry/shard_store_gen.go 60 introduced LOC · 13 ranges

Open complete file

33 logger log.Logger,
34 tracer trace.Tracer,
35 > ) telemetryMetadataStore { shard_store_gen.go
36 > return telemetryMetadataStore{
37 > MetadataStore: base,
38 > tracer: tracer,
39 > debugMode: telemetry.DebugMode(),
40 > }
41 > }
42
43 // CreateNamespace wraps MetadataStore.CreateNamespace.
44 > func (d telemetryMetadataStore) CreateNamespace(ctx context.Context, request *_sourcePersistence.InternalCreateNamespaceRequest) (cp1 *_sourcePersistence.CreateNamespaceResponse, err error) { shard_store_gen.go
45 > ctx, span := d.tracer.Start(
46 > ctx,
47 > "persistence.MetadataStore/CreateNamespace",
48 > trace.WithAttributes(
49 > attribute.Key("persistence.store").String("MetadataStore"),
50 > attribute.Key("persistence.method").String("CreateNamespace"),
51 > ))
52 > defer span.End()
53 >
54 > if deadline, ok := ctx.Deadline(); ok {
55 > span.SetAttributes(attribute.String("deadline", deadline.Format(time.RFC3339Nano)))
56 > span.SetAttributes(attribute.String("timeout", time.Until(deadline).String()))
57 > }
58
59 > cp1, err = d.MetadataStore.CreateNamespace(ctx, request) shard_store_gen.go
60 > if err != nil {
61 span.RecordError(err)
62 }
63
64 > if d.debugMode { shard_store_gen.go
65
66 requestPayload, err := json.MarshalIndent(request, "", " ")
189
190 // GetNamespace wraps MetadataStore.GetNamespace.
191 > func (d telemetryMetadataStore) GetNamespace(ctx context.Context, request *_sourcePersistence.GetNamespaceRequest) (ip1 *_sourcePersistence.InternalGetNamespaceResponse, err error) { shard_store_gen.go
192 > ctx, span := d.tracer.Start(
193 > ctx,
194 > "persistence.MetadataStore/GetNamespace",
195 > trace.WithAttributes(
196 > attribute.Key("persistence.store").String("MetadataStore"),
197 > attribute.Key("persistence.method").String("GetNamespace"),
198 > ))
199 > defer span.End()
200 >
201 > if deadline, ok := ctx.Deadline(); ok {
202 > span.SetAttributes(attribute.String("deadline", deadline.Format(time.RFC3339Nano)))
203 > span.SetAttributes(attribute.String("timeout", time.Until(deadline).String()))
204 > }
205
206 > ip1, err = d.MetadataStore.GetNamespace(ctx, request) shard_store_gen.go
207 > if err != nil {
208 > span.RecordError(err)
209 > }
210
211 > if d.debugMode { shard_store_gen.go
212
213 requestPayload, err := json.MarshalIndent(request, "", " ")
227 }
228
229 > return shard_store_gen.go
230 }
231
232 // ListNamespaces wraps MetadataStore.ListNamespaces.
233 > func (d telemetryMetadataStore) ListNamespaces(ctx context.Context, request *_sourcePersistence.InternalListNamespacesRequest) (ip1 *_sourcePersistence.InternalListNamespacesResponse, err error) { shard_store_gen.go
234 > ctx, span := d.tracer.Start(
235 > ctx,
236 > "persistence.MetadataStore/ListNamespaces",
237 > trace.WithAttributes(
238 > attribute.Key("persistence.store").String("MetadataStore"),
239 > attribute.Key("persistence.method").String("ListNamespaces"),
240 > ))
241 > defer span.End()
242 >
243 > if deadline, ok := ctx.Deadline(); ok {
244 span.SetAttributes(attribute.String("deadline", deadline.Format(time.RFC3339Nano)))
245 span.SetAttributes(attribute.String("timeout", time.Until(deadline).String()))
246 }
247
248 > ip1, err = d.MetadataStore.ListNamespaces(ctx, request) shard_store_gen.go
249 > if err != nil {
250 span.RecordError(err)
251 }
252
253 > if d.debugMode { shard_store_gen.go
254
255 requestPayload, err := json.MarshalIndent(request, "", " ")
go.temporal.io/server/common/persistence/telemetry/shared_store_gen.go 58 introduced LOC · 13 ranges

Open complete file

33 logger log.Logger,
34 tracer trace.Tracer,
35 > ) telemetryShardStore { shared_store_gen.go
36 > return telemetryShardStore{
37 > ShardStore: base,
38 > tracer: tracer,
39 > debugMode: telemetry.DebugMode(),
40 > }
41 > }
42
43 // AssertShardOwnership wraps ShardStore.AssertShardOwnership.
44 > func (d telemetryShardStore) AssertShardOwnership(ctx context.Context, request *_sourcePersistence.AssertShardOwnershipRequest) (err error) { shared_store_gen.go
45 > ctx, span := d.tracer.Start(
46 > ctx,
47 > "persistence.ShardStore/AssertShardOwnership",
48 > trace.WithAttributes(
49 > attribute.Key("persistence.store").String("ShardStore"),
50 > attribute.Key("persistence.method").String("AssertShardOwnership"),
51 > ))
52 > defer span.End()
53 >
54 > if deadline, ok := ctx.Deadline(); ok {
55 span.SetAttributes(attribute.String("deadline", deadline.Format(time.RFC3339Nano)))
56 span.SetAttributes(attribute.String("timeout", time.Until(deadline).String()))
57 }
58
59 > err = d.ShardStore.AssertShardOwnership(ctx, request) shared_store_gen.go
60 > if err != nil {
61 span.RecordError(err)
62 }
63
64 > if d.debugMode { shared_store_gen.go
65
66 requestPayload, err := json.MarshalIndent(request, "", " ")
73 }
74
75 > return shared_store_gen.go
76 }
77
78 // GetOrCreateShard wraps ShardStore.GetOrCreateShard.
79 > func (d telemetryShardStore) GetOrCreateShard(ctx context.Context, request *_sourcePersistence.InternalGetOrCreateShardRequest) (ip1 *_sourcePersistence.InternalGetOrCreateShardResponse, err error) { shared_store_gen.go
80 > ctx, span := d.tracer.Start(
81 > ctx,
82 > "persistence.ShardStore/GetOrCreateShard",
83 > trace.WithAttributes(
84 > attribute.Key("persistence.store").String("ShardStore"),
85 > attribute.Key("persistence.method").String("GetOrCreateShard"),
86 > ))
87 > defer span.End()
88 >
89 > if deadline, ok := ctx.Deadline(); ok {
90 > span.SetAttributes(attribute.String("deadline", deadline.Format(time.RFC3339Nano)))
91 > span.SetAttributes(attribute.String("timeout", time.Until(deadline).String()))
92 > }
93
94 > ip1, err = d.ShardStore.GetOrCreateShard(ctx, request) shared_store_gen.go
95 > if err != nil {
96 span.RecordError(err)
97 }
98
99 > if d.debugMode { shared_store_gen.go
100
101 requestPayload, err := json.MarshalIndent(request, "", " ")
115 }
116
117 > return shared_store_gen.go
118 }
119
120 // UpdateShard wraps ShardStore.UpdateShard.
121 > func (d telemetryShardStore) UpdateShard(ctx context.Context, request *_sourcePersistence.InternalUpdateShardRequest) (err error) { shared_store_gen.go
122 > ctx, span := d.tracer.Start(
123 > ctx,
124 > "persistence.ShardStore/UpdateShard",
125 > trace.WithAttributes(
126 > attribute.Key("persistence.store").String("ShardStore"),
127 > attribute.Key("persistence.method").String("UpdateShard"),
128 > ))
129 > defer span.End()
130 >
131 > if deadline, ok := ctx.Deadline(); ok {
132 > span.SetAttributes(attribute.String("deadline", deadline.Format(time.RFC3339Nano)))
133 > span.SetAttributes(attribute.String("timeout", time.Until(deadline).String()))
134 > }
135
136 > err = d.ShardStore.UpdateShard(ctx, request) shared_store_gen.go
137 > if err != nil {
138 span.RecordError(err)
139 }
140
141 > if d.debugMode { shared_store_gen.go
142
143 requestPayload, err := json.MarshalIndent(request, "", " ")
go.temporal.io/server/common/rpc/interceptor/logtags/workflow_service_server_gen.go 35 introduced LOC · 9 ranges

Open complete file

104 case *workflowservice.DescribeDeploymentResponse:
105 return nil
106 > case *workflowservice.DescribeNamespaceRequest: workflow_service_server_gen.go
107 > return nil
108 > case *workflowservice.DescribeNamespaceResponse:
109 > return nil
110 case *workflowservice.DescribeNexusOperationExecutionRequest:
111 return []tag.Tag{
172 case *workflowservice.GetSearchAttributesResponse:
173 return nil
174 > case *workflowservice.GetSystemInfoRequest: workflow_service_server_gen.go
175 > return nil
176 > case *workflowservice.GetSystemInfoResponse:
177 > return nil
178 case *workflowservice.GetWorkerBuildIdCompatibilityRequest:
179 return nil
193 tag.WorkflowRunID(r.GetExecution().GetRunId()),
194 }
195 > case *workflowservice.GetWorkflowExecutionHistoryResponse: workflow_service_server_gen.go
196 > return nil
197 case *workflowservice.GetWorkflowExecutionHistoryReverseRequest:
198 return []tag.Tag{
297 tag.WorkflowRunID(r.GetRunId()),
298 }
299 > case *workflowservice.PollActivityTaskQueueRequest: workflow_service_server_gen.go
300 > return nil
301 case *workflowservice.PollActivityTaskQueueResponse:
302 return []tag.Tag{
327 tag.WorkflowRunID(r.GetUpdateRef().GetWorkflowExecution().GetRunId()),
328 }
329 > case *workflowservice.PollWorkflowTaskQueueRequest: workflow_service_server_gen.go
330 > return nil
331 > case *workflowservice.PollWorkflowTaskQueueResponse:
332 > return []tag.Tag{
333 > tag.WorkflowID(r.GetWorkflowExecution().GetWorkflowId()),
334 > tag.WorkflowRunID(r.GetWorkflowExecution().GetRunId()),
335 > }
336 case *workflowservice.QueryWorkflowRequest:
337 return []tag.Tag{
357 case *workflowservice.RecordWorkerHeartbeatResponse:
358 return nil
359 > case *workflowservice.RegisterNamespaceRequest: workflow_service_server_gen.go
360 > return nil
361 > case *workflowservice.RegisterNamespaceResponse:
362 > return nil
363 case *workflowservice.RequestCancelActivityExecutionRequest:
364 return []tag.Tag{
461 case *workflowservice.RespondQueryTaskCompletedResponse:
462 return nil
463 > case *workflowservice.RespondWorkflowTaskCompletedRequest: workflow_service_server_gen.go
464 > return wt.fromTaskToken(r.GetTaskToken())
465 > case *workflowservice.RespondWorkflowTaskCompletedResponse:
466 > return nil
467 case *workflowservice.RespondWorkflowTaskFailedRequest:
468 return wt.fromTaskToken(r.GetTaskToken())
489 case *workflowservice.SetWorkerDeploymentRampingVersionResponse:
490 return nil
491 > case *workflowservice.ShutdownWorkerRequest: workflow_service_server_gen.go
492 > return nil
493 > case *workflowservice.ShutdownWorkerResponse:
494 > return nil
495 case *workflowservice.SignalWithStartWorkflowExecutionRequest:
496 return []tag.Tag{
532 tag.WorkflowID(r.GetWorkflowId()),
533 }
534 > case *workflowservice.StartWorkflowExecutionResponse: workflow_service_server_gen.go
535 > return []tag.Tag{
536 > tag.WorkflowRunID(r.GetRunId()),
537 > }
538 case *workflowservice.StopBatchOperationRequest:
539 return nil
go.temporal.io/server/common/rpc/interceptor/logtags/matching_service_server_gen.go 34 introduced LOC · 8 ranges

Open complete file

17 case *matchingservice.AddActivityTaskResponse:
18 return nil
19 > case *matchingservice.AddWorkflowTaskRequest: matching_service_server_gen.go
20 > return []tag.Tag{
21 > tag.WorkflowID(r.GetExecution().GetWorkflowId()),
22 > tag.WorkflowRunID(r.GetExecution().GetRunId()),
23 > }
24 > case *matchingservice.AddWorkflowTaskResponse:
25 > return nil
26 case *matchingservice.ApplyTaskQueueUserDataReplicationEventRequest:
27 return nil
28 case *matchingservice.ApplyTaskQueueUserDataReplicationEventResponse:
29 return nil
30 > case *matchingservice.CancelOutstandingPollRequest: matching_service_server_gen.go
31 > return nil
32 > case *matchingservice.CancelOutstandingPollResponse:
33 > return nil
34 case *matchingservice.CancelOutstandingWorkerPollsRequest:
35 return nil
88 case *matchingservice.ForceUnloadTaskQueueResponse:
89 return nil
90 > case *matchingservice.ForceUnloadTaskQueuePartitionRequest: matching_service_server_gen.go
91 > return nil
92 > case *matchingservice.ForceUnloadTaskQueuePartitionResponse:
93 > return nil
94 case *matchingservice.GetBuildIdTaskQueueMappingRequest:
95 return nil
98 case *matchingservice.GetTaskQueueUserDataRequest:
99 return nil
100 > case *matchingservice.GetTaskQueueUserDataResponse: matching_service_server_gen.go
101 > return nil
102 case *matchingservice.GetWorkerBuildIdCompatibilityRequest:
103 return nil
108 case *matchingservice.GetWorkerVersioningRulesResponse:
109 return nil
110 > case *matchingservice.ListNexusEndpointsRequest: matching_service_server_gen.go
111 > return nil
112 > case *matchingservice.ListNexusEndpointsResponse:
113 > return nil
114 case *matchingservice.ListTaskQueuePartitionsRequest:
115 return nil
120 case *matchingservice.ListWorkersResponse:
121 return nil
122 > case *matchingservice.PollActivityTaskQueueRequest: matching_service_server_gen.go
123 > return nil
124 case *matchingservice.PollActivityTaskQueueResponse:
125 return []tag.Tag{
131 case *matchingservice.PollNexusTaskQueueResponse:
132 return nil
133 > case *matchingservice.PollWorkflowTaskQueueRequest: matching_service_server_gen.go
134 > return nil
135 > case *matchingservice.PollWorkflowTaskQueueResponseWithRawHistory:
136 > return []tag.Tag{
137 > tag.WorkflowID(r.GetWorkflowExecution().GetWorkflowId()),
138 > tag.WorkflowRunID(r.GetWorkflowExecution().GetRunId()),
139 > }
140 case *matchingservice.QueryWorkflowRequest:
141 return []tag.Tag{
145 case *matchingservice.QueryWorkflowResponse:
146 return nil
147 > case *matchingservice.RecordWorkerHeartbeatRequest: matching_service_server_gen.go
148 > return nil
149 > case *matchingservice.RecordWorkerHeartbeatResponse:
150 > return nil
151 case *matchingservice.ReplicateTaskQueueUserDataRequest:
152 return nil
go.temporal.io/server/temporal/fx.go 29 introduced LOC · 7 ranges

Open complete file

959 var tracingReady atomic.Bool
960 otel.SetErrorHandler(otel.ErrorHandlerFunc(func(err error) {
961 > if tracingReady.Load() { // ignore errors during startup fx.go
962 inputs.Logger.Warn("OTEL error", tag.Error(err), tag.ServiceErrorType(err))
963 }
1023 sps := make([]otelsdktrace.SpanProcessor, 0, len(exps))
1024 for _, exp := range exps {
1025 > sps = append(sps, otelsdktrace.NewBatchSpanProcessor(exp, opts...)) fx.go
1026 > }
1027 return sps
1028 },
1056 return telemetry.NoopTracerProvider
1057 }
1058 > opts := make([]otelsdktrace.TracerProviderOption, 0, len(sps)+1) fx.go
1059 > opts = append(opts, otelsdktrace.WithResource(r))
1060 > for _, sp := range sps {
1061 > opts = append(opts, otelsdktrace.WithSpanProcessor(sp))
1062 > }
1063 > tp := otelsdktrace.NewTracerProvider(opts...)
1064 > lc.Append(fx.Hook{
1065 > OnStop: func(ctx context.Context) error {
1066 > otel.SetErrorHandler(otel.ErrorHandlerFunc(func(err error) {
1067 > // ignore errors during shutdown
1068 > }))
1069
1070 > shutdownCtx, cancel := context.WithTimeout(context.Background(), 1*time.Second) fx.go
1071 > defer cancel()
1072 >
1073 > err := tp.Shutdown(shutdownCtx)
1074 > if errors.Is(err, context.DeadlineExceeded) {
1075 > // Ignore timeouts since it's okay to drop OTEL traces on shutdown.
1076 > // Either there's no collector, or there are too many traces left to export.
1077 > return nil
1078 > }
1079 return err
1080 }})
1081 > return tp fx.go
1082 }),
1083 // Haven't had use for baggage propagation yet
1092 return func(ctx context.Context) error {
1093 for _, e := range exporters {
1094 > if starter, ok := e.(starter); ok { fx.go
1095 > err := starter.Start(ctx)
1096 > if err != nil {
1097 return err
1098 }
1109
1110 for _, e := range exporters {
1111 > err := e.Shutdown(shutdownCtx) fx.go
1112 > if errors.Is(err, context.DeadlineExceeded) {
1113 // Ignore timeouts since it's okay to drop OTEL traces on shutdown.
1114 // Either there's no collector, or there are too many traces left to export.
go.temporal.io/server/common/persistence/telemetry/nexus_endpoint_store_gen.go 25 introduced LOC · 5 ranges

Open complete file

33 logger log.Logger,
34 tracer trace.Tracer,
35 > ) telemetryNexusEndpointStore { nexus_endpoint_store_gen.go
36 > return telemetryNexusEndpointStore{
37 > NexusEndpointStore: base,
38 > tracer: tracer,
39 > debugMode: telemetry.DebugMode(),
40 > }
41 > }
42
43 // CreateOrUpdateNexusEndpoint wraps NexusEndpointStore.CreateOrUpdateNexusEndpoint.
154
155 // ListNexusEndpoints wraps NexusEndpointStore.ListNexusEndpoints.
156 > func (d telemetryNexusEndpointStore) ListNexusEndpoints(ctx context.Context, request *_sourcePersistence.ListNexusEndpointsRequest) (ip1 *_sourcePersistence.InternalListNexusEndpointsResponse, err error) { nexus_endpoint_store_gen.go
157 > ctx, span := d.tracer.Start(
158 > ctx,
159 > "persistence.NexusEndpointStore/ListNexusEndpoints",
160 > trace.WithAttributes(
161 > attribute.Key("persistence.store").String("NexusEndpointStore"),
162 > attribute.Key("persistence.method").String("ListNexusEndpoints"),
163 > ))
164 > defer span.End()
165 >
166 > if deadline, ok := ctx.Deadline(); ok {
167 > span.SetAttributes(attribute.String("deadline", deadline.Format(time.RFC3339Nano)))
168 > span.SetAttributes(attribute.String("timeout", time.Until(deadline).String()))
169 > }
170
171 > ip1, err = d.NexusEndpointStore.ListNexusEndpoints(ctx, request) nexus_endpoint_store_gen.go
172 > if err != nil {
173 span.RecordError(err)
174 }
175
176 > if d.debugMode { nexus_endpoint_store_gen.go
177
178 requestPayload, err := json.MarshalIndent(request, "", " ")
go.temporal.io/server/common/rpc/interceptor/logtags/history_service_server_gen.go 24 introduced LOC · 4 ranges

Open complete file

142 case *historyservice.GetShardResponse:
143 return nil
144 > case *historyservice.GetWorkflowExecutionHistoryRequest: history_service_server_gen.go
145 > return []tag.Tag{
146 > tag.WorkflowID(r.GetRequest().GetExecution().GetWorkflowId()),
147 > tag.WorkflowRunID(r.GetRequest().GetExecution().GetRunId()),
148 > }
149 > case *historyservice.GetWorkflowExecutionHistoryResponseWithRaw:
150 > return nil
151 case *historyservice.GetWorkflowExecutionHistoryReverseRequest:
152 return []tag.Tag{
284 case *historyservice.RecordChildExecutionCompletedResponse:
285 return nil
286 > case *historyservice.RecordWorkflowTaskStartedRequest: history_service_server_gen.go
287 > return []tag.Tag{
288 > tag.WorkflowID(r.GetWorkflowExecution().GetWorkflowId()),
289 > tag.WorkflowRunID(r.GetWorkflowExecution().GetRunId()),
290 > }
291 > case *historyservice.RecordWorkflowTaskStartedResponseWithRawHistory:
292 > return nil
293 case *historyservice.RefreshWorkflowTasksRequest:
294 return []tag.Tag{
367 case *historyservice.RespondWorkflowTaskCompletedRequest:
368 return wt.fromTaskToken(r.GetCompleteRequest().GetTaskToken())
369 > case *historyservice.RespondWorkflowTaskCompletedResponse: history_service_server_gen.go
370 > return nil
371 case *historyservice.RespondWorkflowTaskFailedRequest:
372 return wt.fromTaskToken(r.GetFailedRequest().GetTaskToken())
399 case *historyservice.StartNexusOperationResponse:
400 return nil
401 > case *historyservice.StartWorkflowExecutionRequest: history_service_server_gen.go
402 > return []tag.Tag{
403 > tag.WorkflowID(r.GetStartRequest().GetWorkflowId()),
404 > }
405 > case *historyservice.StartWorkflowExecutionResponse:
406 > return []tag.Tag{
407 > tag.WorkflowRunID(r.GetRunId()),
408 > }
409 case *historyservice.SyncActivityRequest:
410 return []tag.Tag{
go.temporal.io/server/service/history/queues/executable.go 23 introduced LOC · 4 ranges

Open complete file

282 // Wrapped in if block to avoid unnecessary allocations when OTEL is disabled.
283 if telemetry.IsEnabled(e.tracer) {
284 > var span trace.Span executable.go
285 >
286 > // Set defaults assuming workflow task
287 > entityID := e.GetWorkflowID()
288 > idKey := telemetry.WorkflowIDKey
289 > taskLabel := e.GetType().String()
290 >
291 > // Override defaults if CHASM task
292 > if _, ok := e.GetTask().(tasks.HasArchetypeID); ok {
293 idKey = telemetry.BusinessIDKey
294 if name := e.GetTask().GetCategory().Name(); name != "" {
297 }
298
299 > ctx, span = e.tracer.Start( executable.go
300 > ctx,
301 > fmt.Sprintf("queue.Execute/%v", taskLabel),
302 > trace.WithSpanKind(trace.SpanKindConsumer),
303 > trace.WithAttributes(
304 > attribute.Key(idKey).String(entityID),
305 > attribute.Key(telemetry.RunIDKey).String(e.GetRunID()),
306 > attribute.Key("queue.task.type").String(e.GetType().String()),
307 > attribute.Key("queue.task.id").Int64(e.GetTaskID())))
308 >
309 > if telemetry.DebugMode() {
310 if taskPayload, err := json.Marshal(e.GetTask()); err != nil {
311 e.logger.Error("failed to serialize task payload for OTEL span", tag.Error(err))
315 }
316
317 > defer func() { executable.go
318 > if retErr != nil {
319 span.RecordError(retErr)
320 }
321 > span.End() executable.go
322 }()
323 }
go.temporal.io/server/common/persistence/telemetry/queue_gen.go 22 introduced LOC · 5 ranges

Open complete file

34 logger log.Logger,
35 tracer trace.Tracer,
36 > ) telemetryQueue { queue_gen.go
37 > return telemetryQueue{
38 > Queue: base,
39 > tracer: tracer,
40 > debugMode: telemetry.DebugMode(),
41 > }
42 > }
43
44 // DeleteMessageFromDLQ wraps Queue.DeleteMessageFromDLQ.
260
261 // Init wraps Queue.Init.
262 > func (d telemetryQueue) Init(ctx context.Context, blob *commonpb.DataBlob) (err error) { queue_gen.go
263 > ctx, span := d.tracer.Start(
264 > ctx,
265 > "persistence.Queue/Init",
266 > trace.WithAttributes(
267 > attribute.Key("persistence.store").String("Queue"),
268 > attribute.Key("persistence.method").String("Init"),
269 > ))
270 > defer span.End()
271 >
272 > if deadline, ok := ctx.Deadline(); ok {
273 span.SetAttributes(attribute.String("deadline", deadline.Format(time.RFC3339Nano)))
274 span.SetAttributes(attribute.String("timeout", time.Until(deadline).String()))
275 }
276
277 > err = d.Queue.Init(ctx, blob) queue_gen.go
278 > if err != nil {
279 span.RecordError(err)
280 }
281
282 > if d.debugMode { queue_gen.go
283
284 requestPayload, err := json.MarshalIndent(blob, "", " ")
go.temporal.io/server/common/telemetry/grpc.go 14 introduced LOC · 5 ranges

Open complete file

69 }
70
71 > return otelgrpc.NewClientHandler( grpc.go
72 > otelgrpc.WithPropagators(tmp),
73 > otelgrpc.WithTracerProvider(tp),
74 > )
75 }
76
112
113 switch s := stat.(type) {
114 > case *stats.InHeader: grpc.go
115 > if c.isDebug {
116 span := trace.SpanFromContext(ctx)
117 for key, values := range s.Header {
136 span.SetAttributes(attribute.Key("rpc.request.type").String(msgType))
137 }
138 > case *stats.OutHeader: grpc.go
139 > if c.isDebug {
140 span := trace.SpanFromContext(ctx)
141 for key, values := range s.Header {
184 }
185
186 > func (c *customServerStatsHandler) TagConn(ctx context.Context, info *stats.ConnTagInfo) context.Context { grpc.go
187 > return c.wrapped.TagConn(ctx, info)
188 > }
189
190 > func (c *customServerStatsHandler) HandleConn(ctx context.Context, stat stats.ConnStats) { grpc.go
191 > c.wrapped.HandleConn(ctx, stat)
192 > }
193
194 func isEnabled(tp trace.TracerProvider) bool {
go.temporal.io/server/common/persistence/telemetry/queue_v2_gen.go 7 introduced LOC · 1 range

Open complete file

33 logger log.Logger,
34 tracer trace.Tracer,
35 > ) telemetryQueueV2 { queue_v2_gen.go
36 > return telemetryQueueV2{
37 > QueueV2: base,
38 > tracer: tracer,
39 > debugMode: telemetry.DebugMode(),
40 > }
41 > }
42
43 // CreateQueue wraps QueueV2.CreateQueue.
go.temporal.io/server/common/rpc/interceptor/logtags/admin_service_server_gen.go 6 introduced LOC · 2 ranges

Open complete file

8 )
9
10 > func (wt *WorkflowTags) extractFromAdminServiceServerMessage(message any) []tag.Tag { admin_service_server_gen.go
11 > switch r := message.(type) {
12 case *adminservice.AddOrUpdateRemoteClusterRequest:
13 return nil
110 case *adminservice.GetShardResponse:
111 return nil
112 > case *adminservice.GetTaskQueueTasksRequest: admin_service_server_gen.go
113 > return nil
114 > case *adminservice.GetTaskQueueTasksResponse:
115 > return nil
116 case *adminservice.GetTaskQueueUserDataRequest:
117 return nil
go.temporal.io/server/common/persistence/client/fx.go 2 introduced LOC · 1 range

Open complete file

214 tracer := tracerProvider.Tracer(otel.ComponentPersistence)
215 if otel.IsEnabled(tracer) {
216 > dataStoreFactory = telemetry.NewTelemetryDataStoreFactory(dataStoreFactory, logger, tracer) fx.go
217 > }
218
219 return dataStoreFactory
go.temporal.io/server/common/resource/fx.go 2 introduced LOC · 1 range

Open complete file

441 var options []grpc.DialOption
442 if tracingStatsHandler != nil {
443 > options = append(options, grpc.WithStatsHandler(tracingStatsHandler)) fx.go
444 > }
445 enableServerKeepalive := dynamicconfig.EnableInternodeServerKeepAlive.Get(dc)()
446 enableClientKeepalive := dynamicconfig.EnableInternodeClientKeepAlive.Get(dc)()
go.temporal.io/server/common/rpc/interceptor/logtags/workflow_tags.go 2 introduced LOC · 1 range

Open complete file

40 // OperatorService doesn't have a single API with workflow tags.
41 return nil
42 > case strings.HasPrefix(fullMethod, api.AdminServicePrefix): workflow_tags.go
43 > return wt.extractFromAdminServiceServerMessage(req)
44 case strings.HasPrefix(fullMethod, api.HistoryServicePrefix):
45 return wt.extractFromHistoryServiceServerMessage(req)
go.temporal.io/server/service/frontend/fx.go 2 introduced LOC · 1 range

Open complete file

329 multiStats := rpc.MultiStatsHandler{}
330 if traceStatsHandler != nil {
331 > multiStats = append(multiStats, traceStatsHandler) fx.go
332 > }
333 if metricsStatsHandler != nil {
334 multiStats = append(multiStats, metricsStatsHandler)