159
pageSize int,
160
pageToken []byte,
162
>
// Reading replication tasks need to be quorum level consistent, otherwise we could lose tasks
163
>
// Use negative queue type as the dlq type
164
>
query := q.session.Query(templateGetMessagesFromDLQQuery,
165
>
q.getDLQTypeFromQueueType(),
166
>
firstMessageID,
167
>
lastMessageID,
168
>
).WithContext(ctx)
169
>
iter := query.PageSize(pageSize).PageState(pageToken).Iter()
170
>
171
>
var result []*persistence.QueueMessage
172
>
message := make(map[string]any)
173
>
for iter.MapScan(message) {
174
>
queueMessage := convertQueueMessage(message)
175
>
result = append(result, queueMessage)
176
>
message = make(map[string]any)
177
>
}
178
180
>
if len(iter.PageState()) > 0 {
181
>
nextPageToken = iter.PageState()
182
>
}
183
>
if err := iter.Close(); err != nil {
184
return nil, nil, gocql.ConvertError("ReadMessagesFromDLQ", err)
185
}
186
188
}
189