122
}
123
125
>
sourceCluster := shuffle.String(testHistoryReplicationTaskDLQSourceCluster)
126
>
shardID := rand.Int31()
127
>
taskID := int64(1)
128
>
129
>
task := s.newRandomReplicationTasksDLQRow(sourceCluster, shardID, taskID)
130
>
result, err := s.store.InsertIntoReplicationDLQTasks(newExecutionContext(), []sqlplugin.ReplicationDLQTasksRow{task})
131
>
s.NoError(err)
132
>
rowsAffected, err := result.RowsAffected()
133
>
s.NoError(err)
134
>
s.Equal(1, int(rowsAffected))
135
>
136
>
rangeFilter := sqlplugin.ReplicationDLQTasksRangeFilter{
137
>
ShardID: shardID,
138
>
SourceClusterName: sourceCluster,
139
>
InclusiveMinTaskID: taskID,
140
>
ExclusiveMaxTaskID: taskID + 1,
141
>
PageSize: 1,
142
>
}
143
>
rows, err := s.store.RangeSelectFromReplicationDLQTasks(newExecutionContext(), rangeFilter)
144
>
s.NoError(err)
145
>
for index := range rows {
146
>
rows[index].ShardID = shardID
147
>
rows[index].SourceClusterName = sourceCluster
148
>
}
149
>
s.Equal([]sqlplugin.ReplicationDLQTasksRow{task}, rows)
150
}
151