workflow_handler.go ×8

Frontier kind: Code frontier

unlabeled · c_f2b99965d23d

5 tests · 41968 LOC · 751 files · introduces 0 tests · 144 LOC · 11 files

Introduces — evidence that enters the hierarchy at this concept

Code
39 ranges144 lines · 11 files
Tests
0 tests

Contains — complete concept membership

All code (extent)
9815 ranges41968 lines · 751 files · Browse complete extent
All tests (intent)
5 testsBrowse complete intent

Neighbourhood graph

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

Introduced files, introduced tests, and structurally relevant concept specialization

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

Graph controls are ready.

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

Native relationship evidence

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

Introduced tests

Every collected test enters the hierarchy at exactly one concept.

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

Introduced code

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

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

go.temporal.io/server/service/frontend/workflow_handler.go 35 introduced LOC · 8 ranges

Open complete file

1624 ctx context.Context,
1625 request *workflowservice.RespondActivityTaskCompletedRequest,
1626 > ) (_ *workflowservice.RespondActivityTaskCompletedResponse, retError error) { workflow_handler.go
1627 >
1628 > defer log.CapturePanic(wh.logger, &retError)
1629 >
1630 > if request == nil {
1631 return nil, errRequestNotSet
1632 }
1633 > taskToken, err := wh.tokenSerializer.Deserialize(request.TaskToken) workflow_handler.go
1634 > if err != nil {
1635 return nil, errDeserializingToken
1636 }
1637 > namespaceId := namespace.ID(taskToken.GetNamespaceId()) workflow_handler.go
1638 > namespaceEntry, err := wh.namespaceRegistry.GetNamespaceByID(namespaceId)
1639 > if err != nil {
1640 return nil, err
1641 }
1642 > namespaceName := namespaceEntry.Name().String() workflow_handler.go
1643 >
1644 > if len(taskToken.GetComponentRef()) > 0 && !wh.IsStandaloneActivityEnabled(namespaceName) {
1645 return nil, activity.ErrStandaloneActivityDisabled
1646 }
1647
1648 > if len(request.GetIdentity()) > wh.config.MaxIDLengthLimit() { workflow_handler.go
1649 return nil, errIdentityTooLong
1650 }
1651
1652 > sizeLimitError := wh.config.BlobSizeLimitError(namespaceName) workflow_handler.go
1653 > sizeLimitWarn := wh.config.BlobSizeLimitWarn(namespaceName)
1654 >
1655 > if err := common.CheckEventBlobSizeLimit(
1656 > request.GetResult().Size(),
1657 > sizeLimitWarn,
1658 > sizeLimitError,
1659 > namespaceId.String(),
1660 > taskToken.GetWorkflowId(),
1661 > taskToken.GetRunId(),
1662 > wh.metricsScope(ctx).WithTags(metrics.CommandTypeTag(enumspb.COMMAND_TYPE_UNSPECIFIED.String())),
1663 > wh.throttledLogger,
1664 > "RespondActivityTaskCompleted",
1665 > ); err != nil {
1666 // result exceeds blob size limit, we would record it as failure
1667 failRequest := &workflowservice.RespondActivityTaskFailedRequest{
1677 return nil, err
1678 }
1679 > } else { workflow_handler.go
1680 > _, err = wh.historyClient.RespondActivityTaskCompleted(ctx, &historyservice.RespondActivityTaskCompletedRequest{
1681 > NamespaceId: namespaceId.String(),
1682 > CompleteRequest: request,
1683 > })
1684 > if err != nil {
1685 return nil, err
1686 }
1687 }
1688
1689 > return &workflowservice.RespondActivityTaskCompletedResponse{}, nil workflow_handler.go
1690 }
1691
go.temporal.io/server/api/historyservice/v1/request_response.pb.go 22 introduced LOC · 4 ranges

Open complete file

2482 }
2483
2484 > func (x *RespondActivityTaskCompletedRequest) Reset() { request_response.pb.go
2485 > *x = RespondActivityTaskCompletedRequest{}
2486 > mi := &file_temporal_server_api_historyservice_v1_request_response_proto_msgTypes[24]
2487 > ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
2488 > ms.StoreMessageInfo(mi)
2489 > }
2490
2491 func (x *RespondActivityTaskCompletedRequest) String() string {
2498 mi := &file_temporal_server_api_historyservice_v1_request_response_proto_msgTypes[24]
2499 if x != nil {
2500 > ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) request_response.pb.go
2501 > if ms.LoadMessageInfo() == nil {
2502 > ms.StoreMessageInfo(mi)
2503 > }
2504 > return ms
2505 }
2506 return mi.MessageOf(x)
2532 }
2533
2534 > func (x *RespondActivityTaskCompletedResponse) Reset() { request_response.pb.go
2535 > *x = RespondActivityTaskCompletedResponse{}
2536 > mi := &file_temporal_server_api_historyservice_v1_request_response_proto_msgTypes[25]
2537 > ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
2538 > ms.StoreMessageInfo(mi)
2539 > }
2540
2541 func (x *RespondActivityTaskCompletedResponse) String() string {
2548 mi := &file_temporal_server_api_historyservice_v1_request_response_proto_msgTypes[25]
2549 if x != nil {
2550 > ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) request_response.pb.go
2551 > if ms.LoadMessageInfo() == nil {
2552 > ms.StoreMessageInfo(mi)
2553 > }
2554 > return ms
2555 }
2556 return mi.MessageOf(x)
go.temporal.io/server/client/history/client_gen.go 22 introduced LOC · 4 ranges

Open complete file

1080 request *historyservice.RespondActivityTaskCompletedRequest,
1081 opts ...grpc.CallOption,
1082 > ) (*historyservice.RespondActivityTaskCompletedResponse, error) { client_gen.go
1083 > taskToken, err := c.tokenSerializer.Deserialize(request.GetCompleteRequest().GetTaskToken())
1084 > if err != nil {
1085 return nil, serviceerror.NewInvalidArgument("error deserializing task token")
1086 }
1087 > var namespaceID string client_gen.go
1088 > var businessID string
1089 > if len(taskToken.GetComponentRef()) > 0 {
1090 ref, err := c.tokenSerializer.DeserializeChasmComponentRef(taskToken.GetComponentRef())
1091 if err != nil {
1094 namespaceID = ref.GetNamespaceId()
1095 businessID = ref.GetBusinessId()
1096 > } else { client_gen.go
1097 > namespaceID = request.GetNamespaceId()
1098 > businessID = taskToken.GetWorkflowId()
1099 > }
1100 > shardID := c.shardIDFromWorkflowID(namespaceID, businessID)
1101 >
1102 > var response *historyservice.RespondActivityTaskCompletedResponse
1103 > op := func(ctx context.Context, client historyservice.HistoryServiceClient) error {
1104 > var err error
1105 > ctx, cancel := c.createContext(ctx)
1106 > defer cancel()
1107 > response, err = client.RespondActivityTaskCompleted(ctx, request, opts...)
1108 > return err
1109 > }
1110 > if err := c.executeWithRedirect(ctx, shardID, op); err != nil {
1111 return nil, err
1112 }
1113 > return response, nil client_gen.go
1114 }
1115
go.temporal.io/server/api/historyservice/v1/service_grpc.pb.go 17 introduced LOC · 5 ranges

Open complete file

490 }
491
492 > func (c *historyServiceClient) RespondActivityTaskCompleted(ctx context.Context, in *RespondActivityTaskCompletedRequest, opts ...grpc.CallOption) (*RespondActivityTaskCompletedResponse, error) { service_grpc.pb.go
493 > out := new(RespondActivityTaskCompletedResponse)
494 > err := c.cc.Invoke(ctx, HistoryService_RespondActivityTaskCompleted_FullMethodName, in, out, opts...)
495 > if err != nil {
496 return nil, err
497 }
498 > return out, nil service_grpc.pb.go
499 }
500
1836 }
1837
1838 > func _HistoryService_RespondActivityTaskCompleted_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { service_grpc.pb.go
1839 > in := new(RespondActivityTaskCompletedRequest)
1840 > if err := dec(in); err != nil {
1841 return nil, err
1842 }
1843 > if interceptor == nil { service_grpc.pb.go
1844 return srv.(HistoryServiceServer).RespondActivityTaskCompleted(ctx, in)
1845 }
1846 > info := &grpc.UnaryServerInfo{ service_grpc.pb.go
1847 > Server: srv,
1848 > FullMethod: HistoryService_RespondActivityTaskCompleted_FullMethodName,
1849 > }
1850 > handler := func(ctx context.Context, req interface{}) (interface{}, error) {
1851 > return srv.(HistoryServiceServer).RespondActivityTaskCompleted(ctx, req.(*RespondActivityTaskCompletedRequest))
1852 > }
1853 > return interceptor(ctx, in, info, handler)
1854 }
1855
go.temporal.io/server/service/history/handler.go 14 introduced LOC · 8 ranges

Open complete file

395
396 // RespondActivityTaskCompleted - records completion of an activity task
397 > func (h *Handler) RespondActivityTaskCompleted(ctx context.Context, request *historyservice.RespondActivityTaskCompletedRequest) (*historyservice.RespondActivityTaskCompletedResponse, error) { handler.go
398 > taskToken, err := h.tokenSerializer.Deserialize(request.CompleteRequest.GetTaskToken())
399 > if err != nil {
400 return nil, consts.ErrDeserializingToken
401 }
402
403 > if err := validateTaskToken(taskToken); err != nil { handler.go
404 return nil, h.convertError(err)
405 }
406
407 // Handle standalone activity if component ref is present in the token.
408 > if componentRef := taskToken.GetComponentRef(); len(componentRef) > 0 { handler.go
409 response, _, err := chasm.UpdateComponent(
410 ctx,
423
424 // Handle worklow activity (mutable state backed implementation).
425 > namespaceID := namespace.ID(request.GetNamespaceId()) handler.go
426 > if namespaceID == "" {
427 return nil, h.convertError(errNamespaceNotSet)
428 }
429
430 > shardContext, err := h.controller.GetShardByNamespaceWorkflow(namespaceID, taskToken.GetWorkflowId()) handler.go
431 > if err != nil {
432 return nil, h.convertError(err)
433 }
434 > engine, err := shardContext.GetEngine(ctx) handler.go
435 > if err != nil {
436 return nil, h.convertError(err)
437 }
438
439 > resp, err := engine.RespondActivityTaskCompleted(ctx, request) handler.go
440 > if err != nil {
441 return nil, h.convertError(err)
442 }
443
444 > return resp, nil handler.go
445 }
446
go.temporal.io/server/client/history/retryable_client_gen.go 9 introduced LOC · 1 range

Open complete file

856 request *historyservice.RespondActivityTaskCompletedRequest,
857 opts ...grpc.CallOption,
858 > ) (*historyservice.RespondActivityTaskCompletedResponse, error) { retryable_client_gen.go
859 > var resp *historyservice.RespondActivityTaskCompletedResponse
860 > op := func(ctx context.Context) error {
861 > var err error
862 > resp, err = c.client.RespondActivityTaskCompleted(ctx, request, opts...)
863 > return err
864 > }
865 > err := backoff.ThrottleRetryContext(ctx, op, c.policy, c.isRetryable)
866 > return resp, err
867 }
868
go.temporal.io/server/service/history/workflow/mutable_state_impl.go 9 introduced LOC · 4 ranges

Open complete file

3878 for _, existingValue := range existingValues {
3879 if existingValue == buildId {
3880 > foundBuildId = true mutable_state_impl.go
3881 > }
3882 if !worker_versioning.IsUnversionedOrAssignedBuildIdSearchAttribute(existingValue) &&
3883 !strings.HasPrefix(existingValue, worker_versioning.BuildIdSearchAttributePrefixPinned) {
3884 > newValues = append(newValues, existingValue) mutable_state_impl.go
3885 > }
3886 }
3887
4253
4254 if attributes.UseWorkflowBuildId {
4255 > if ms.GetAssignedBuildId() != "" { mutable_state_impl.go
4256 // only set when using new versioning
4257 ai.BuildIdInfo = &persistencespb.ActivityInfo_UseWorkflowBuildIdInfo_{
4258 UseWorkflowBuildIdInfo: &persistencespb.ActivityInfo_UseWorkflowBuildIdInfo{},
4259 }
4260 > } else { mutable_state_impl.go
4261 > // only set when using old versioning
4262 > ai.UseCompatibleVersion = true
4263 > }
4264 }
4265
go.temporal.io/server/client/history/metric_client_gen.go 7 introduced LOC · 2 ranges

Open complete file

784 request *historyservice.RespondActivityTaskCompletedRequest,
785 opts ...grpc.CallOption,
786 > ) (_ *historyservice.RespondActivityTaskCompletedResponse, retError error) { metric_client_gen.go
787 >
788 > metricsHandler, startTime := c.startMetricsRecording(ctx, "HistoryClientRespondActivityTaskCompleted")
789 > defer func() {
790 > c.finishMetricsRecording(metricsHandler, startTime, retError)
791 > }()
792
793 > return c.client.RespondActivityTaskCompleted(ctx, request, opts...) metric_client_gen.go
794 }
795
go.temporal.io/server/service/history/workflow/workflow_task_state_machine.go 4 introduced LOC · 1 range

Open complete file

1026 // Record persisted workflow task timeout tasks for deletion after successful persistence update.
1027 if task := workflowTask.ScheduleToStartTimeoutTask; task != nil {
1028 > key := task.GetKey() workflow_task_state_machine.go
1029 > if key.FireTime.Sub(workflowTask.ScheduledTime) < maxWorkflowTaskTimeoutToDelete {
1030 > m.ms.BestEffortDeleteTasks[tasks.CategoryTimer] = append(m.ms.BestEffortDeleteTasks[tasks.CategoryTimer], key)
1031 > }
1032 }
1033 if task := workflowTask.StartToCloseTimeoutTask; task != nil {
go.temporal.io/server/temporaltest/logger.go 3 introduced LOC · 1 range

Open complete file

25 }
26
27 > func (tl *testLogger) Debug(msg string, keyvals ...any) { logger.go
28 > tl.logLevel("DEBUG", msg, keyvals)
29 > }
30
31 func (tl *testLogger) Info(msg string, keyvals ...any) {
go.temporal.io/server/client/matching/client.go 2 introduced LOC · 1 range

Open complete file

197 opts ...grpc.CallOption) (*matchingservice.AddWorkflowTaskResponse, error) {
198 if !isPartitionAwareKind(request.GetTaskQueue().GetKind()) {
199 > return c.addWorkflowTask(ctx, PartitionCounts{}, request, opts) client.go
200 > }
201 pkey := c.partitionCache.makeKey(
202 request.GetNamespaceId(),