( ctx context.Context, blob *commonpb.DataBlob, )
| 46 | } |
| 47 | |
| 48 | func (q *sqlQueue) EnqueueMessage( |
| 49 | ctx context.Context, |
| 50 | blob *commonpb.DataBlob, |
| 51 | ) error { |
| 52 | err := q.txExecute(ctx, "EnqueueMessage", func(tx sqlplugin.Tx) error { |
| 53 | lastMessageID, err := tx.GetLastEnqueuedMessageIDForUpdate(ctx, q.queueType) |
| 54 | switch err { |
| 55 | case nil: |
| 56 | _, err = tx.InsertIntoMessages(ctx, []sqlplugin.QueueMessageRow{ |
| 57 | newQueueRow(q.queueType, lastMessageID+1, blob), |
| 58 | }) |
| 59 | return err |
| 60 | case sql.ErrNoRows: |
| 61 | _, err = tx.InsertIntoMessages(ctx, []sqlplugin.QueueMessageRow{ |
| 62 | newQueueRow(q.queueType, persistence.EmptyQueueMessageID+1, blob), |
| 63 | }) |
| 64 | return err |
| 65 | default: |
| 66 | return fmt.Errorf("failed to get last enqueued message id: %v", err) |
| 67 | } |
| 68 | }) |
| 69 | if err != nil { |
| 70 | return serviceerror.NewUnavailable(err.Error()) |
| 71 | } |
| 72 | return nil |
| 73 | } |
| 74 | |
| 75 | func (q *sqlQueue) ReadMessages( |
| 76 | ctx context.Context, |
nothing calls this directly
no test coverage detected