122
}
123
124
>
tqId, tqHash := taskQueueIdAndHash(nidBytes, request.TaskQueue, request.TaskType, request.Subqueue)
task_v1.go
125
>
rows, err := m.DB.SelectFromTasks(ctx, sqlplugin.TasksFilter{
126
>
RangeHash: tqHash,
127
>
TaskQueueID: tqId,
128
>
InclusiveMinTaskID: &inclusiveMinTaskID,
129
>
ExclusiveMaxTaskID: &exclusiveMaxTaskID,
130
>
PageSize: &request.PageSize,
131
>
})
132
>
if err != nil {
133
return nil, serviceerror.NewUnavailablef("GetTasks operation failed. Failed to get rows. Error: %v", err)
134
}
135
136
>
response := &persistence.InternalGetTasksResponse{
task_v1.go
137
>
Tasks: make([]*commonpb.DataBlob, len(rows)),
138
>
}
139
>
for i, v := range rows {
140
response.Tasks[i] = persistence.NewDataBlob(v.Data, v.DataEncoding)
141
}
142
>
if len(rows) == request.PageSize {
task_v1.go
143
nextTaskID := rows[len(rows)-1].TaskID + 1
144
if nextTaskID < exclusiveMaxTaskID {