go.temporal.io/server/common/telemetry/grpc.go

197 LOC · 98 covered · 99 uncovered · 24 ranges · 77 concepts · 11 introducers · 17 tests

File neighbourhood

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

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

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

Focused file, its introducer and connector concepts, their introduced files, and tests that run code from the filego.temporal.io/server/common/persistence/client/fx.go · 241 LOCclient/fx.gogo.temporal.io/server/common/persistence/telemetry/cluster_metadata_store_gen.go · 314 LOCtelemetry/cluster_metada…go.temporal.io/server/common/persistence/telemetry/data_store_factory.go · 140 LOCtelemetry/data_store_fac…go.temporal.io/server/common/persistence/telemetry/execution_store_gen.go · 986 LOCtelemetry/execution_stor…go.temporal.io/server/common/persistence/telemetry/nexus_endpoint_store_gen.go · 195 LOCtelemetry/nexus_endpoint…go.temporal.io/server/common/persistence/telemetry/queue_gen.go · 483 LOCtelemetry/queue_gen.gogo.temporal.io/server/common/persistence/telemetry/queue_v2_gen.go · 251 LOCtelemetry/queue_v2_gen.g…go.temporal.io/server/common/persistence/telemetry/shard_store_gen.go · 342 LOCtelemetry/shard_store_ge…go.temporal.io/server/common/persistence/telemetry/shared_store_gen.go · 153 LOCtelemetry/shared_store_g…go.temporal.io/server/common/persistence/telemetry/task_store_gen.go · 566 LOCtelemetry/task_store_gen…go.temporal.io/server/common/resource/fx.go · 533 LOCresource/fx.gogo.temporal.io/server/common/rpc/interceptor/logtags/admin_service_server_gen.go · 231 LOClogtags/admin_service_se…go.temporal.io/server/common/rpc/interceptor/logtags/history_service_server_gen.go · 486 LOClogtags/history_service_…go.temporal.io/server/common/rpc/interceptor/logtags/matching_service_server_gen.go · 198 LOClogtags/matching_service…go.temporal.io/server/common/rpc/interceptor/logtags/workflow_service_server_gen.go · 663 LOClogtags/workflow_service…go.temporal.io/server/common/rpc/interceptor/logtags/workflow_tags.go · 63 LOClogtags/workflow_tags.gogo.temporal.io/server/common/telemetry/config.go · 422 LOCtelemetry/config.gogo.temporal.io/server/service/frontend/fx.go · 1057 LOCfrontend/fx.gogo.temporal.io/server/service/fx.go · 180 LOCservice/fx.gogo.temporal.io/server/service/history/queues/executable.go · 979 LOCqueues/executable.gogo.temporal.io/server/temporal/fx.go · 1273 LOCtemporal/fx.goworkflow_handler.go ×11 · 117 introduced LOCworkflow_handler.go ×11common.go ×1 · 11 introduced LOCcommon.go ×1request_response.pb.go ×6 · 246 introduced LOCrequest_response.pb.go ×…visibility_store.go ×17 · 171 introduced LOCvisibility_store.go ×17telemetry.go ×2 · 15 introduced LOCtelemetry.go ×2fx.go ×1 · 4 introduced LOCfx.go ×1collector.go ×7 · 33 introduced LOCcollector.go ×7data_store_factory.go ×29 · 703 introduced LOCdata_store_factory.go ×2…TestNewServer · 0 introduced LOCTestNewServerpri_matcher.go ×1 · 2 introduced LOCpri_matcher.go ×1timer_queue_active_task_executor.go ×1 · 3 introduced LOCtimer_queue_active_task_…TestNewServer · 0 introduced LOCTestNewServermetric_client.go ×3 · 18 introduced LOCmetric_client.go ×3request_response.pb.go ×12 · 137 introduced LOCrequest_response.pb.go ×…workflow_task_completed_handler.go ×9 · 75 introduced LOCworkflow_task_completed_…pri_forwarder.go ×2 · 4 introduced LOCpri_forwarder.go ×2request_response.pb.go ×1 · 7 introduced LOCrequest_response.pb.go ×…logger.go ×1 · 5 introduced LOClogger.go ×1connections.go ×1 · 9 introduced LOCconnections.go ×1persistence_rate_limited_clients.go ×2 · 17 introduced LOCpersistence_rate_limited…logger.go ×1 · 2 introduced LOClogger.go ×1metric_client.go ×2 · 7 introduced LOCmetric_client.go ×2pri_matcher.go ×1 · 4 introduced LOCpri_matcher.go ×1task_queue_partition_manager.go ×2 · 10 introduced LOCtask_queue_partition_man…connections.go ×2 · 12 introduced LOCconnections.go ×2workflow_handler.go ×8 · 144 introduced LOCworkflow_handler.go ×8pri_matcher.go ×1 · 2 introduced LOCpri_matcher.go ×1handler.go ×1 · 20 introduced LOChandler.go ×1pri_matcher.go ×8 · 55 introduced LOCpri_matcher.go ×8matching_engine.go ×1 · 8 introduced LOCmatching_engine.go ×1server.go ×1 · 3 introduced LOCserver.go ×1logger.go ×2 · 21 introduced LOClogger.go ×2queue_scheduled.go ×1 · 2 introduced LOCqueue_scheduled.go ×1handler.go ×25 · 728 introduced LOChandler.go ×25onebox.go ×75 · 1256 introduced LOConebox.go ×75metric_client_gen.go ×4 · 50 introduced LOCmetric_client_gen.go ×4pri_forwarder.go ×1 · 6 introduced LOCpri_forwarder.go ×1matching_service_server_gen.go ×1 · 2 introduced LOCmatching_service_server_…http_api_server.go ×23 · 195 introduced LOChttp_api_server.go ×23db.go ×1 · 2 introduced LOCdb.go ×1pri_task_writer.go ×3 · 14 introduced LOCpri_task_writer.go ×3namespace_handover.go ×3 · 8 introduced LOCnamespace_handover.go ×3task_queue_partition_manager.go ×2 · 4 introduced LOCtask_queue_partition_man…matching_engine.go ×3 · 5 introduced LOCmatching_engine.go ×3service_grpc.pb.go ×19 · 357 introduced LOCservice_grpc.pb.go ×19request_response.pb.go ×6 · 114 introduced LOCrequest_response.pb.go ×…reader.go ×2 · 8 introduced LOCreader.go ×2service_grpc.pb.go ×20 · 740 introduced LOCservice_grpc.pb.go ×20endpoint_registry.go ×2 · 26 introduced LOCendpoint_registry.go ×2server.go ×3 · 13 introduced LOCserver.go ×3scanner.go ×1 · 2 introduced LOCscanner.go ×1mask_internal_error.go ×1 · 2 introduced LOCmask_internal_error.go ×…matching_engine.go ×2 · 5 introduced LOCmatching_engine.go ×2service_resolver.go ×4 · 49 introduced LOCservice_resolver.go ×4lite_server.go ×25 · 303 introduced LOClite_server.go ×25fx.go ×44 · 705 introduced LOCfx.go ×44adaptive_pool.go ×1 · 3 introduced LOCadaptive_pool.go ×1queue_immediate.go ×1 · 1 introduced LOCqueue_immediate.go ×1fx.go ×1 · 1 introduced LOCfx.go ×1service.go ×8 · 413 introduced LOCservice.go ×8rpc.go ×1 · 13 introduced LOCrpc.go ×1queue_scheduled.go ×1 · 1 introduced LOCqueue_scheduled.go ×1fx.go ×1 · 2 introduced LOCfx.go ×1fx.go ×44 · 4693 introduced LOCfx.go ×44grpc.go ×1 · 7 introduced LOCgrpc.go ×1grpc.go ×1 · 3 introduced LOCgrpc.go ×1annotate_span_with_workflow_tags · 0 introduced LOCannotate_span_with_workf…grpc.go ×2 · 13 introduced LOCgrpc.go ×2workflow_service_server_gen.go ×1 · 2 introduced LOCworkflow_service_server_…grpc.go ×1 · 6 introduced LOCgrpc.go ×1workflow_service_server_gen.go ×1 · 5 introduced LOCworkflow_service_server_…grpc.go ×9 · 44 introduced LOCgrpc.go ×9grpc.go ×1 · 2 introduced LOCgrpc.go ×1grpc.go ×1 · 2 introduced LOCgrpc.go ×1grpc.go ×1 · 2 introduced LOCgrpc.go ×1grpc.go ×1 · 2 introduced LOCgrpc.go ×1grpc.go ×1 · 4 introduced LOCgrpc.go ×1skip_if_noop_trace_provider · introduced test · go.temporal.io/server/common/telemetry/Test_ClientStatsHandler/skip_if_noop_trace_providerskip_if_noop_trace_provi…response_payload_in_debug_mode · introduced test · go.temporal.io/server/common/telemetry/Test_ServerStatsHandler/annotate_span_with_request/response_payload_in_debug_moderesponse_payload_in_debu…annotate_span_with_response_error_payload_in_debug_mode · introduced test · go.temporal.io/server/common/telemetry/Test_ServerStatsHandler/annotate_span_with_response_error_payload_in_debug_modeannotate_span_with_respo…annotate_span_with_workflow_tags · introduced test · go.temporal.io/server/common/telemetry/Test_ServerStatsHandler/annotate_span_with_workflow_tagsannotate_span_with_workf…skip_if_noop_trace_provider · introduced test · go.temporal.io/server/common/telemetry/Test_ServerStatsHandler/skip_if_noop_trace_providerskip_if_noop_trace_provi…TestNewServer · introduced test · go.temporal.io/server/temporal/TestNewServerTestNewServerTestNewServerWithJSONEncoding · introduced test · go.temporal.io/server/temporal/TestNewServerWithJSONEncodingTestNewServerWithJSONEnc…with_OTEL_Collector_running · introduced test · go.temporal.io/server/temporal/TestNewServerWithOTEL/with_OTEL_Collector_runningwith_OTEL_Collector_runn…without_OTEL_Collector_running · introduced test · go.temporal.io/server/temporal/TestNewServerWithOTEL/without_OTEL_Collector_runningwithout_OTEL_Collector_r…ExampleNewServer · introduced test · go.temporal.io/server/temporaltest/ExampleNewServerExampleNewServerTestBaseServerOptions · introduced test · go.temporal.io/server/temporaltest/TestBaseServerOptionsTestBaseServerOptionsTestClientWithCustomInterceptor · introduced test · go.temporal.io/server/temporaltest/TestClientWithCustomInterceptorTestClientWithCustomInte…TestDefaultWorkerOptions · introduced test · go.temporal.io/server/temporaltest/TestDefaultWorkerOptionsTestDefaultWorkerOptionsTestNewServer · introduced test · go.temporal.io/server/temporaltest/TestNewServerTestNewServerTestNewWorkerWithOptions · introduced test · go.temporal.io/server/temporaltest/TestNewWorkerWithOptionsTestNewWorkerWithOptionsTestSearchAttributeRegistration · introduced test · go.temporal.io/server/temporaltest/TestSearchAttributeRegistrationTestSearchAttributeRegis…TestWorkerServiceHealthCheck · introduced test · go.temporal.io/server/tests/testcore/TestFunctionalTestBaseSuite/TestWorkerServiceHealthCheckTestWorkerServiceHealthC…Focused file · go.temporal.io/server/common/telemetry/grpc.go · 197 LOCtelemetry/grpc.go

Graph controls are ready.

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

1 package telemetry
2
3 import (
4 "context"
5 "time"
6
7 "go.opentelemetry.io/contrib/instrumentation/google.golang.org/grpc/otelgrpc"
8 "go.opentelemetry.io/otel/attribute"
9 "go.opentelemetry.io/otel/propagation"
10 "go.opentelemetry.io/otel/trace"
11 otelnoop "go.opentelemetry.io/otel/trace/noop"
12 "go.temporal.io/server/common/log"
13 "go.temporal.io/server/common/log/tag"
14 "go.temporal.io/server/common/rpc/interceptor/logtags"
15 "go.temporal.io/server/common/tasktoken"
16 "google.golang.org/grpc/stats"
17 "google.golang.org/grpc/status"
18 "google.golang.org/protobuf/encoding/protojson"
19 "google.golang.org/protobuf/proto"
20 )
21
22 type methodNameKey struct{}
23
24 type (
25 // ServerStatsHandler gives a named type to the stats.Handler implementation provided by otelgrpc.
26 ServerStatsHandler stats.Handler
27
28 // ClientStatsHandler gives a named type to the grpc.UnaryClientInterceptor implementation provided by otelgrpc.
29 ClientStatsHandler stats.Handler
30
31 customServerStatsHandler struct {
32 isDebug bool
33 wrapped stats.Handler
34 tags *logtags.WorkflowTags
35 }
36 )
37
38 // NewServerStatsHandler creates a new gRPC stats handler that tracks each request with an encapsulating span
39 // using the provided TracerProvider and TextMapPropagator.
40 //
41 // NOTE: If the TracerProvider is `noop.TracerProvider`, it returns `nil`.
42 func NewServerStatsHandler(
43 tp trace.TracerProvider,
44 tmp propagation.TextMapPropagator,
45 logger log.Logger,
46 > ) ServerStatsHandler { grpc.go ×1
47 > if !isEnabled(tp) {
48 > return nil grpc.go ×1
49 > }
50
51 > return newCustomServerStatsHandler( grpc.go ×9
52 > otelgrpc.NewServerHandler(
53 > otelgrpc.WithPropagators(tmp),
54 > otelgrpc.WithTracerProvider(tp),
55 > ),
56 > logger)
57 }
58
59 // NewClientStatsHandler creates a new gRPC stats handler that tracks each request with an encapsulating span
60 // using the provided TracerProvider and TextMapPropagator.
61 //
62 // NOTE: If the TracerProvider is `noop.TracerProvider`, it returns `nil`.
63 func NewClientStatsHandler(
64 tp trace.TracerProvider,
65 tmp propagation.TextMapPropagator,
66 > ) ClientStatsHandler { grpc.go ×1
67 > if !isEnabled(tp) {
68 > return nil grpc.go ×1
69 > }
70
71 > return otelgrpc.NewClientHandler( data_store_factory.go ×29
72 > otelgrpc.WithPropagators(tmp),
73 > otelgrpc.WithTracerProvider(tp),
74 > )
75 }
76
77 func newCustomServerStatsHandler(
78 handler stats.Handler,
79 logger log.Logger,
80 > ) *customServerStatsHandler { grpc.go ×9
81 > return &customServerStatsHandler{
82 > wrapped: handler,
83 > isDebug: DebugMode(),
84 > tags: logtags.NewWorkflowTags(tasktoken.NewSerializer(), logger),
85 > }
86 > }
87
88 > func (c *customServerStatsHandler) TagRPC(ctx context.Context, info *stats.RPCTagInfo) context.Context { grpc.go ×9
89 > return c.wrapped.TagRPC(
90 > context.WithValue(ctx, methodNameKey{}, info.FullMethodName),
91 > info)
92 > }
93
94 > func (c *customServerStatsHandler) HandleRPC(ctx context.Context, stat stats.RPCStats) { grpc.go ×9
95 > // handling `End` before wrapped stats.Handler since it closes the span
96 > switch s := stat.(type) {
97 > case *stats.End:
98 > // annotate with gRPC error payload
99 > if c.isDebug {
100 > span := trace.SpanFromContext(ctx) grpc.go ×2
101 >
102 > //revive:disable-next-line:unchecked-type-assertion
103 > statusErr, ok := status.FromError(s.Error)
104 > if ok && statusErr != nil {
105 > payload, _ := protojson.Marshal(statusErr.Proto()) grpc.go ×1
106 > span.SetAttributes(attribute.Key("rpc.response.error").String(string(payload)))
107 > }
108 }
109 }
110
111 > c.wrapped.HandleRPC(ctx, stat) grpc.go ×9
112 >
113 > switch s := stat.(type) {
114 > case *stats.InHeader: data_store_factory.go ×29
115 > if c.isDebug {
116 span := trace.SpanFromContext(ctx)
117 for key, values := range s.Header {
118 span.SetAttributes(attribute.StringSlice("rpc.request.headers."+key, values))
119 }
120 if deadline, ok := ctx.Deadline(); ok {
121 span.SetAttributes(attribute.String("rpc.request.deadline", deadline.Format(time.RFC3339Nano)))
122 span.SetAttributes(attribute.String("rpc.request.timeout", time.Until(deadline).String()))
123 }
124 }
125 > case *stats.InPayload: grpc.go ×9
126 > span := trace.SpanFromContext(ctx)
127 > c.annotateTags(ctx, span, s.Payload)
128 >
129 > // annotate with gRPC request payload
130 > if c.isDebug {
131 > //revive:disable-next-line:unchecked-type-assertion grpc.go ×2
132 > reqMsg := s.Payload.(proto.Message)
133 > payload, _ := protojson.Marshal(reqMsg)
134 > msgType := string(proto.MessageName(reqMsg).Name())
135 > span.SetAttributes(attribute.Key("rpc.request.payload").String(string(payload)))
136 > span.SetAttributes(attribute.Key("rpc.request.type").String(msgType))
137 > }
138 > case *stats.OutHeader: data_store_factory.go ×29
139 > if c.isDebug {
140 span := trace.SpanFromContext(ctx)
141 for key, values := range s.Header {
142 span.SetAttributes(attribute.StringSlice("rpc.response.headers."+key, values))
143 }
144 }
145 > case *stats.OutPayload: grpc.go ×1
146 > span := trace.SpanFromContext(ctx)
147 > c.annotateTags(ctx, span, s.Payload)
148 >
149 > // annotate with gRPC response payload
150 > if c.isDebug {
151 > //revive:disable-next-line:unchecked-type-assertion grpc.go ×1
152 > respMsg := s.Payload.(proto.Message)
153 > payload, _ := protojson.Marshal(respMsg)
154 > msgType := string(proto.MessageName(respMsg).Name())
155 > span.SetAttributes(attribute.Key("rpc.response.payload").String(string(payload)))
156 > span.SetAttributes(attribute.Key("rpc.response.type").String(msgType))
157 > }
158 }
159 }
160
161 func (c *customServerStatsHandler) annotateTags(
162 ctx context.Context,
163 span trace.Span,
164 payload any,
165 > ) { grpc.go ×9
166 > methodName, ok := ctx.Value(methodNameKey{}).(string)
167 > if !ok {
168 methodName = "unknown"
169 }
170
171 // annotate span with workflow tags (same ones the Temporal SDKs use)
172 > for _, logTag := range c.tags.Extract(payload, methodName) { grpc.go ×9
173 > var k string
174 > switch logTag.Key() {
175 > case tag.WorkflowIDKey:
176 > k = WorkflowIDKey
177 > case tag.WorkflowRunIDKey:
178 > k = RunIDKey
179 default:
180 continue
181 }
182 > span.SetAttributes(attribute.Key(k).String(logTag.Value().(string))) grpc.go ×9
183 }
184 }
185
186 > func (c *customServerStatsHandler) TagConn(ctx context.Context, info *stats.ConnTagInfo) context.Context { data_store_factory.go ×29
187 > return c.wrapped.TagConn(ctx, info)
188 > }
189
190 > func (c *customServerStatsHandler) HandleConn(ctx context.Context, stat stats.ConnStats) { data_store_factory.go ×29
191 > c.wrapped.HandleConn(ctx, stat)
192 > }
193
194 > func isEnabled(tp trace.TracerProvider) bool { grpc.go ×1
195 > _, isNoop := tp.(otelnoop.TracerProvider)
196 > return !isNoop
197 > }