service_grpc.pb.go ×20

Frontier kind: Code frontier

unlabeled · c_a8fd0cd3fe96

11 tests · 23685 LOC · 619 files · introduces 0 tests · 740 LOC · 37 files

Introduces — evidence that enters the hierarchy at this concept

Code
182 ranges740 lines · 37 files
Tests
0 tests

Contains — complete concept membership

All code (extent)
4762 ranges23685 lines · 619 files · Browse complete extent
All tests (intent)
11 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 37 files by introduced lines: 689 of 740 introduced LOC and 160 of 182 ranges. Expand a file to inspect source; the > gutter marks introduced lines.

go.temporal.io/server/api/matchingservice/v1/request_response.pb.go 110 introduced LOC · 19 ranges

Open complete file

54 }
55
56 > func (x *PollWorkflowTaskQueueRequest) Reset() { request_response.pb.go
57 > *x = PollWorkflowTaskQueueRequest{}
58 > mi := &file_temporal_server_api_matchingservice_v1_request_response_proto_msgTypes[0]
59 > ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
60 > ms.StoreMessageInfo(mi)
61 > }
62
63 func (x *PollWorkflowTaskQueueRequest) String() string {
67 func (*PollWorkflowTaskQueueRequest) ProtoMessage() {}
68
69 > func (x *PollWorkflowTaskQueueRequest) ProtoReflect() protoreflect.Message { request_response.pb.go
70 > mi := &file_temporal_server_api_matchingservice_v1_request_response_proto_msgTypes[0]
71 > if x != nil {
72 > ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
73 > if ms.LoadMessageInfo() == nil {
74 > ms.StoreMessageInfo(mi)
75 > }
76 > return ms
77 }
78 return mi.MessageOf(x)
105 }
106
107 > func (x *PollWorkflowTaskQueueRequest) GetForwardedSource() string { request_response.pb.go
108 > if x != nil {
109 > return x.ForwardedSource
110 > }
111 return ""
112 }
151 }
152
153 > func (x *PollWorkflowTaskQueueResponse) Reset() { request_response.pb.go
154 > *x = PollWorkflowTaskQueueResponse{}
155 > mi := &file_temporal_server_api_matchingservice_v1_request_response_proto_msgTypes[1]
156 > ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
157 > ms.StoreMessageInfo(mi)
158 > }
159
160 func (x *PollWorkflowTaskQueueResponse) String() string {
164 func (*PollWorkflowTaskQueueResponse) ProtoMessage() {}
165
166 > func (x *PollWorkflowTaskQueueResponse) ProtoReflect() protoreflect.Message { request_response.pb.go
167 > mi := &file_temporal_server_api_matchingservice_v1_request_response_proto_msgTypes[1]
168 > if x != nil {
169 > ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
170 > if ms.LoadMessageInfo() == nil {
171 > ms.StoreMessageInfo(mi)
172 > }
173 > return ms
174 }
175 return mi.MessageOf(x)
569 }
570
571 > func (x *PollActivityTaskQueueRequest) Reset() { request_response.pb.go
572 > *x = PollActivityTaskQueueRequest{}
573 > mi := &file_temporal_server_api_matchingservice_v1_request_response_proto_msgTypes[3]
574 > ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
575 > ms.StoreMessageInfo(mi)
576 > }
577
578 func (x *PollActivityTaskQueueRequest) String() string {
582 func (*PollActivityTaskQueueRequest) ProtoMessage() {}
583
584 > func (x *PollActivityTaskQueueRequest) ProtoReflect() protoreflect.Message { request_response.pb.go
585 > mi := &file_temporal_server_api_matchingservice_v1_request_response_proto_msgTypes[3]
586 > if x != nil {
587 > ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
588 > if ms.LoadMessageInfo() == nil {
589 > ms.StoreMessageInfo(mi)
590 > }
591 > return ms
592 }
593 return mi.MessageOf(x)
620 }
621
622 > func (x *PollActivityTaskQueueRequest) GetForwardedSource() string { request_response.pb.go
623 > if x != nil {
624 > return x.ForwardedSource
625 > }
626 return ""
627 }
2789 }
2790
2791 > func (x *GetTaskQueueUserDataRequest) Reset() { request_response.pb.go
2792 > *x = GetTaskQueueUserDataRequest{}
2793 > mi := &file_temporal_server_api_matchingservice_v1_request_response_proto_msgTypes[35]
2794 > ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
2795 > ms.StoreMessageInfo(mi)
2796 > }
2797
2798 func (x *GetTaskQueueUserDataRequest) String() string {
2819 }
2820
2821 > func (x *GetTaskQueueUserDataRequest) GetNamespaceId() string { request_response.pb.go
2822 > if x != nil {
2823 > return x.NamespaceId
2824 > }
2825 return ""
2826 }
2833 }
2834
2835 > func (x *GetTaskQueueUserDataRequest) GetTaskQueueType() v19.TaskQueueType { request_response.pb.go
2836 > if x != nil {
2837 > return x.TaskQueueType
2838 > }
2839 return v19.TaskQueueType(0)
2840 }
2878 }
2879
2880 > func (x *GetTaskQueueUserDataResponse) Reset() { request_response.pb.go
2881 > *x = GetTaskQueueUserDataResponse{}
2882 > mi := &file_temporal_server_api_matchingservice_v1_request_response_proto_msgTypes[36]
2883 > ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
2884 > ms.StoreMessageInfo(mi)
2885 > }
2886
2887 func (x *GetTaskQueueUserDataResponse) String() string {
2891 func (*GetTaskQueueUserDataResponse) ProtoMessage() {}
2892
2893 > func (x *GetTaskQueueUserDataResponse) ProtoReflect() protoreflect.Message { request_response.pb.go
2894 > mi := &file_temporal_server_api_matchingservice_v1_request_response_proto_msgTypes[36]
2895 > if x != nil {
2896 > ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
2897 > if ms.LoadMessageInfo() == nil {
2898 > ms.StoreMessageInfo(mi)
2899 > }
2900 > return ms
2901 }
2902 return mi.MessageOf(x)
4790 }
4791
4792 > func (x *ListNexusEndpointsRequest) Reset() { request_response.pb.go
4793 > *x = ListNexusEndpointsRequest{}
4794 > mi := &file_temporal_server_api_matchingservice_v1_request_response_proto_msgTypes[69]
4795 > ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
4796 > ms.StoreMessageInfo(mi)
4797 > }
4798
4799 func (x *ListNexusEndpointsRequest) String() string {
4803 func (*ListNexusEndpointsRequest) ProtoMessage() {}
4804
4805 > func (x *ListNexusEndpointsRequest) ProtoReflect() protoreflect.Message { request_response.pb.go
4806 > mi := &file_temporal_server_api_matchingservice_v1_request_response_proto_msgTypes[69]
4807 > if x != nil {
4808 > ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
4809 > if ms.LoadMessageInfo() == nil {
4810 > ms.StoreMessageInfo(mi)
4811 > }
4812 > return ms
4813 }
4814 return mi.MessageOf(x)
4858 }
4859
4860 > func (x *ListNexusEndpointsResponse) Reset() { request_response.pb.go
4861 > *x = ListNexusEndpointsResponse{}
4862 > mi := &file_temporal_server_api_matchingservice_v1_request_response_proto_msgTypes[70]
4863 > ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
4864 > ms.StoreMessageInfo(mi)
4865 > }
4866
4867 func (x *ListNexusEndpointsResponse) String() string {
4871 func (*ListNexusEndpointsResponse) ProtoMessage() {}
4872
4873 > func (x *ListNexusEndpointsResponse) ProtoReflect() protoreflect.Message { request_response.pb.go
4874 > mi := &file_temporal_server_api_matchingservice_v1_request_response_proto_msgTypes[70]
4875 > if x != nil {
4876 > ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
4877 > if ms.LoadMessageInfo() == nil {
4878 > ms.StoreMessageInfo(mi)
4879 > }
4880 > return ms
4881 }
4882 return mi.MessageOf(x)
5698 func (*PollConditions) ProtoMessage() {}
5699
5700 > func (x *PollConditions) ProtoReflect() protoreflect.Message { request_response.pb.go
5701 > mi := &file_temporal_server_api_matchingservice_v1_request_response_proto_msgTypes[85]
5702 > if x != nil {
5703 ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
5704 if ms.LoadMessageInfo() == nil {
go.temporal.io/server/client/matching/client.go 82 introduced LOC · 17 ranges

Open complete file

233 request *matchingservice.PollActivityTaskQueueRequest,
234 opts ...grpc.CallOption,
235 > ) (*matchingservice.PollActivityTaskQueueResponse, error) { client.go
236 > if !isPartitionAwareKind(request.GetPollRequest().GetTaskQueue().GetKind()) {
237 return c.pollActivityTaskQueue(ctx, PartitionCounts{}, request, opts)
238 }
239 > pkey := c.partitionCache.makeKey( client.go
240 > request.GetNamespaceId(),
241 > request.GetPollRequest().GetTaskQueue().GetName(),
242 > enumspb.TASK_QUEUE_TYPE_ACTIVITY,
243 > )
244 > return invokeWithPartitionCounts(ctx, c.logger, c.partitionCache, pkey, request, opts, c.pollActivityTaskQueue)
245 }
246
250 request *matchingservice.PollActivityTaskQueueRequest,
251 opts []grpc.CallOption,
252 > ) (*matchingservice.PollActivityTaskQueueResponse, error) { client.go
253 > request = common.CloneProto(request)
254 > client, release, err := c.pickClientForRead(
255 > request.GetPollRequest().GetTaskQueue(),
256 > request.GetNamespaceId(),
257 > enumspb.TASK_QUEUE_TYPE_ACTIVITY,
258 > request.GetForwardedSource(),
259 > pc,
260 > )
261 > if err != nil {
262 return nil, err
263 }
264 > if release != nil { client.go
265 > defer release()
266 > }
267 > ctx, cancel := c.createLongPollContext(ctx)
268 > defer cancel()
269 > return client.PollActivityTaskQueue(ctx, request, opts...)
270 }
271
274 request *matchingservice.PollWorkflowTaskQueueRequest,
275 opts ...grpc.CallOption,
276 > ) (*matchingservice.PollWorkflowTaskQueueResponse, error) { client.go
277 > if !isPartitionAwareKind(request.GetPollRequest().GetTaskQueue().GetKind()) {
278 > return c.pollWorkflowTaskQueue(ctx, PartitionCounts{}, request, opts)
279 > }
280 > pkey := c.partitionCache.makeKey(
281 > request.GetNamespaceId(),
282 > request.GetPollRequest().GetTaskQueue().GetName(),
283 > enumspb.TASK_QUEUE_TYPE_WORKFLOW,
284 > )
285 > return invokeWithPartitionCounts(ctx, c.logger, c.partitionCache, pkey, request, opts, c.pollWorkflowTaskQueue)
286 }
287
291 request *matchingservice.PollWorkflowTaskQueueRequest,
292 opts []grpc.CallOption,
293 > ) (*matchingservice.PollWorkflowTaskQueueResponse, error) { client.go
294 > request = common.CloneProto(request)
295 > client, release, err := c.pickClientForRead(
296 > request.GetPollRequest().GetTaskQueue(),
297 > request.GetNamespaceId(),
298 > enumspb.TASK_QUEUE_TYPE_WORKFLOW,
299 > request.GetForwardedSource(),
300 > pc,
301 > )
302 > if err != nil {
303 return nil, err
304 }
305 > if release != nil { client.go
306 > defer release()
307 > }
308 > ctx, cancel := c.createLongPollContext(ctx)
309 > defer cancel()
310 > return client.PollWorkflowTaskQueue(ctx, request, opts...)
311 }
312
444 // processInputPartition returns a partition in certain cases that load balancer involvement is not necessary,
445 // otherwise, returns a task queue to pass down to the load balancer.
446 > func (c *clientImpl) processInputPartition(proto *taskqueuepb.TaskQueue, nsid string, taskType enumspb.TaskQueueType, forwardedFrom string) (tqid.Partition, *tqid.TaskQueue) { client.go
447 > partition, err := tqid.PartitionFromProto(proto, nsid, taskType)
448 > if err != nil {
449 // We preserve the old logic (not returning error in case of invalid proto info) until it's verified that
450 // clients are not sending invalid names.
454 }
455
456 > if forwardedFrom != "" || !partition.IsRoot() { client.go
457 > return partition, nil
458 > }
459
460 > switch p := partition.(type) { client.go
461 > case *tqid.NormalPartition:
462 > return nil, p.TaskQueue()
463 default:
464 return partition, nil
489 forwardedFrom string,
490 pc PartitionCounts,
491 > ) (client matchingservice.MatchingServiceClient, release func(), err error) { client.go
492 > p, tq := c.processInputPartition(proto, nsid, taskType, forwardedFrom)
493 > if tq != nil {
494 > token := c.loadBalancer.PickReadPartition(tq, pc)
495 > p = token.TQPartition
496 > release = token.Release
497 > }
498
499 > proto.Name = p.RpcName() client.go
500 > client, err = c.getClientForTaskQueuePartition(p)
501 > return client, release, err
502 }
503
504 > func (c *clientImpl) createContext(parent context.Context) (context.Context, context.CancelFunc) { client.go
505 > return context.WithTimeout(parent, c.timeout)
506 > }
507
508 > func (c *clientImpl) createLongPollContext(parent context.Context) (context.Context, context.CancelFunc) { client.go
509 > return context.WithTimeout(parent, c.longPollTimeout)
510 > }
511
512 func (c *clientImpl) Route(p tqid.Partition) (string, error) {
523 return nil, err
524 }
525 > client, err := c.clients.GetClientForClientKey(addr) client.go
526 > if err != nil {
527 return nil, err
528 }
529 > return client.(matchingservice.MatchingServiceClient), nil client.go
530 }
531
532 > func isPartitionAwareKind(kind enumspb.TaskQueueKind) bool { client.go
533 > // only normal partitions participate in scaling
534 > return kind == enumspb.TASK_QUEUE_KIND_NORMAL
535 > }
go.temporal.io/server/api/matchingservice/v1/service_grpc.pb.go 70 introduced LOC · 20 ranges

Open complete file

259 }
260
261 > func NewMatchingServiceClient(cc grpc.ClientConnInterface) MatchingServiceClient { service_grpc.pb.go
262 > return &matchingServiceClient{cc}
263 > }
264
265 > func (c *matchingServiceClient) PollWorkflowTaskQueue(ctx context.Context, in *PollWorkflowTaskQueueRequest, opts ...grpc.CallOption) (*PollWorkflowTaskQueueResponse, error) { service_grpc.pb.go
266 > out := new(PollWorkflowTaskQueueResponse)
267 > err := c.cc.Invoke(ctx, MatchingService_PollWorkflowTaskQueue_FullMethodName, in, out, opts...)
268 > if err != nil {
269 return nil, err
270 }
271 > return out, nil service_grpc.pb.go
272 }
273
274 > func (c *matchingServiceClient) PollActivityTaskQueue(ctx context.Context, in *PollActivityTaskQueueRequest, opts ...grpc.CallOption) (*PollActivityTaskQueueResponse, error) { service_grpc.pb.go
275 > out := new(PollActivityTaskQueueResponse)
276 > err := c.cc.Invoke(ctx, MatchingService_PollActivityTaskQueue_FullMethodName, in, out, opts...)
277 > if err != nil {
278 return nil, err
279 }
434 }
435
436 > func (c *matchingServiceClient) GetTaskQueueUserData(ctx context.Context, in *GetTaskQueueUserDataRequest, opts ...grpc.CallOption) (*GetTaskQueueUserDataResponse, error) { service_grpc.pb.go
437 > out := new(GetTaskQueueUserDataResponse)
438 > err := c.cc.Invoke(ctx, MatchingService_GetTaskQueueUserData_FullMethodName, in, out, opts...)
439 > if err != nil {
440 return nil, err
441 }
442 > return out, nil service_grpc.pb.go
443 }
444
569 }
570
571 > func (c *matchingServiceClient) ListNexusEndpoints(ctx context.Context, in *ListNexusEndpointsRequest, opts ...grpc.CallOption) (*ListNexusEndpointsResponse, error) { service_grpc.pb.go
572 > out := new(ListNexusEndpointsResponse)
573 > err := c.cc.Invoke(ctx, MatchingService_ListNexusEndpoints_FullMethodName, in, out, opts...)
574 > if err != nil {
575 return nil, err
576 }
577 > return out, nil service_grpc.pb.go
578 }
579
975 }
976
977 > func _MatchingService_PollWorkflowTaskQueue_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { service_grpc.pb.go
978 > in := new(PollWorkflowTaskQueueRequest)
979 > if err := dec(in); err != nil {
980 return nil, err
981 }
982 > if interceptor == nil { service_grpc.pb.go
983 return srv.(MatchingServiceServer).PollWorkflowTaskQueue(ctx, in)
984 }
985 > info := &grpc.UnaryServerInfo{ service_grpc.pb.go
986 > Server: srv,
987 > FullMethod: MatchingService_PollWorkflowTaskQueue_FullMethodName,
988 > }
989 > handler := func(ctx context.Context, req interface{}) (interface{}, error) {
990 > return srv.(MatchingServiceServer).PollWorkflowTaskQueue(ctx, req.(*PollWorkflowTaskQueueRequest))
991 > }
992 > return interceptor(ctx, in, info, handler)
993 }
994
995 > func _MatchingService_PollActivityTaskQueue_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { service_grpc.pb.go
996 > in := new(PollActivityTaskQueueRequest)
997 > if err := dec(in); err != nil {
998 return nil, err
999 }
1000 > if interceptor == nil { service_grpc.pb.go
1001 return srv.(MatchingServiceServer).PollActivityTaskQueue(ctx, in)
1002 }
1003 > info := &grpc.UnaryServerInfo{ service_grpc.pb.go
1004 > Server: srv,
1005 > FullMethod: MatchingService_PollActivityTaskQueue_FullMethodName,
1006 > }
1007 > handler := func(ctx context.Context, req interface{}) (interface{}, error) {
1008 > return srv.(MatchingServiceServer).PollActivityTaskQueue(ctx, req.(*PollActivityTaskQueueRequest))
1009 > }
1010 > return interceptor(ctx, in, info, handler)
1011 }
1012
1317 }
1318
1319 > func _MatchingService_GetTaskQueueUserData_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { service_grpc.pb.go
1320 > in := new(GetTaskQueueUserDataRequest)
1321 > if err := dec(in); err != nil {
1322 return nil, err
1323 }
1324 > if interceptor == nil { service_grpc.pb.go
1325 return srv.(MatchingServiceServer).GetTaskQueueUserData(ctx, in)
1326 }
1327 > info := &grpc.UnaryServerInfo{ service_grpc.pb.go
1328 > Server: srv,
1329 > FullMethod: MatchingService_GetTaskQueueUserData_FullMethodName,
1330 > }
1331 > handler := func(ctx context.Context, req interface{}) (interface{}, error) {
1332 > return srv.(MatchingServiceServer).GetTaskQueueUserData(ctx, req.(*GetTaskQueueUserDataRequest))
1333 > }
1334 > return interceptor(ctx, in, info, handler)
1335 }
1336
1587 }
1588
1589 > func _MatchingService_ListNexusEndpoints_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { service_grpc.pb.go
1590 > in := new(ListNexusEndpointsRequest)
1591 > if err := dec(in); err != nil {
1592 return nil, err
1593 }
1594 > if interceptor == nil { service_grpc.pb.go
1595 return srv.(MatchingServiceServer).ListNexusEndpoints(ctx, in)
1596 }
1597 > info := &grpc.UnaryServerInfo{ service_grpc.pb.go
1598 > Server: srv,
1599 > FullMethod: MatchingService_ListNexusEndpoints_FullMethodName,
1600 > }
1601 > handler := func(ctx context.Context, req interface{}) (interface{}, error) {
1602 > return srv.(MatchingServiceServer).ListNexusEndpoints(ctx, req.(*ListNexusEndpointsRequest))
1603 > }
1604 > return interceptor(ctx, in, info, handler)
1605 }
1606
go.temporal.io/server/service/frontend/workflow_handler.go 62 introduced LOC · 18 ranges

Open complete file

1069 }
1070
1071 > if len(request.GetIdentity()) > wh.config.MaxIDLengthLimit() { workflow_handler.go
1072 return nil, errIdentityTooLong
1073 }
1074
1075 //nolint:staticcheck // SA1019: worker versioning v0.31
1076 > if err := wh.validateVersioningInfo(request.Namespace, request.WorkerVersionCapabilities, request.DeploymentOptions, request.TaskQueue); err != nil { workflow_handler.go
1077 return nil, err
1078 }
1079
1080 > if request.TaskQueue.GetKind() == enumspb.TASK_QUEUE_KIND_UNSPECIFIED { workflow_handler.go
1081 wh.logger.Warn("Unspecified task queue kind",
1082 tag.WorkflowTaskQueueName(request.TaskQueue.GetName()),
1085 }
1086
1087 > if err := tqid.NormalizeAndValidate(request.TaskQueue, "", wh.config.MaxIDLengthLimit()); err != nil { workflow_handler.go
1088 return nil, err
1089 }
1090
1091 > callTime := time.Now().UTC() workflow_handler.go
1092 >
1093 > namespaceEntry, err := wh.namespaceRegistry.GetNamespace(namespace.Name(request.GetNamespace()))
1094 > if err != nil {
1095 return nil, err
1096 }
1097 > namespaceID := namespaceEntry.ID() workflow_handler.go
1098 >
1099 > wh.logger.Debug("Poll workflow task queue.", tag.WorkflowNamespace(namespaceEntry.Name().String()), tag.WorkflowNamespaceID(namespaceID.String()))
1100 > if err := wh.checkBadBinary(namespaceEntry, request.GetBinaryChecksum()); err != nil {
1101 return nil, err
1102 }
1103
1104 > if contextNearDeadline(ctx, longPollTailRoom) { workflow_handler.go
1105 return &workflowservice.PollWorkflowTaskQueueResponse{}, nil
1106 }
1107
1108 > pollerID := uuid.NewString() workflow_handler.go
1109 > childCtx := wh.registerOutstandingPollContext(ctx, pollerID, namespaceID.String())
1110 > defer wh.unregisterOutstandingPollContext(pollerID, namespaceID.String())
1111 >
1112 > matchingResp, err := wh.matchingClient.PollWorkflowTaskQueue(childCtx, &matchingservice.PollWorkflowTaskQueueRequest{
1113 > NamespaceId: namespaceID.String(),
1114 > PollerId: pollerID,
1115 > PollRequest: request,
1116 > })
1117 > if err != nil {
1118 contextWasCanceled := wh.cancelOutstandingPoll(childCtx, namespaceID, enumspb.TASK_QUEUE_TYPE_WORKFLOW, request.TaskQueue, pollerID)
1119 if contextWasCanceled {
1150 // through to RawHistory field. The matching client auto-deserializes the repeated bytes into
1151 // a History message via gRPC wire compatibility.
1152 > history := matchingResp.History workflow_handler.go
1153 > if matchingResp.RawHistory != nil {
1154 history = matchingResp.RawHistory
1155 // Process search attributes for raw history since it bypasses the normal processing path.
1165 }
1166
1167 > return &workflowservice.PollWorkflowTaskQueueResponse{ workflow_handler.go
1168 > TaskToken: matchingResp.TaskToken,
1169 > WorkflowExecution: matchingResp.WorkflowExecution,
1170 > WorkflowType: matchingResp.WorkflowType,
1171 > PreviousStartedEventId: matchingResp.PreviousStartedEventId,
1172 > StartedEventId: matchingResp.StartedEventId,
1173 > Query: matchingResp.Query,
1174 > BacklogCountHint: matchingResp.BacklogCountHint,
1175 > Attempt: matchingResp.Attempt,
1176 > History: history,
1177 > NextPageToken: matchingResp.NextPageToken,
1178 > WorkflowExecutionTaskQueue: matchingResp.WorkflowExecutionTaskQueue,
1179 > ScheduledTime: matchingResp.ScheduledTime,
1180 > StartedTime: matchingResp.StartedTime,
1181 > Queries: matchingResp.Queries,
1182 > Messages: matchingResp.Messages,
1183 > PollerScalingDecision: matchingResp.PollerScalingDecision,
1184 > }, nil
1185 }
1186
1339 }
1340
1341 > namespaceName := namespace.Name(request.GetNamespace()) workflow_handler.go
1342 > if err := tqid.NormalizeAndValidate(request.TaskQueue, "", wh.config.MaxIDLengthLimit()); err != nil {
1343 return nil, err
1344 }
1345 > if len(request.GetIdentity()) > wh.config.MaxIDLengthLimit() { workflow_handler.go
1346 return nil, errIdentityTooLong
1347 }
1348
1349 //nolint:staticcheck // SA1019: worker versioning v0.31
1350 > if err := wh.validateVersioningInfo(request.Namespace, request.WorkerVersionCapabilities, request.DeploymentOptions, request.TaskQueue); err != nil { workflow_handler.go
1351 return nil, err
1352 }
1353
1354 > namespaceID, err := wh.namespaceRegistry.GetNamespaceID(namespaceName) workflow_handler.go
1355 > if err != nil {
1356 return nil, err
1357 }
1358
1359 > if contextNearDeadline(ctx, longPollTailRoom) { workflow_handler.go
1360 return &workflowservice.PollActivityTaskQueueResponse{}, nil
1361 }
1362
1363 > pollerID := uuid.NewString() workflow_handler.go
1364 > childCtx := wh.registerOutstandingPollContext(ctx, pollerID, namespaceID.String())
1365 > defer wh.unregisterOutstandingPollContext(pollerID, namespaceID.String())
1366 > matchingResponse, err := wh.matchingClient.PollActivityTaskQueue(childCtx, &matchingservice.PollActivityTaskQueueRequest{
1367 > NamespaceId: namespaceID.String(),
1368 > PollerId: pollerID,
1369 > PollRequest: request,
1370 > })
1371 > if err != nil {
1372 contextWasCanceled := wh.cancelOutstandingPoll(childCtx, namespaceID, enumspb.TASK_QUEUE_TYPE_ACTIVITY, request.TaskQueue, pollerID)
1373 if contextWasCanceled {
6764 }
6765
6766 > func (wh *WorkflowHandler) checkBadBinary(namespaceEntry *namespace.Namespace, binaryChecksum string) error { workflow_handler.go
6767 > if err := namespaceEntry.VerifyBinaryChecksum(binaryChecksum); err != nil {
6768 return serviceerror.NewInvalidArgumentf("Binary %v already marked as bad deployment.", binaryChecksum)
6769 }
6770 > return nil workflow_handler.go
6771 }
6772
go.temporal.io/server/service/matching/nexus_endpoint_client.go 53 introduced LOC · 12 ranges

Open complete file

243 ctx context.Context,
244 request *matchingservice.ListNexusEndpointsRequest,
245 > ) (*matchingservice.ListNexusEndpointsResponse, chan struct{}, error) { nexus_endpoint_client.go
246 > m.RLock()
247 > if request.LastKnownTableVersion > m.tableVersion {
248 // indicates we may have lost table ownership, so need to reload from persistence
249 m.hasLoadedEndpoints.Store(false)
250 }
251 > m.RUnlock() nexus_endpoint_client.go
252 >
253 > if !m.hasLoadedEndpoints.Load() {
254 > if err := m.loadEndpoints(ctx); err != nil {
255 return nil, nil, fmt.Errorf("error loading nexus endpoints cache: %w", err)
256 }
257 }
258
259 > m.RLock() nexus_endpoint_client.go
260 > defer m.RUnlock()
261 >
262 > if request.LastKnownTableVersion != 0 && request.LastKnownTableVersion != m.tableVersion {
263 return nil, nil, serviceerror.NewFailedPreconditionf("nexus endpoints table version mismatch. received: %v expected %v", request.LastKnownTableVersion, m.tableVersion)
264 }
265
266 > startIdx := 0 nexus_endpoint_client.go
267 > if request.NextPageToken != nil {
268 nextEndpointID := string(request.NextPageToken)
269
281 }
282
283 > endIdx := min(startIdx+int(request.PageSize), len(m.endpointEntries)) nexus_endpoint_client.go
284 >
285 > var nextPageToken []byte
286 > if endIdx < len(m.endpointEntries) {
287 nextPageToken = []byte(m.endpointEntries[endIdx].Id)
288 }
289
290 > resp := &matchingservice.ListNexusEndpointsResponse{ nexus_endpoint_client.go
291 > TableVersion: m.tableVersion,
292 > NextPageToken: nextPageToken,
293 > Entries: slices.Clone(m.endpointEntries[startIdx:endIdx]),
294 > }
295 >
296 > return resp, m.tableVersionChanged, nil
297 }
298
299 > func (m *nexusEndpointClient) loadEndpoints(ctx context.Context) error { nexus_endpoint_client.go
300 > m.Lock()
301 > defer m.Unlock()
302 >
303 > if m.hasLoadedEndpoints.Load() {
304 // check whether endpoints were loaded while waiting for write lock
305 return nil
307
308 // reset cached view since we will be paging from the start
309 > m.resetCacheStateLocked() nexus_endpoint_client.go
310 >
311 > var pageToken []byte
312 >
313 > for ctx.Err() == nil {
314 > resp, err := m.persistence.ListNexusEndpoints(ctx, &p.ListNexusEndpointsRequest{
315 > LastKnownTableVersion: m.tableVersion,
316 > NextPageToken: pageToken,
317 > PageSize: loadEndpointsPageSize,
318 > })
319 > if err != nil {
320 if errors.Is(err, p.ErrNexusTableVersionConflict) {
321 // indicates table was updated during paging, so reset and start from the beginning
327 }
328
329 > pageToken = resp.NextPageToken nexus_endpoint_client.go
330 > m.tableVersion = resp.TableVersion
331 > for _, entry := range resp.Entries {
332 m.endpointEntries = append(m.endpointEntries, entry)
333 m.endpointsByID[entry.Id] = entry
335 }
336
337 > if len(pageToken) == 0 { nexus_endpoint_client.go
338 > break
339 }
340 }
341
342 > m.hasLoadedEndpoints.Store(ctx.Err() == nil) nexus_endpoint_client.go
343 > return ctx.Err()
344 }
345
346 > func (m *nexusEndpointClient) resetCacheStateLocked() { nexus_endpoint_client.go
347 > m.tableVersion = 0
348 > m.endpointEntries = []*persistencespb.NexusEndpointEntry{}
349 > m.endpointsByID = make(map[string]*persistencespb.NexusEndpointEntry)
350 > m.endpointsByName = make(map[string]*persistencespb.NexusEndpointEntry)
351 > }
352
353 // notifyOwnershipChanged starts or stops a background routine which watches the Nexus endpoints table version for
go.temporal.io/server/service/matching/handler.go 41 introduced LOC · 9 ranges

Open complete file

227 ctx context.Context,
228 request *matchingservice.PollActivityTaskQueueRequest,
229 > ) (_ *matchingservice.PollActivityTaskQueueResponse, retError error) { handler.go
230 > defer log.CapturePanic(h.logger, &retError)
231 > opMetrics := h.opMetricsHandler(
232 > request.GetNamespaceId(),
233 > request.GetPollRequest().GetTaskQueue(),
234 > enumspb.TASK_QUEUE_TYPE_ACTIVITY,
235 > metrics.MatchingPollActivityTaskQueueScope,
236 > )
237 >
238 > if request.GetForwardedSource() != "" {
239 h.reportForwardedPerTaskQueueCounter(opMetrics, namespace.ID(request.GetNamespaceId()))
240 }
241
242 > if _, err := common.ValidateLongPollContextTimeoutIsSet( handler.go
243 > ctx,
244 > "PollActivityTaskQueue",
245 > h.throttledLogger,
246 > ); err != nil {
247 return nil, err
248 }
249
250 > return h.engine.PollActivityTaskQueue(ctx, request, opMetrics) handler.go
251 }
252
255 ctx context.Context,
256 request *matchingservice.PollWorkflowTaskQueueRequest,
257 > ) (_ *matchingservice.PollWorkflowTaskQueueResponseWithRawHistory, retError error) { handler.go
258 > defer log.CapturePanic(h.logger, &retError)
259 > opMetrics := h.opMetricsHandler(
260 > request.GetNamespaceId(),
261 > request.GetPollRequest().GetTaskQueue(),
262 > enumspb.TASK_QUEUE_TYPE_WORKFLOW,
263 > metrics.MatchingPollWorkflowTaskQueueScope,
264 > )
265 >
266 > if request.GetForwardedSource() != "" {
267 h.reportForwardedPerTaskQueueCounter(opMetrics, namespace.ID(request.GetNamespaceId()))
268 }
269
270 > if _, err := common.ValidateLongPollContextTimeoutIsSet( handler.go
271 > ctx,
272 > "PollWorkflowTaskQueue",
273 > h.throttledLogger,
274 > ); err != nil {
275 return nil, err
276 }
277
278 > return h.engine.PollWorkflowTaskQueue(ctx, request, opMetrics) handler.go
279 }
280
434 ctx context.Context,
435 request *matchingservice.GetTaskQueueUserDataRequest,
436 > ) (_ *matchingservice.GetTaskQueueUserDataResponse, retError error) { handler.go
437 > defer log.CapturePanic(h.logger, &retError)
438 > return h.engine.GetTaskQueueUserData(ctx, request)
439 > }
440
441 func (h *Handler) SyncDeploymentUserData(
593 }
594
595 > func (h *Handler) ListNexusEndpoints(ctx context.Context, request *matchingservice.ListNexusEndpointsRequest) (_ *matchingservice.ListNexusEndpointsResponse, retError error) { handler.go
596 > defer log.CapturePanic(h.logger, &retError)
597 > return h.engine.ListNexusEndpoints(ctx, request)
598 > }
599
600 // RecordWorkerHeartbeat receive heartbeat request from the worker.
689 return ""
690 }
691 > return entry.Name() handler.go
692 }
693
go.temporal.io/server/common/persistence/persistence_metric_clients.go 36 introduced LOC · 4 ranges

Open complete file

612 ctx context.Context,
613 request *GetTasksRequest,
614 > ) (_ *GetTasksResponse, retErr error) { persistence_metric_clients.go
615 > caller := headers.GetCallerInfo(ctx).CallerName
616 > startTime := time.Now().UTC()
617 > defer func() {
618 > p.healthSignals.Record(CallerSegmentMissing, time.Since(startTime), retErr)
619 > p.recordRequestMetrics(metrics.PersistenceGetTasksScope, caller, time.Since(startTime), retErr)
620 > p.recordDataLossMetrics(metrics.PersistenceGetTasksScope, caller, retErr, "", "")
621 > }()
622 > return p.persistence.GetTasks(ctx, request)
623 }
624
640 ctx context.Context,
641 request *CreateTaskQueueRequest,
642 > ) (_ *CreateTaskQueueResponse, retErr error) { persistence_metric_clients.go
643 > caller := headers.GetCallerInfo(ctx).CallerName
644 > startTime := time.Now().UTC()
645 > defer func() {
646 > p.healthSignals.Record(CallerSegmentMissing, time.Since(startTime), retErr)
647 > p.recordRequestMetrics(metrics.PersistenceCreateTaskQueueScope, caller, time.Since(startTime), retErr)
648 > p.recordDataLossMetrics(metrics.PersistenceCreateTaskQueueScope, caller, retErr, "", "")
649 > }()
650 > return p.persistence.CreateTaskQueue(ctx, request)
651 }
652
668 ctx context.Context,
669 request *GetTaskQueueRequest,
670 > ) (_ *GetTaskQueueResponse, retErr error) { persistence_metric_clients.go
671 > caller := headers.GetCallerInfo(ctx).CallerName
672 > startTime := time.Now().UTC()
673 > defer func() {
674 > p.healthSignals.Record(CallerSegmentMissing, time.Since(startTime), retErr)
675 > p.recordRequestMetrics(metrics.PersistenceGetTaskQueueScope, caller, time.Since(startTime), retErr)
676 > p.recordDataLossMetrics(metrics.PersistenceGetTaskQueueScope, caller, retErr, "", "")
677 > }()
678 > return p.persistence.GetTaskQueue(ctx, request)
679 }
680
710 ctx context.Context,
711 request *GetTaskQueueUserDataRequest,
712 > ) (_ *GetTaskQueueUserDataResponse, retErr error) { persistence_metric_clients.go
713 > caller := headers.GetCallerInfo(ctx).CallerName
714 > startTime := time.Now().UTC()
715 > defer func() {
716 > p.healthSignals.Record(CallerSegmentMissing, time.Since(startTime), retErr)
717 > p.recordRequestMetrics(metrics.PersistenceGetTaskQueueUserDataScope, caller, time.Since(startTime), retErr)
718 > p.recordDataLossMetrics(metrics.PersistenceGetTaskQueueUserDataScope, caller, retErr, "", "")
719 > }()
720 > return p.persistence.GetTaskQueueUserData(ctx, request)
721 }
722
go.temporal.io/server/common/persistence/persistence_retryable_clients.go 36 introduced LOC · 8 ranges

Open complete file

624 ctx context.Context,
625 request *GetTasksRequest,
626 > ) (*GetTasksResponse, error) { persistence_retryable_clients.go
627 > var response *GetTasksResponse
628 > op := func(ctx context.Context) error {
629 > var err error
630 > response, err = p.persistence.GetTasks(ctx, request)
631 > return err
632 > }
633
634 > err := backoff.ThrottleRetryContext(ctx, op, p.policy, p.isRetryable) persistence_retryable_clients.go
635 > return response, err
636 }
637
654 ctx context.Context,
655 request *CreateTaskQueueRequest,
656 > ) (*CreateTaskQueueResponse, error) { persistence_retryable_clients.go
657 > var response *CreateTaskQueueResponse
658 > op := func(ctx context.Context) error {
659 > var err error
660 > response, err = p.persistence.CreateTaskQueue(ctx, request)
661 > return err
662 > }
663
664 > err := backoff.ThrottleRetryContext(ctx, op, p.policy, p.isRetryable) persistence_retryable_clients.go
665 > return response, err
666 }
667
684 ctx context.Context,
685 request *GetTaskQueueRequest,
686 > ) (*GetTaskQueueResponse, error) { persistence_retryable_clients.go
687 > var response *GetTaskQueueResponse
688 > op := func(ctx context.Context) error {
689 > var err error
690 > response, err = p.persistence.GetTaskQueue(ctx, request)
691 > return err
692 > }
693
694 > err := backoff.ThrottleRetryContext(ctx, op, p.policy, p.isRetryable) persistence_retryable_clients.go
695 > return response, err
696 }
697
725 ctx context.Context,
726 request *GetTaskQueueUserDataRequest,
727 > ) (*GetTaskQueueUserDataResponse, error) { persistence_retryable_clients.go
728 > var response *GetTaskQueueUserDataResponse
729 > op := func(ctx context.Context) error {
730 > var err error
731 > response, err = p.persistence.GetTaskQueueUserData(ctx, request)
732 > return err
733 > }
734
735 > err := backoff.ThrottleRetryContext(ctx, op, p.policy, p.isRetryable) persistence_retryable_clients.go
736 > return response, err
737 }
738
go.temporal.io/server/client/matching/metric_client.go 34 introduced LOC · 9 ranges

Open complete file

86 request *matchingservice.PollActivityTaskQueueRequest,
87 opts ...grpc.CallOption,
88 > ) (_ *matchingservice.PollActivityTaskQueueResponse, retError error) { metric_client.go
89 >
90 > scope, stopwatch := c.startMetricsRecording(ctx, metrics.MatchingClientPollActivityTaskQueueScope)
91 > defer func() {
92 c.finishMetricsRecording(scope, stopwatch, retError)
93 }()
94
95 > if request.PollRequest != nil { metric_client.go
96 > c.emitForwardedSourceStats(
97 > scope,
98 > request.GetForwardedSource(),
99 > request.PollRequest.TaskQueue,
100 > )
101 > }
102
103 > return c.client.PollActivityTaskQueue(ctx, request, opts...) metric_client.go
104 }
105
108 request *matchingservice.PollWorkflowTaskQueueRequest,
109 opts ...grpc.CallOption,
110 > ) (_ *matchingservice.PollWorkflowTaskQueueResponse, retError error) { metric_client.go
111 >
112 > scope, stopwatch := c.startMetricsRecording(ctx, metrics.MatchingClientPollWorkflowTaskQueueScope)
113 > defer func() {
114 > c.finishMetricsRecording(scope, stopwatch, retError)
115 > }()
116
117 > if request.PollRequest != nil { metric_client.go
118 > c.emitForwardedSourceStats(
119 > scope,
120 > request.GetForwardedSource(),
121 > request.PollRequest.TaskQueue,
122 > )
123 > }
124
125 > return c.client.PollWorkflowTaskQueue(ctx, request, opts...) metric_client.go
126 }
127
190 forwardedFrom string,
191 taskQueue *taskqueuepb.TaskQueue,
192 > ) { metric_client.go
193 > if taskQueue == nil {
194 return
195 }
196
197 > switch { metric_client.go
198 case forwardedFrom != "":
199 metrics.MatchingClientForwardedCounter.With(metricsHandler).Record(1)
200 > default: metric_client.go
201 > // TODO: confirmed from metrics, it seems this error does happen at the moment...
202 > // it means some mangled name come here; need to check why
203 > _, err := tqid.NewTaskQueueFamily("", taskQueue.GetName())
204 > if err != nil {
205 c.logger.Info("invalid tq name", tag.Error(err), tag.String("proto", taskQueue.GetName()))
206 metrics.MatchingClientInvalidTaskQueueName.With(metricsHandler).Record(1)
go.temporal.io/server/service/matching/matching_engine.go 29 introduced LOC · 10 ranges

Open complete file

215
216 // Remove unregisters a poller. Cleans up empty worker entries to prevent memory leak. Thread-safe.
217 > func (t *workerPollerTracker) Remove(workerKey, pollerID string) { matching_engine.go
218 > t.lock.Lock()
219 > defer t.lock.Unlock()
220 > util.DeleteFromMap(t.pollers, workerKey, pollerID)
221 > }
222
223 // CancelAll cancels all pollers for a worker and removes the worker entry. Returns cancelled count. Thread-safe.
2868 }
2869
2870 > func (e *matchingEngineImpl) ListNexusEndpoints(ctx context.Context, request *matchingservice.ListNexusEndpointsRequest) (*matchingservice.ListNexusEndpointsResponse, error) { matching_engine.go
2871 > lastKnownVersion := request.LastKnownTableVersion
2872 > // Read API, verify table ownership via membership.
2873 > isOwner, ownershipLostCh, err := e.checkNexusEndpointsOwnership()
2874 > if err != nil {
2875 e.logger.Error("Failed to check Nexus endpoints ownership", tag.Error(err))
2876 return nil, serviceerror.NewAbortedf("cannot verify ownership of Nexus endpoints table: %v", err)
2877 }
2878 > if !isOwner { matching_engine.go
2879 e.logger.Error("Matching node doesn't think it's the Nexus endpoints table owner", tag.Error(err))
2880 return nil, serviceerror.NewAborted("matching node doesn't think it's the Nexus endpoints table owner")
2881 }
2882
2883 > if request.Wait { matching_engine.go
2884 > if request.NextPageToken != nil {
2885 return nil, serviceerror.NewInvalidArgument("request Wait=true and NextPageToken!=nil on ListNexusEndpoints request. waiting is only allowed on first page")
2886 }
2887
2888 // if waiting, send request with unknown table version so we get the newest view of the table
2889 > request.LastKnownTableVersion = 0 matching_engine.go
2890 >
2891 > var cancel context.CancelFunc
2892 > ctx, cancel = contextutil.WithDeadlineBuffer(ctx, e.config.ListNexusEndpointsLongPollTimeout(), returnEmptyTaskTimeBudget)
2893 > defer cancel()
2894 }
2895
2896 > for { matching_engine.go
2897 > resp, tableVersionChanged, err := e.nexusEndpointClient.ListNexusEndpoints(ctx, request)
2898 > if err != nil {
2899 return resp, err
2900 }
2901
2902 > if request.Wait && lastKnownVersion == resp.TableVersion { matching_engine.go
2903 > // long-poll: wait for data to change/appear
2904 > select {
2905 case <-ownershipLostCh:
2906 return nil, serviceerror.NewAborted("Nexus endpoints table ownership lost")
3055 pollerTrackerKey := uuid.NewString()
3056 if workerInstanceKey != "" {
3057 > e.workerInstancePollers.Add(workerInstanceKey, pollerTrackerKey, cancel) matching_engine.go
3058 > }
3059
3060 defer func() {
3061 e.outstandingPollers.Delete(pollerID)
3062 if workerInstanceKey != "" {
3063 > e.workerInstancePollers.Remove(workerInstanceKey, pollerTrackerKey) matching_engine.go
3064 > }
3065 }()
3066 }
go.temporal.io/server/service/matching/task_queue_partition_manager.go 23 introduced LOC · 2 ranges

Open complete file

155 var scaleManager *scaleManager
156 if partition.IsRoot() && e.partitionScalerFactory != nil {
157 > partitionScaler := e.partitionScalerFactory.New( task_queue_partition_manager.go
158 > ns.Name(),
159 > partition.TaskQueue().Name(),
160 > partition.TaskQueue().TaskType(),
161 > )
162 > if partitionScaler != nil {
163 > baseCtx := headers.SetCallerInfo(context.Background(), headers.NewBackgroundLowCallerInfo(ns.Name().String()))
164 > scaleManager = newScaleManager(
165 > baseCtx,
166 > partition,
167 > logger,
168 > metricsHandler,
169 > userDataManager,
170 > e.matchingRawClient,
171 > partitionScaler,
172 > e.timeSource,
173 > tqConfig.PartitionScaleManagerSettings,
174 > tqConfig.NumWritePartitions,
175 > tqConfig.BreakdownMetricsByTaskQueue,
176 > )
177 > }
178 }
179
1908 }
1909 if partitions <= 1 {
1911 > }
1912
1913 // record total-1 as we won't try to load the root partition.
go.temporal.io/server/client/matching/retryable_client_gen.go 18 introduced LOC · 2 ranges

Open complete file

406 request *matchingservice.PollActivityTaskQueueRequest,
407 opts ...grpc.CallOption,
408 > ) (*matchingservice.PollActivityTaskQueueResponse, error) { retryable_client_gen.go
409 > var resp *matchingservice.PollActivityTaskQueueResponse
410 > op := func(ctx context.Context) error {
411 > var err error
412 > resp, err = c.client.PollActivityTaskQueue(ctx, request, opts...)
413 > return err
414 > }
415 > err := backoff.ThrottleRetryContext(ctx, op, c.pollPolicy, c.isRetryable)
416 > return resp, err
417 }
418
436 request *matchingservice.PollWorkflowTaskQueueRequest,
437 opts ...grpc.CallOption,
438 > ) (*matchingservice.PollWorkflowTaskQueueResponse, error) { retryable_client_gen.go
439 > var resp *matchingservice.PollWorkflowTaskQueueResponse
440 > op := func(ctx context.Context) error {
441 > var err error
442 > resp, err = c.client.PollWorkflowTaskQueue(ctx, request, opts...)
443 > return err
444 > }
445 > err := backoff.ThrottleRetryContext(ctx, op, c.pollPolicy, c.isRetryable)
446 > return resp, err
447 }
448
go.temporal.io/server/common/client_cache.go 16 introduced LOC · 4 ranges

Open complete file

75 }
76
77 > func (c *clientCacheImpl) GetClientForClientKey(clientKey string) (any, error) { client_cache.go
78 > c.cacheLock.RLock()
79 > entry, ok := c.clients[clientKey]
80 > c.cacheLock.RUnlock()
81 > if ok {
82 > return entry.client, nil
83 > }
84
85 > c.cacheLock.Lock() client_cache.go
86 > defer c.cacheLock.Unlock()
87 >
88 > entry, ok = c.clients[clientKey]
89 > if ok {
90 return entry.client, nil
91 }
92
93 > client, release, err := c.clientProvider(clientKey) client_cache.go
94 > if err != nil {
95 return nil, err
96 }
97 > c.clients[clientKey] = cachedEntry{client: client, release: release} client_cache.go
98 > return client, nil
99 }
100
go.temporal.io/server/client/matching/loadbalancer.go 15 introduced LOC · 4 ranges

Open complete file

104 taskQueue *tqid.TaskQueue,
105 pc PartitionCounts,
106 > ) *pollToken { loadbalancer.go
107 > tqlb := lb.getTaskQueueLoadBalancer(taskQueue)
108 >
109 > // For read path it's safer to return global default partition count instead of root partition, when we fail to
110 > // map namespace ID to name.
111 > var partitionCount = dynamicconfig.GlobalDefaultNumTaskQueuePartitions
112 >
113 > if pc.Read > 0 {
114 partitionCount = int(pc.Read)
115 > } else { loadbalancer.go
116 > namespaceName, err := lb.namespaceIDToName(namespace.ID(taskQueue.NamespaceId()))
117 > if err == nil {
118 > partitionCount = lb.nReadPartitions(string(namespaceName), taskQueue.Name(), taskQueue.TaskType())
119 > }
120 }
121
122 > if n, ok := testhooks.Get(lb.testHooks, testhooks.MatchingLBForceReadPartition, namespace.ID(taskQueue.NamespaceId())); ok { loadbalancer.go
123 return tqlb.forceReadPartition(partitionCount, n)
124 }
125
126 > return tqlb.pickReadPartition(partitionCount) loadbalancer.go
127 }
128
go.temporal.io/server/common/persistence/persistence_rate_limited_clients.go 13 introduced LOC · 8 ranges

Open complete file

524 ctx context.Context,
525 request *GetTasksRequest,
526 > ) (*GetTasksResponse, error) { persistence_rate_limited_clients.go
527 > if err := allow(ctx, "GetTasks", CallerSegmentMissing, p.systemRateLimiter, p.namespaceRateLimiter, p.shardRateLimiter); err != nil {
528 return nil, err
529 }
530
531 > response, err := p.persistence.GetTasks(ctx, request) persistence_rate_limited_clients.go
532 > return response, err
533 }
534
546 ctx context.Context,
547 request *CreateTaskQueueRequest,
548 > ) (*CreateTaskQueueResponse, error) { persistence_rate_limited_clients.go
549 > if err := allow(ctx, "CreateTaskQueue", CallerSegmentMissing, p.systemRateLimiter, p.namespaceRateLimiter, p.shardRateLimiter); err != nil {
550 return nil, err
551 }
552 > return p.persistence.CreateTaskQueue(ctx, request) persistence_rate_limited_clients.go
553 }
554
566 ctx context.Context,
567 request *GetTaskQueueRequest,
568 > ) (*GetTaskQueueResponse, error) { persistence_rate_limited_clients.go
569 > if err := allow(ctx, "GetTaskQueue", CallerSegmentMissing, p.systemRateLimiter, p.namespaceRateLimiter, p.shardRateLimiter); err != nil {
570 return nil, err
571 }
572 > return p.persistence.GetTaskQueue(ctx, request) persistence_rate_limited_clients.go
573 }
574
596 ctx context.Context,
597 request *GetTaskQueueUserDataRequest,
598 > ) (*GetTaskQueueUserDataResponse, error) { persistence_rate_limited_clients.go
599 > if err := allow(ctx, "GetTaskQueueUserData", CallerSegmentMissing, p.systemRateLimiter, p.namespaceRateLimiter, p.shardRateLimiter); err != nil {
600 return nil, err
601 }
602 > return p.persistence.GetTaskQueueUserData(ctx, request) persistence_rate_limited_clients.go
603 }
604
go.temporal.io/server/service/matching/configs/quotas.go 13 introduced LOC · 1 range

Open complete file

80 return quotas.NewNamespaceRequestRateLimiter(
81 func(req quotas.Request) quotas.RequestRateLimiter {
82 > return quotas.NewPriorityRateLimiterHelper( quotas.go
83 > quotas.NewNamespaceRateBurst(
84 > req.Caller,
85 > namespaceRateFn,
86 > // TODO: We can consider adding a separate burst ratio dynamic config
87 > // on namespace level rate limiter if needed.
88 > quotas.DefaultIncomingNamespaceBurstRatioFn,
89 > ),
90 > operatorRPSRatio,
91 > RequestToPriority,
92 > APIPrioritiesOrdered,
93 > )
94 > },
95 )
96 }
go.temporal.io/server/client/matching/client_gen.go 12 introduced LOC · 4 ranges

Open complete file

346 request *matchingservice.GetTaskQueueUserDataRequest,
347 opts ...grpc.CallOption,
348 > ) (*matchingservice.GetTaskQueueUserDataResponse, error) { client_gen.go
349 >
350 > p, err := tqid.NormalPartitionFromRpcName(request.GetTaskQueue(), request.GetNamespaceId(), request.GetTaskQueueType())
351 > if err != nil {
352 return nil, err
353 }
354
355 > client, err := c.getClientForTaskQueuePartition(p) client_gen.go
356 > if err != nil {
357 return nil, err
358 }
359 > ctx, cancel := c.createLongPollContext(ctx) client_gen.go
360 > defer cancel()
361 > return client.GetTaskQueueUserData(ctx, request, opts...)
362 }
363
417 return nil, err
418 }
419 > ctx, cancel := c.createLongPollContext(ctx) client_gen.go
420 > defer cancel()
421 > return client.ListNexusEndpoints(ctx, request, opts...)
422 }
423
go.temporal.io/server/common/namespace/nsregistry/registry.go 10 introduced LOC · 4 ranges

Open complete file

364 func (r *registry) GetNamespaceID(
365 name namespace.Name,
366 > ) (namespace.ID, error) { registry.go
367 >
368 > ns, err := r.GetNamespace(name)
369 > if err != nil {
370 return "", err
371 }
372 > return ns.ID(), nil registry.go
373 }
374
376 func (r *registry) GetNamespaceName(
377 id namespace.ID,
378 > ) (namespace.Name, error) { registry.go
379 >
380 > ns, err := r.getOrReadthroughNamespaceByID(id)
381 > if err != nil {
382 return "", err
383 }
384 > return ns.Name(), nil registry.go
385 }
386
go.temporal.io/server/common/rpc/rpc.go 9 introduced LOC · 3 ranges

Open complete file

268
269 // createInternodeGRPCConnection creates connection for gRPC calls
270 > func (d *RPCFactory) createInternodeGRPCConnection(hostName string, serviceName primitives.ServiceName) *grpc.ClientConn { rpc.go
271 > var tlsClientConfig *tls.Config
272 > var err error
273 > if d.tlsFactory != nil {
274 tlsClientConfig, err = d.tlsFactory.GetInternodeClientConfig()
275 if err != nil {
278 }
279 }
280 > additionalDialOptions := append([]grpc.DialOption{}, d.perServiceDialOptions[serviceName]...) rpc.go
281 > return d.dial(hostName, tlsClientConfig, append(additionalDialOptions, d.getClientKeepAliveConfig(serviceName))...)
282 }
283
286 }
287
288 > func (d *RPCFactory) CreateMatchingGRPCConnection(rpcAddress string) *grpc.ClientConn { rpc.go
289 > return d.createInternodeGRPCConnection(rpcAddress, primitives.MatchingService)
290 > }
291
292 func (d *RPCFactory) dial(hostName string, tlsClientConfig *tls.Config, dialOptions ...grpc.DialOption) *grpc.ClientConn {
go.temporal.io/server/client/matching/metric_client_gen.go 7 introduced LOC · 2 ranges

Open complete file

252 request *matchingservice.GetTaskQueueUserDataRequest,
253 opts ...grpc.CallOption,
254 > ) (_ *matchingservice.GetTaskQueueUserDataResponse, retError error) { metric_client_gen.go
255 >
256 > metricsHandler, startTime := c.startMetricsRecording(ctx, "MatchingClientGetTaskQueueUserData")
257 > defer func() {
258 > c.finishMetricsRecording(metricsHandler, startTime, retError)
259 > }()
260
261 > return c.client.GetTaskQueueUserData(ctx, request, opts...) metric_client_gen.go
262 }
263