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