129
}
130
131
>
tqId, tqHash := taskQueueIdAndHash(nidBytes, request.TaskQueue, request.TaskType, request.Subqueue)
task_v2.go
132
>
rows, err := m.DB.SelectFromTasksV2(ctx, sqlplugin.TasksFilterV2{
133
>
RangeHash: tqHash,
134
>
TaskQueueID: tqId,
135
>
InclusiveMinLevel: &inclusiveMinLevel,
136
>
PageSize: &request.PageSize,
137
>
})
138
>
if err != nil {
139
return nil, serviceerror.NewUnavailablef("GetTasks operation failed. Failed to get rows. Error: %v", err)
140
}
141
142
>
response := &persistence.InternalGetTasksResponse{
task_v2.go
143
>
Tasks: make([]*commonpb.DataBlob, len(rows)),
144
>
}
145
>
for i, v := range rows {
146
response.Tasks[i] = persistence.NewDataBlob(v.Data, v.DataEncoding)
147
}
148
>
if len(rows) == request.PageSize {
task_v2.go
149
token, err := serializePageTokenJson(&matchingTaskPageToken{
150
TaskPass: rows[len(rows)-1].TaskPass,