73
}
74
76
>
queueID := shuffle.Bytes(testMatchingTaskTaskQueueID)
77
>
78
>
tasks := []sqlplugin.TasksRowV2{
79
>
s.newRandomTasksRow(queueID, 1, 1),
80
>
s.newRandomTasksRow(queueID, 2, 2),
81
>
s.newRandomTasksRow(queueID, 3, 3),
82
>
}
83
>
result, err := s.store.InsertIntoTasksV2(newExecutionContext(), tasks)
84
>
s.NoError(err)
85
>
rowsAffected, err := result.RowsAffected()
86
>
s.NoError(err)
87
>
s.Equal(3, int(rowsAffected))
88
>
89
>
filter := sqlplugin.TasksFilterV2{
90
>
RangeHash: testMatchingTaskRangeHash,
91
>
TaskQueueID: queueID,
92
>
ExclusiveMaxLevel: &sqlplugin.FairLevel{TaskPass: 3, TaskID: 0},
93
>
Limit: new(10),
94
>
}
95
>
result, err = s.store.DeleteFromTasksV2(newExecutionContext(), filter)
96
>
s.NoError(err)
97
>
rowsAffected, err = result.RowsAffected()
98
>
s.NoError(err)
99
>
s.Equal(2, int(rowsAffected))
100
>
101
>
filter = sqlplugin.TasksFilterV2{
102
>
RangeHash: testMatchingTaskRangeHash,
103
>
TaskQueueID: queueID,
104
>
InclusiveMinLevel: &sqlplugin.FairLevel{TaskPass: 1, TaskID: 0},
105
>
PageSize: new(10),
106
>
}
107
>
rows, err := s.store.SelectFromTasksV2(newExecutionContext(), filter)
108
>
s.NoError(err)
109
>
ids := make([]int64, len(rows))
110
>
for i, r := range rows {
111
>
ids[i] = r.TaskID
112
>
}
113
>
s.Equal([]int64{3}, ids)
114
}
115