309
}
310
311
>
func (q *sqlQueue) getDLQTypeFromQueueType() persistence.QueueType {
queue.go
312
>
return -q.queueType
313
>
}
314
315
func (q *sqlQueue) initializeQueueMetadata(
316
ctx context.Context,
317
blob *commonpb.DataBlob,
319
>
_, err := q.DB.SelectFromQueueMetadata(ctx, sqlplugin.QueueMetadataFilter{
320
>
QueueType: q.queueType,
321
>
})
322
>
switch err {
323
case nil:
324
return nil
326
>
result, err := q.DB.InsertIntoQueueMetadata(ctx, &sqlplugin.QueueMetadataRow{
327
>
QueueType: q.queueType,
328
>
Data: blob.Data,
329
>
DataEncoding: blob.EncodingType.String(),
330
>
})
331
>
if err != nil {
332
return serviceerror.NewUnavailablef("initializeQueueMetadata operation failed. Error %v", err)
333
}
334
>
rowsAffected, err := result.RowsAffected()
queue.go
335
>
if err != nil {
336
return fmt.Errorf("rowsAffected returned error when initializing queue metadata %v: %v", q.queueType, err)
337
}
339
return fmt.Errorf("rowsAffected returned %v queue metadata instead of one", rowsAffected)
340
}
342
default:
343
return err