(ctx context.Context, rows *sqlx.Rows)
| 595 | |
| 596 | var msgs []*MessageIncoming |
| 597 | for rows.Next() { |
| 598 | msg, err := c.parseRow(ctx, rows) |
| 599 | if err != nil { |
| 600 | return nil, errors.WithStack(err) |
| 601 | } |
| 602 | msgs = append(msgs, msg) |
| 603 | } |
| 604 | if err := rows.Err(); err != nil { |
| 605 | return nil, errors.WithStack(err) |
| 606 | } |
| 607 | if len(msgs) == 0 { |
| 608 | return nil, sql.ErrNoRows |
| 609 | } |
| 610 | if err := tx.Commit(); err != nil { |
| 611 | return nil, errors.Wrap(err, "commit message consumption") |
| 612 | } |
| 613 | return msgs, nil |
| 614 | } |
| 615 | |
| 616 | func (c *Consumer) parseRow(ctx context.Context, rows *sqlx.Rows) (*MessageIncoming, error) { |
| 617 | var pgMsg pgMessage |
| 618 | if err := rows.Scan( |
| 619 | &pgMsg.ID, |
| 620 | &pgMsg.Payload, |
| 621 | &pgMsg.Metadata, |
| 622 | &pgMsg.Attempt, |
| 623 | &pgMsg.LockedUntil, |
| 624 | ); err != nil { |
| 625 | if isErrorCode(err, undefinedTableErrCode, undefinedColumnErrCode) { |
| 626 | return nil, fatalError{Err: err} |
no test coverage detected