( ctx context.Context, tc sqlplugin.TableCRUD, queueType persistence.QueueV2Type, queueName string, )
| 316 | } |
| 317 | |
| 318 | func (q *queueV2) getQueueMetadata( |
| 319 | ctx context.Context, |
| 320 | tc sqlplugin.TableCRUD, |
| 321 | queueType persistence.QueueV2Type, |
| 322 | queueName string, |
| 323 | ) (*persistencespb.Queue, error) { |
| 324 | |
| 325 | filter := sqlplugin.QueueV2MetadataFilter{ |
| 326 | QueueType: queueType, |
| 327 | QueueName: queueName, |
| 328 | } |
| 329 | var ( |
| 330 | metadata *sqlplugin.QueueV2MetadataRow |
| 331 | err error |
| 332 | ) |
| 333 | switch tc.(type) { |
| 334 | case sqlplugin.Tx: |
| 335 | metadata, err = tc.SelectFromQueueV2MetadataForUpdate(ctx, filter) |
| 336 | default: |
| 337 | metadata, err = tc.SelectFromQueueV2Metadata(ctx, filter) |
| 338 | } |
| 339 | if err != nil { |
| 340 | if errors.Is(err, sql.ErrNoRows) { |
| 341 | return nil, persistence.NewQueueNotFoundError(queueType, queueName) |
| 342 | } |
| 343 | return nil, serviceerror.NewUnavailablef( |
| 344 | "failed to get metadata for queue with type: %v and name: %v. Error: %v", queueType, queueName, err, |
| 345 | ) |
| 346 | } |
| 347 | return q.extractQueueMetadata(metadata) |
| 348 | } |
| 349 | |
| 350 | func (q queueV2) extractQueueMetadata(metadataRow *sqlplugin.QueueV2MetadataRow) (*persistencespb.Queue, error) { |
| 351 | if metadataRow.MetadataEncoding != enumspb.ENCODING_TYPE_PROTO3.String() { |
no test coverage detected