39
ctx context.Context,
40
request *persistence.InternalCreateTasksRequest,
41
>
) (*persistence.CreateTasksResponse, error) {
task_v2.go
42
>
nidBytes, err := primitives.ParseUUID(request.NamespaceID)
43
>
if err != nil {
44
return nil, serviceerror.NewUnavailable(err.Error())
45
}
46
47
// cache by subqueue to minimize calls to taskQueueIdAndHash
49
>
id []byte
50
>
hash uint32
51
>
}
52
>
cache := make(map[int]pair)
53
>
idAndHash := func(subqueue int) ([]byte, uint32) {
54
>
if pair, ok := cache[subqueue]; ok {
55
>
return pair.id, pair.hash
56
>
}
57
>
id, hash := taskQueueIdAndHash(nidBytes, request.TaskQueue, request.TaskType, subqueue)
58
>
cache[subqueue] = pair{id: id, hash: hash}
59
>
return id, hash
60
}
61
62
>
tasksRows := make([]sqlplugin.TasksRowV2, len(request.Tasks))
task_v2.go
63
>
for i, v := range request.Tasks {
64
>
tqId, tqHash := idAndHash(v.Subqueue)
65
>
tasksRows[i] = sqlplugin.TasksRowV2{
66
>
RangeHash: tqHash,
67
>
TaskQueueID: tqId,
68
>
FairLevel: sqlplugin.FairLevel{
69
>
TaskPass: v.TaskPass,
70
>
TaskID: v.TaskId,
71
>
},
72
>
Data: v.Task.Data,
73
>
DataEncoding: v.Task.EncodingType.String(),
74
>
}
75
>
}
76
>
var resp *persistence.CreateTasksResponse
77
>
err = m.SqlStore.txExecute(ctx, "CreateTasks", func(tx sqlplugin.Tx) error {
78
>
if _, err1 := tx.InsertIntoTasksV2(ctx, tasksRows); err1 != nil {
79
return err1
80
}
81
// Lock task queue before committing.
82
>
tqId, tqHash := idAndHash(persistence.SubqueueZero)
task_v2.go
83
>
if err := lockTaskQueue(ctx,
84
>
tx,
85
>
tqHash,
86
>
tqId,
87
>
request.RangeID,
88
>
sqlplugin.MatchingTaskVersion2,
89
>
); err != nil {
90
return err
91
}
92
>
resp = &persistence.CreateTasksResponse{UpdatedMetadata: false}
task_v2.go
93
>
return nil
94
})
96
}
97