20
shard historyi.ShardContext,
21
workflowConsistencyChecker api.WorkflowConsistencyChecker,
22
>
) (resp *historyservice.RecordActivityTaskHeartbeatResponse, retError error) {
api.go
23
>
request := req.HeartbeatRequest
24
>
tokenSerializer := tasktoken.NewSerializer()
25
>
token, err0 := tokenSerializer.Deserialize(request.TaskToken)
26
>
if err0 != nil {
27
return nil, consts.ErrDeserializingToken
28
}
29
30
>
_, err := api.GetActiveNamespace(shard, namespace.ID(req.GetNamespaceId()), token.WorkflowId)
api.go
31
>
if err != nil {
32
return nil, err
33
}
34
>
if err := api.SetActivityTaskRunID(ctx, token, workflowConsistencyChecker); err != nil {
api.go
35
return nil, err
36
}
37
38
>
var cancelRequested bool
api.go
39
>
var activityPaused bool
40
>
var activityReset bool
41
>
err = api.GetAndUpdateWorkflowWithNew(
42
>
ctx,
43
>
token.Clock,
44
>
definition.NewWorkflowKey(
45
>
token.NamespaceId,
46
>
token.WorkflowId,
47
>
token.RunId,
48
>
),
49
>
func(workflowLease api.WorkflowLease) (*api.UpdateWorkflowAction, error) {
50
>
mutableState := workflowLease.GetMutableState()
51
>
if !mutableState.IsWorkflowExecutionRunning() {
52
return nil, consts.ErrWorkflowCompleted
53
}
54
55
>
scheduledEventID := token.GetScheduledEventId()
api.go
56
>
if scheduledEventID == common.EmptyEventID { // client call RecordActivityHeartbeatByID, so get scheduledEventID by activityID
57
scheduledEventID, err0 = api.GetActivityScheduledEventID(token.GetActivityId(), mutableState)
58
if err0 != nil {