196
pageSize int,
197
pageToken []byte,
198
>
) ([]*persistence.QueueMessage, []byte, error) {
queue.go
199
>
if len(pageToken) != 0 {
200
>
lastReadMessageID, err := deserializePageToken(pageToken)
201
>
if err != nil {
202
return nil, nil, serviceerror.NewInternalf("invalid next page token %v", pageToken)
203
}
204
>
firstMessageID = lastReadMessageID
queue.go
205
}
206
207
>
rows, err := q.DB.RangeSelectFromMessages(ctx, sqlplugin.QueueMessagesRangeFilter{
queue.go
208
>
QueueType: q.getDLQTypeFromQueueType(),
209
>
MinMessageID: firstMessageID,
210
>
MaxMessageID: lastMessageID,
211
>
PageSize: pageSize,
212
>
})
213
>
if err != nil {
214
return nil, nil, serviceerror.NewUnavailablef("ReadMessagesFromDLQ operation failed. Error %v", err)
215
}
216
217
>
var messages []*persistence.QueueMessage
queue.go
218
>
for _, row := range rows {
219
>
messages = append(messages, &persistence.QueueMessage{
220
>
QueueType: q.getDLQTypeFromQueueType(),
221
>
ID: row.MessageID,
222
>
Data: row.MessagePayload,
223
>
Encoding: row.MessageEncoding,
224
>
})
225
>
}
226
227
>
var newPagingToken []byte
queue.go
228
>
if messages != nil && len(messages) >= pageSize {
229
>
lastReadMessageID := messages[len(messages)-1].ID
230
>
newPagingToken = serializePageToken(lastReadMessageID)
231
>
}
232
>
return messages, newPagingToken, nil
233
}
234