163
ctx context.Context,
164
request *persistence.CompleteTasksLessThanRequest,
166
>
// Require starting from pass 1.
167
>
if request.ExclusiveMaxPass < 1 {
168
return 0, serviceerror.NewInternal("invalid CompleteTasksLessThan request on fair queue")
169
}
170
171
>
nidBytes, err := primitives.ParseUUID(request.NamespaceID)
task_v2.go
172
>
if err != nil {
173
return 0, serviceerror.NewUnavailable(err.Error())
174
}
175
>
tqId, tqHash := taskQueueIdAndHash(nidBytes, request.TaskQueueName, request.TaskType, request.Subqueue)
task_v2.go
176
>
exclusiveMaxLevel := sqlplugin.FairLevel{
177
>
TaskPass: request.ExclusiveMaxPass,
178
>
TaskID: request.ExclusiveMaxTaskID,
179
>
}
180
>
result, err := m.DB.DeleteFromTasksV2(ctx, sqlplugin.TasksFilterV2{
181
>
RangeHash: tqHash,
182
>
TaskQueueID: tqId,
183
>
ExclusiveMaxLevel: &exclusiveMaxLevel,
184
>
Limit: &request.Limit,
185
>
})
186
>
if err != nil {
187
return 0, serviceerror.NewUnavailable(err.Error())
188
}
189
>
nRows, err := result.RowsAffected()
task_v2.go
190
>
if err != nil {
191
return 0, serviceerror.NewUnavailablef("rowsAffected returned error: %v", err)
192
}
194
}