( ctx context.Context, blob *commonpb.DataBlob, )
| 345 | } |
| 346 | |
| 347 | func (q *sqlQueue) initializeDLQMetadata( |
| 348 | ctx context.Context, |
| 349 | blob *commonpb.DataBlob, |
| 350 | ) error { |
| 351 | _, err := q.DB.SelectFromQueueMetadata(ctx, sqlplugin.QueueMetadataFilter{ |
| 352 | QueueType: q.getDLQTypeFromQueueType(), |
| 353 | }) |
| 354 | switch err { |
| 355 | case nil: |
| 356 | return nil |
| 357 | case sql.ErrNoRows: |
| 358 | result, err := q.DB.InsertIntoQueueMetadata(ctx, &sqlplugin.QueueMetadataRow{ |
| 359 | QueueType: q.getDLQTypeFromQueueType(), |
| 360 | Data: blob.Data, |
| 361 | DataEncoding: blob.EncodingType.String(), |
| 362 | }) |
| 363 | if err != nil { |
| 364 | return serviceerror.NewUnavailablef("initializeDLQMetadata operation failed. Error %v", err) |
| 365 | } |
| 366 | rowsAffected, err := result.RowsAffected() |
| 367 | if err != nil { |
| 368 | return fmt.Errorf("rowsAffected returned error when initializing DLQ metadata %v: %v", q.queueType, err) |
| 369 | } |
| 370 | if rowsAffected != 1 { |
| 371 | return fmt.Errorf("rowsAffected returned %v DLQ metadata instead of one", rowsAffected) |
| 372 | } |
| 373 | return nil |
| 374 | default: |
| 375 | return err |
| 376 | } |
| 377 | } |
| 378 | |
| 379 | func newQueueRow( |
| 380 | queueType persistence.QueueType, |
no test coverage detected