( ctx context.Context, request *persistence.InternalEnqueueMessageRequest, )
| 43 | } |
| 44 | |
| 45 | func (q *queueV2) EnqueueMessage( |
| 46 | ctx context.Context, |
| 47 | request *persistence.InternalEnqueueMessageRequest, |
| 48 | ) (*persistence.InternalEnqueueMessageResponse, error) { |
| 49 | |
| 50 | _, err := q.getQueueMetadata(ctx, q.DB, request.QueueType, request.QueueName) |
| 51 | if err != nil { |
| 52 | return nil, err |
| 53 | } |
| 54 | tx, err := q.DB.BeginTx(ctx) |
| 55 | if err != nil { |
| 56 | return nil, serviceerror.NewUnavailablef( |
| 57 | "EnqueueMessage failed for queue with type: %v and name: %v. BeginTx operation failed. Error: %v", |
| 58 | request.QueueType, |
| 59 | request.QueueName, |
| 60 | err, |
| 61 | ) |
| 62 | } |
| 63 | lastMessageID, ok, err := q.getMaxMessageID(ctx, request.QueueType, request.QueueName, tx) |
| 64 | if err != nil { |
| 65 | rollBackErr := tx.Rollback() |
| 66 | if rollBackErr != nil { |
| 67 | q.logger.Error("transaction rollback error", tag.Error(rollBackErr)) |
| 68 | } |
| 69 | return nil, serviceerror.NewUnavailablef( |
| 70 | "EnqueueMessage failed for queue with type: %v and name: %v. failed to get last messageId. Error: %v", |
| 71 | request.QueueType, |
| 72 | request.QueueName, |
| 73 | err, |
| 74 | ) |
| 75 | } |
| 76 | nextMessageID := int64(persistence.FirstQueueMessageID) |
| 77 | if ok { |
| 78 | nextMessageID = lastMessageID + 1 |
| 79 | } |
| 80 | _, err = tx.InsertIntoQueueV2Messages(ctx, []sqlplugin.QueueV2MessageRow{ |
| 81 | newQueueV2Row(request.QueueType, request.QueueName, nextMessageID, request.Blob), |
| 82 | }) |
| 83 | if err != nil { |
| 84 | rollBackErr := tx.Rollback() |
| 85 | if rollBackErr != nil { |
| 86 | q.logger.Error("transaction rollback error", tag.Error(rollBackErr)) |
| 87 | } |
| 88 | return nil, serviceerror.NewUnavailablef( |
| 89 | "EnqueueMessage failed for queue with type: %v and name: %v. InsertIntoQueueV2Messages operation failed. Error: %v", |
| 90 | request.QueueType, |
| 91 | request.QueueName, |
| 92 | err, |
| 93 | ) |
| 94 | } |
| 95 | |
| 96 | if err := tx.Commit(); err != nil { |
| 97 | return nil, serviceerror.NewUnavailablef( |
| 98 | "EnqueueMessage failed for queue with type: %v and name: %v. commit operation failed. Error: %v", |
| 99 | request.QueueType, |
| 100 | request.QueueName, |
| 101 | err, |
| 102 | ) |
nothing calls this directly
no test coverage detected