266
ctx context.Context,
267
metadata *persistence.InternalQueueMetadata,
269
>
err := q.txExecute(ctx, "UpdateDLQAckLevel", func(tx sqlplugin.Tx) error {
270
>
271
>
result, err := tx.UpdateQueueMetadata(ctx, &sqlplugin.QueueMetadataRow{
272
>
QueueType: q.getDLQTypeFromQueueType(),
273
>
Data: metadata.Blob.Data,
274
>
DataEncoding: metadata.Blob.EncodingType.String(),
275
>
})
276
>
if err != nil {
277
return serviceerror.NewUnavailablef("UpdateDLQAckLevel operation failed. Error %v", err)
278
}
279
>
rowsAffected, err := result.RowsAffected()
queue.go
280
>
if err != nil {
281
return fmt.Errorf("rowsAffected returned error for DLQ metadata %v: %v", q.queueType, err)
282
}
284
return fmt.Errorf("rowsAffected returned %v DLQ metadata instead of one", rowsAffected)
285
}
287
})
288
290
return serviceerror.NewUnavailable(err.Error())
291
}
293
}
294
295
func (q *sqlQueue) GetDLQAckLevels(
296
ctx context.Context,
297
>
) (*persistence.InternalQueueMetadata, error) {
queue.go
298
>
row, err := q.DB.SelectFromQueueMetadata(ctx, sqlplugin.QueueMetadataFilter{
299
>
QueueType: q.getDLQTypeFromQueueType(),
300
>
})
301
>
if err != nil {
302
return nil, serviceerror.NewUnavailablef("GetDLQAckLevels operation failed. Error %v", err)
303
}
304
305
>
return &persistence.InternalQueueMetadata{
queue.go
306
>
Blob: persistence.NewDataBlob(row.Data, row.DataEncoding),
307
>
Version: row.Version,
308
>
}, nil
309
}
310